Modifier and Type | Method and Description |
---|---|
DataStream<CommittableMessage<FileSinkCommittable>> |
FileSink.addPreCommitTopology(DataStream<CommittableMessage<FileSinkCommittable>> committableStream) |
Modifier and Type | Method and Description |
---|---|
DataStream<CommittableMessage<FileSinkCommittable>> |
FileSink.addPreCommitTopology(DataStream<CommittableMessage<FileSinkCommittable>> committableStream) |
Modifier and Type | Method and Description |
---|---|
<T extends StreamOperator<CommittableMessage<FileSinkCommittable>>> |
CompactorOperatorFactory.createStreamOperator(StreamOperatorParameters<CommittableMessage<FileSinkCommittable>> parameters) |
<T extends StreamOperator<CommittableMessage<FileSinkCommittable>>> |
CompactorOperatorStateHandlerFactory.createStreamOperator(StreamOperatorParameters<CommittableMessage<FileSinkCommittable>> parameters) |
<T extends StreamOperator<Either<CommittableMessage<FileSinkCommittable>,CompactorRequest>>> |
CompactCoordinatorStateHandlerFactory.createStreamOperator(StreamOperatorParameters<Either<CommittableMessage<FileSinkCommittable>,CompactorRequest>> parameters) |
Modifier and Type | Class and Description |
---|---|
class |
CommittableSummary<CommT>
This class tracks the information about committables belonging to one checkpoint coming from one
subtask.
|
class |
CommittableWithLineage<CommT>
Provides metadata.
|
Modifier and Type | Method and Description |
---|---|
CommittableMessage<CommT> |
CommittableMessageSerializer.deserialize(int version,
byte[] serialized) |
Modifier and Type | Method and Description |
---|---|
DataStream<CommittableMessage<CommT>> |
WithPreCommitTopology.addPreCommitTopology(DataStream<CommittableMessage<CommT>> committables)
Intercepts and modifies the committables sent on checkpoint or at end of input.
|
TypeSerializer<CommittableMessage<CommT>> |
CommittableMessageTypeInfo.createSerializer(ExecutionConfig config) |
Class<CommittableMessage<CommT>> |
CommittableMessageTypeInfo.getTypeClass() |
static TypeInformation<CommittableMessage<Void>> |
CommittableMessageTypeInfo.noOutput()
Returns the type information for a
CommittableMessage with no committable. |
static <CommT> TypeInformation<CommittableMessage<CommT>> |
CommittableMessageTypeInfo.of(SerializableSupplier<SimpleVersionedSerializer<CommT>> committableSerializerFactory)
Returns the type information based on the serializer for a
CommittableMessage . |
Modifier and Type | Method and Description |
---|---|
byte[] |
CommittableMessageSerializer.serialize(CommittableMessage<CommT> obj) |
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.
|
void |
WithPostCommitTopology.addPostCommitTopology(DataStream<CommittableMessage<CommT>> committables)
Adds a custom post-commit topology where all committables can be processed.
|
DataStream<CommittableMessage<CommT>> |
WithPreCommitTopology.addPreCommitTopology(DataStream<CommittableMessage<CommT>> committables)
Intercepts and modifies the committables sent on checkpoint or at end of input.
|
Modifier and Type | Method and Description |
---|---|
<T extends StreamOperator<CommittableMessage<CommT>>> |
CommitterOperatorFactory.createStreamOperator(StreamOperatorParameters<CommittableMessage<CommT>> parameters) |
<T extends StreamOperator<CommittableMessage<CommT>>> |
SinkWriterOperatorFactory.createStreamOperator(StreamOperatorParameters<CommittableMessage<CommT>> parameters) |
Modifier and Type | Method and Description |
---|---|
<T extends StreamOperator<CommittableMessage<CommT>>> |
CommitterOperatorFactory.createStreamOperator(StreamOperatorParameters<CommittableMessage<CommT>> parameters) |
<T extends StreamOperator<CommittableMessage<CommT>>> |
SinkWriterOperatorFactory.createStreamOperator(StreamOperatorParameters<CommittableMessage<CommT>> parameters) |
Modifier and Type | Method and Description |
---|---|
void |
CommittableCollector.addMessage(CommittableMessage<CommT> message)
Adds a
CommittableMessage to the collector to hold it until emission. |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.