Modifier and Type | Method and Description |
---|---|
Optional<SimpleVersionedSerializer<CommT>> |
Sink.getCommittableSerializer()
Deprecated.
Returns the serializer of the committable type.
|
Optional<SimpleVersionedSerializer<GlobalCommT>> |
Sink.getGlobalCommittableSerializer()
Deprecated.
Returns the serializer of the aggregated committable type.
|
Optional<SimpleVersionedSerializer<WriterStateT>> |
Sink.getWriterStateSerializer()
Deprecated.
Any stateful sink needs to provide this state serializer and implement
SinkWriter.snapshotState(long) properly. |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<CommittableT> |
SupportsCommitter.getCommittableSerializer()
Returns the serializer of the committable type.
|
SimpleVersionedSerializer<WriterStateT> |
SupportsWriterState.getWriterStateSerializer()
Any stateful sink needs to provide this state serializer and implement
StatefulSinkWriter.snapshotState(long) properly. |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<EnumChkT> |
Source.getEnumeratorCheckpointSerializer()
Creates the serializer for the
SplitEnumerator checkpoint. |
SimpleVersionedSerializer<SplitT> |
Source.getSplitSerializer()
Creates a serializer for the source splits.
|
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<Collection<NumberSequenceSource.NumberSequenceSplit>> |
NumberSequenceSource.getEnumeratorCheckpointSerializer() |
SimpleVersionedSerializer<NumberSequenceSource.NumberSequenceSplit> |
NumberSequenceSource.getSplitSerializer() |
Modifier and Type | Class and Description |
---|---|
class |
AsyncSinkWriterStateSerializer<RequestEntryT extends Serializable>
Serializer class for
AsyncSinkWriter state. |
Modifier and Type | Class and Description |
---|---|
class |
HybridSourceEnumeratorStateSerializer
The
Serializer for the enumerator state. |
class |
HybridSourceSplitSerializer
Serializes splits by delegating to the source-indexed underlying split serializer.
|
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<HybridSourceEnumeratorState> |
HybridSource.getEnumeratorCheckpointSerializer() |
SimpleVersionedSerializer<HybridSourceSplit> |
HybridSource.getSplitSerializer() |
Modifier and Type | Method and Description |
---|---|
static <SplitT extends SourceSplit,C extends Collection<SplitT>> |
SerdeUtils.deserializeSplitAssignments(byte[] serialized,
SimpleVersionedSerializer<SplitT> splitSerializer,
Function<Integer,C> collectionSupplier)
Deserialize the given bytes returned by
SerdeUtils.serializeSplitAssignments(Map,
SimpleVersionedSerializer) . |
static <SplitT extends SourceSplit,C extends Collection<SplitT>> |
SerdeUtils.serializeSplitAssignments(Map<Integer,C> splitAssignments,
SimpleVersionedSerializer<SplitT> splitSerializer)
Serialize a mapping from subtask ids to lists of assigned splits.
|
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<Collection<NumberSequenceSource.NumberSequenceSplit>> |
DataGeneratorSource.getEnumeratorCheckpointSerializer() |
SimpleVersionedSerializer<NumberSequenceSource.NumberSequenceSplit> |
DataGeneratorSource.getSplitSerializer() |
Modifier and Type | Class and Description |
---|---|
class |
FileSinkCommittableSerializer
Versioned serializer for
FileSinkCommittable . |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<FileSinkCommittable> |
FileSink.getCommittableSerializer() |
SimpleVersionedSerializer<String> |
FileSinkProgram.KeyBucketAssigner.getSerializer() |
SimpleVersionedSerializer<FileSinkCommittable> |
FileSink.getWriteResultSerializer() |
SimpleVersionedSerializer<FileWriterBucketState> |
FileSink.getWriterStateSerializer() |
Constructor and Description |
---|
FileSinkCommittableSerializer(SimpleVersionedSerializer<InProgressFileWriter.PendingFileRecoverable> pendingFileSerializer,
SimpleVersionedSerializer<InProgressFileWriter.InProgressFileRecoverable> inProgressFileSerializer) |
FileSinkCommittableSerializer(SimpleVersionedSerializer<InProgressFileWriter.PendingFileRecoverable> pendingFileSerializer,
SimpleVersionedSerializer<InProgressFileWriter.InProgressFileRecoverable> inProgressFileSerializer) |
Modifier and Type | Class and Description |
---|---|
class |
CompactorRequestSerializer
Versioned serializer for
CompactorRequest . |
Constructor and Description |
---|
CompactCoordinator(FileCompactStrategy strategy,
SimpleVersionedSerializer<FileSinkCommittable> committableSerializer) |
CompactCoordinatorStateHandler(SimpleVersionedSerializer<FileSinkCommittable> committableSerializer) |
CompactorOperator(FileCompactStrategy strategy,
SimpleVersionedSerializer<FileSinkCommittable> committableSerializer,
FileCompactor fileCompactor,
BucketWriter<?,String> bucketWriter) |
CompactorOperatorStateHandler(SimpleVersionedSerializer<FileSinkCommittable> committableSerializer,
BucketWriter<?,String> bucketWriter) |
CompactorRequestSerializer(SimpleVersionedSerializer<FileSinkCommittable> committableSerializer) |
Modifier and Type | Class and Description |
---|---|
class |
FileWriterBucketStateSerializer
A
SimpleVersionedSerializer used to serialize the BucketState . |
Constructor and Description |
---|
FileWriterBucketStateSerializer(SimpleVersionedSerializer<InProgressFileWriter.InProgressFileRecoverable> inProgressFileRecoverableSerializer,
SimpleVersionedSerializer<InProgressFileWriter.PendingFileRecoverable> pendingFileRecoverableSerializer) |
FileWriterBucketStateSerializer(SimpleVersionedSerializer<InProgressFileWriter.InProgressFileRecoverable> inProgressFileRecoverableSerializer,
SimpleVersionedSerializer<InProgressFileWriter.PendingFileRecoverable> pendingFileRecoverableSerializer) |
Modifier and Type | Class and Description |
---|---|
class |
FileSourceSplitSerializer
A serializer for the
FileSourceSplit . |
class |
PendingSplitsCheckpointSerializer<T extends FileSourceSplit>
A serializer for the
PendingSplitsCheckpoint . |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<PendingSplitsCheckpoint<SplitT>> |
AbstractFileSource.getEnumeratorCheckpointSerializer() |
abstract SimpleVersionedSerializer<SplitT> |
AbstractFileSource.getSplitSerializer() |
SimpleVersionedSerializer<FileSourceSplit> |
FileSource.getSplitSerializer() |
Constructor and Description |
---|
PendingSplitsCheckpointSerializer(SimpleVersionedSerializer<T> splitSerializer) |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<String> |
FileSystemTableSink.TableBucketAssigner.getSerializer() |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<NoOpEnumState> |
FromElementsSource.getEnumeratorCheckpointSerializer() |
SimpleVersionedSerializer<FromElementsSplit> |
FromElementsSource.getSplitSerializer() |
Modifier and Type | Class and Description |
---|---|
class |
NoOpEnumStateSerializer
Mock enumerator state serializer.
|
Modifier and Type | Class and Description |
---|---|
class |
FromElementsSplitSerializer
The split serializer for the
FromElementsSource . |
Modifier and Type | Class and Description |
---|---|
class |
ContinuousHivePendingSplitsCheckpointSerializer
SerDe for
ContinuousHivePendingSplitsCheckpoint . |
class |
HiveSourceSplitSerializer
SerDe for
HiveSourceSplit . |
class |
HiveTablePartitionSerializer
SerDe for
HiveTablePartition . |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<PendingSplitsCheckpoint<HiveSourceSplit>> |
HiveSource.getEnumeratorCheckpointSerializer() |
SimpleVersionedSerializer<HiveSourceSplit> |
HiveSource.getSplitSerializer() |
Constructor and Description |
---|
ContinuousHivePendingSplitsCheckpointSerializer(SimpleVersionedSerializer<HiveSourceSplit> splitSerDe) |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<RecoverableWriter.CommitRecoverable> |
RecoverableWriter.getCommitRecoverableSerializer()
The serializer for the CommitRecoverable types created in this writer.
|
SimpleVersionedSerializer<RecoverableWriter.ResumeRecoverable> |
RecoverableWriter.getResumeRecoverableSerializer()
The serializer for the ResumeRecoverable types created in this writer.
|
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<RecoverableWriter.CommitRecoverable> |
LocalRecoverableWriter.getCommitRecoverableSerializer() |
SimpleVersionedSerializer<RecoverableWriter.ResumeRecoverable> |
LocalRecoverableWriter.getResumeRecoverableSerializer() |
Modifier and Type | Method and Description |
---|---|
static <T> T |
SimpleVersionedSerialization.readVersionAndDeSerialize(SimpleVersionedSerializer<T> serializer,
byte[] bytes)
Deserializes the version and datum from a byte array.
|
static <T> T |
SimpleVersionedSerialization.readVersionAndDeSerialize(SimpleVersionedSerializer<T> serializer,
DataInputView in)
Deserializes the version and datum from a stream.
|
static <T> List<T> |
SimpleVersionedSerialization.readVersionAndDeserializeList(SimpleVersionedSerializer<T> serializer,
DataInputView in)
Deserializes the version and data from a stream.
|
static <T> byte[] |
SimpleVersionedSerialization.writeVersionAndSerialize(SimpleVersionedSerializer<T> serializer,
T datum)
Serializes the version and datum into a byte array.
|
static <T> void |
SimpleVersionedSerialization.writeVersionAndSerialize(SimpleVersionedSerializer<T> serializer,
T datum,
DataOutputView out)
Serializes the version and datum into a stream.
|
static <T> void |
SimpleVersionedSerialization.writeVersionAndSerializeList(SimpleVersionedSerializer<T> serializer,
List<T> data,
DataOutputView out)
Serializes the version and data into a stream.
|
Constructor and Description |
---|
SimpleVersionedSerializerTypeSerializerProxy(SerializableSupplier<SimpleVersionedSerializer<T>> serializerSupplier) |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<RecoverableWriter.CommitRecoverable> |
GSRecoverableWriter.getCommitRecoverableSerializer() |
SimpleVersionedSerializer<RecoverableWriter.ResumeRecoverable> |
GSRecoverableWriter.getResumeRecoverableSerializer() |
Modifier and Type | Class and Description |
---|---|
class |
OSSRecoverableSerializer
Serializer implementation for a
OSSRecoverable . |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<RecoverableWriter.CommitRecoverable> |
OSSRecoverableWriter.getCommitRecoverableSerializer() |
SimpleVersionedSerializer<RecoverableWriter.ResumeRecoverable> |
OSSRecoverableWriter.getResumeRecoverableSerializer() |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<RecoverableWriter.CommitRecoverable> |
S3RecoverableWriter.getCommitRecoverableSerializer() |
SimpleVersionedSerializer<RecoverableWriter.ResumeRecoverable> |
S3RecoverableWriter.getResumeRecoverableSerializer() |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<T> |
MasterTriggerRestoreHook.createCheckpointDataSerializer()
Creates a the serializer to (de)serializes the data stored by this hook.
|
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<RecoverableWriter.CommitRecoverable> |
HadoopRecoverableWriter.getCommitRecoverableSerializer() |
SimpleVersionedSerializer<RecoverableWriter.ResumeRecoverable> |
HadoopRecoverableWriter.getResumeRecoverableSerializer() |
Modifier and Type | Class and Description |
---|---|
class |
GenericJobEventSerializer
Serializer for
JobEvent instances that uses Flink's InstantiationUtil for
serialization and deserialization. |
Modifier and Type | Method and Description |
---|---|
void |
SplitAssignmentTracker.restoreState(SimpleVersionedSerializer<SplitT> splitSerializer,
byte[] assignmentData)
Restore the state of the SplitAssignmentTracker.
|
byte[] |
SplitAssignmentTracker.snapshotState(SimpleVersionedSerializer<SplitT> splitSerializer)
Take a snapshot of the split assignments.
|
Constructor and Description |
---|
SourceCoordinatorContext(SourceCoordinatorProvider.CoordinatorExecutorThreadFactory coordinatorThreadFactory,
int numWorkerThreads,
OperatorCoordinator.Context operatorCoordinatorContext,
SimpleVersionedSerializer<SplitT> splitSerializer,
boolean supportsConcurrentExecutionAttempts) |
Modifier and Type | Method and Description |
---|---|
List<SplitT> |
AddSplitEvent.splits(SimpleVersionedSerializer<SplitT> splitSerializer) |
Constructor and Description |
---|
AddSplitEvent(List<SplitT> splits,
SimpleVersionedSerializer<SplitT> splitSerializer) |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<String> |
StreamSQLTestProgram.KeyBucketAssigner.getSerializer() |
Modifier and Type | Class and Description |
---|---|
class |
CommittableMessageSerializer<CommT>
The serializer to serialize
CommittableMessage s in custom operators. |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<WriterResultT> |
SupportsPreCommitTopology.getWriteResultSerializer()
Returns the serializer of the WriteResult type.
|
default SimpleVersionedSerializer<CommT> |
WithPreCommitTopology.getWriteResultSerializer()
Deprecated.
Defaults to
SupportsCommitter.getCommittableSerializer() for backward compatibility. |
Modifier and Type | Method and Description |
---|---|
static <CommT> void |
StandardSinkTopologies.addGlobalCommitter(DataStream<CommittableMessage<CommT>> committables,
SerializableSupplier<Committer<CommT>> committerFactory,
SerializableSupplier<SimpleVersionedSerializer<CommT>> committableSerializer)
Adds a global committer to the pipeline that runs as final operator with a parallelism of
one.
|
static <CommT> TypeInformation<CommittableMessage<CommT>> |
CommittableMessageTypeInfo.of(SerializableSupplier<SimpleVersionedSerializer<CommT>> committableSerializerFactory)
Returns the type information based on the serializer for a
CommittableMessage . |
Constructor and Description |
---|
CommittableMessageSerializer(SimpleVersionedSerializer<CommT> committableSerializer) |
Modifier and Type | Class and Description |
---|---|
static class |
OutputStreamBasedPartFileWriter.OutputStreamBasedInProgressFileRecoverableSerializer
The serializer for
OutputStreamBasedPartFileWriter.OutputStreamBasedInProgressFileRecoverable . |
static class |
OutputStreamBasedPartFileWriter.OutputStreamBasedPendingFileRecoverableSerializer
The serializer for
OutputStreamBasedPartFileWriter.OutputStreamBasedPendingFileRecoverable . |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<RecoverableWriter.CommitRecoverable> |
OutputStreamBasedPartFileWriter.OutputStreamBasedPendingFileRecoverableSerializer.getCommitSerializer() |
SimpleVersionedSerializer<InProgressFileWriter.InProgressFileRecoverable> |
WriterProperties.getInProgressFileRecoverableSerializer() |
SimpleVersionedSerializer<InProgressFileWriter.PendingFileRecoverable> |
WriterProperties.getPendingFileRecoverableSerializer() |
SimpleVersionedSerializer<RecoverableWriter.ResumeRecoverable> |
OutputStreamBasedPartFileWriter.OutputStreamBasedInProgressFileRecoverableSerializer.getResumeSerializer() |
SimpleVersionedSerializer<BucketID> |
BucketAssigner.getSerializer() |
Constructor and Description |
---|
WriterProperties(SimpleVersionedSerializer<InProgressFileWriter.InProgressFileRecoverable> inProgressFileRecoverableSerializer,
SimpleVersionedSerializer<InProgressFileWriter.PendingFileRecoverable> pendingFileRecoverableSerializer,
boolean supportsResume) |
WriterProperties(SimpleVersionedSerializer<InProgressFileWriter.InProgressFileRecoverable> inProgressFileRecoverableSerializer,
SimpleVersionedSerializer<InProgressFileWriter.PendingFileRecoverable> pendingFileRecoverableSerializer,
boolean supportsResume) |
Modifier and Type | Class and Description |
---|---|
class |
SimpleVersionedStringSerializer
A
SimpleVersionedSerializer implementation for Strings. |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<String> |
BasePathBucketAssigner.getSerializer() |
SimpleVersionedSerializer<String> |
DateTimeBucketAssigner.getSerializer() |
Constructor and Description |
---|
SourceOperator(FunctionWithException<SourceReaderContext,SourceReader<OUT,SplitT>,Exception> readerFactory,
OperatorEventGateway operatorEventGateway,
SimpleVersionedSerializer<SplitT> splitSerializer,
WatermarkStrategy<OUT> watermarkStrategy,
ProcessingTimeService timeService,
Configuration configuration,
String localHostname,
boolean emitProgressiveWatermarks,
StreamTask.CanEmitBatchOfRecordsChecker canEmitBatchOfRecords) |
Constructor and Description |
---|
SimpleVersionedListState(ListState<byte[]> rawState,
SimpleVersionedSerializer<T> serializer)
Creates a new SimpleVersionedListState that reads and writes bytes from the given raw
ListState with the given serializer.
|
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<GlobalCommT> |
SinkV1Adapter.GlobalCommitterAdapter.getGlobalCommittableSerializer() |
Modifier and Type | Class and Description |
---|---|
class |
CommittableCollectorSerializer<CommT>
The serializer for the
CommittableCollector . |
Modifier and Type | Method and Description |
---|---|
static <T> List<T> |
SinkV1CommittableDeserializer.readVersionAndDeserializeList(SimpleVersionedSerializer<T> serializer,
DataInputView in) |
Constructor and Description |
---|
CommittableCollectorSerializer(SimpleVersionedSerializer<CommT> committableSerializer,
int subtaskId,
int numberOfSubtasks,
SinkCommitterMetricGroup metricGroup) |
Modifier and Type | Method and Description |
---|---|
SimpleVersionedSerializer<SocketSource.DummyCheckpoint> |
SocketSource.getEnumeratorCheckpointSerializer() |
SimpleVersionedSerializer<SocketSource.DummySplit> |
SocketSource.getSplitSerializer() |
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.