Modifier and Type | Method and Description |
---|---|
List<Committable> |
UnawareBucketCompactor.prepareCommit(boolean waitCompaction,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
List<Committable> |
ChangelogCompactTask.doCompact(FileStoreTable table) |
Modifier and Type | Method and Description |
---|---|
void |
ChangelogCompactCoordinateOperator.processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<Committable> record) |
void |
ChangelogCompactWorkerOperator.processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<org.apache.flink.types.Either<Committable,ChangelogCompactTask>> record) |
Modifier and Type | Method and Description |
---|---|
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
StoreCompactOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
AppendBypassCompactWorkerOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
RowDataStoreWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
AppendOnlySingleTableCompactionWorkerOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
DynamicBucketRowWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
Modifier and Type | Method and Description |
---|---|
Committable |
CommittableSerializer.deserialize(int committableVersion,
byte[] bytes) |
Committable |
MultiTableCommittableSerializer.deserializeCommittable(int committableVersion,
byte[] bytes) |
Modifier and Type | Method and Description |
---|---|
protected Committer.Factory<Committable,ManifestCommittable> |
UnawareBucketCompactionSink.createCommitterFactory() |
protected Committer.Factory<Committable,ManifestCommittable> |
CompactorSink.createCommitterFactory() |
protected abstract Committer.Factory<Committable,ManifestCommittable> |
FlinkSink.createCommitterFactory() |
protected Committer.Factory<Committable,ManifestCommittable> |
FlinkWriteSink.createCommitterFactory() |
org.apache.flink.api.common.typeutils.TypeSerializer<Committable> |
CommittableTypeInfo.createSerializer(org.apache.flink.api.common.ExecutionConfig config)
Do not annotate with
@override here to maintain compatibility with Flink 2.0+. |
org.apache.flink.api.common.typeutils.TypeSerializer<Committable> |
CommittableTypeInfo.createSerializer(SerializerConfig config)
Do not annotate with
@override here to maintain compatibility with Flink 1.18-. |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<InternalRow,Committable> |
RowUnawareBucketSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<UnawareAppendCompactionTask,Committable> |
UnawareBucketCompactionSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<org.apache.flink.table.data.RowData,Committable> |
CompactorSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<ManifestEntry,Committable> |
RewriteFileIndexSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected abstract org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<T,Committable> |
FlinkSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<org.apache.flink.api.java.tuple.Tuple2<InternalRow,Integer>,Committable> |
RowDynamicBucketSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<InternalRow,Committable> |
FixedBucketSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
org.apache.flink.streaming.api.datastream.DataStream<Committable> |
UnawareBucketSink.doWrite(org.apache.flink.streaming.api.datastream.DataStream<T> input,
String initialCommitUser,
Integer parallelism) |
org.apache.flink.streaming.api.datastream.DataStream<Committable> |
FlinkSink.doWrite(org.apache.flink.streaming.api.datastream.DataStream<T> input,
String commitUser,
Integer parallelism) |
Class<Committable> |
CommittableTypeInfo.getTypeClass() |
Map<Long,List<Committable>> |
StoreCommitter.groupByCheckpoint(Collection<Committable> committables) |
List<Committable> |
StoreSinkWriteImpl.prepareCommit(boolean waitCompaction,
long checkpointId) |
List<Committable> |
GlobalFullCompactionSinkWrite.prepareCommit(boolean waitCompaction,
long checkpointId) |
protected List<Committable> |
StoreCompactOperator.prepareCommit(boolean waitCompaction,
long checkpointId) |
List<Committable> |
StoreSinkWrite.prepareCommit(boolean waitCompaction,
long checkpointId) |
protected List<Committable> |
AppendCompactWorkerOperator.prepareCommit(boolean waitCompaction,
long checkpointId) |
protected List<Committable> |
RowDataStoreWriteOperator.prepareCommit(boolean waitCompaction,
long checkpointId) |
protected List<Committable> |
TableWriteOperator.prepareCommit(boolean waitCompaction,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
static MultiTableCommittable |
MultiTableCommittable.fromCommittable(Identifier id,
Committable committable) |
byte[] |
CommittableSerializer.serialize(Committable committable) |
Modifier and Type | Method and Description |
---|---|
ManifestCommittable |
StoreCommitter.combine(long checkpointId,
long watermark,
List<Committable> committables) |
ManifestCommittable |
StoreCommitter.combine(long checkpointId,
long watermark,
ManifestCommittable manifestCommittable,
List<Committable> committables) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
StoreCompactOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
AppendBypassCompactWorkerOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
RowDataStoreWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
AppendOnlySingleTableCompactionWorkerOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
DynamicBucketRowWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
protected org.apache.flink.streaming.api.datastream.DataStreamSink<?> |
FlinkSink.doCommit(org.apache.flink.streaming.api.datastream.DataStream<Committable> written,
String commitUser) |
Map<Long,List<Committable>> |
StoreCommitter.groupByCheckpoint(Collection<Committable> committables) |
void |
AppendBypassCompactWorkerOperator.processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<org.apache.flink.types.Either<Committable,UnawareAppendCompactionTask>> element) |
Constructor and Description |
---|
AppendCompactWorkerOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters,
FileStoreTable table,
String commitUser) |
RowDataStoreWriteOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters,
FileStoreTable table,
LogSinkFunction logSinkFunction,
StoreSinkWrite.Provider storeSinkWriteProvider,
String initialCommitUser) |
TableWriteOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters,
FileStoreTable table,
StoreSinkWrite.Provider storeSinkWriteProvider,
String initialCommitUser) |
Modifier and Type | Method and Description |
---|---|
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
CdcRecordStoreWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
CdcDynamicBucketWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
CdcUnawareBucketWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
Modifier and Type | Method and Description |
---|---|
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<CdcRecord,Committable> |
CdcFixedBucketSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<org.apache.flink.api.java.tuple.Tuple2<CdcRecord,Integer>,Committable> |
CdcDynamicBucketSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<CdcRecord,Committable> |
CdcUnawareBucketSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
org.apache.flink.streaming.api.datastream.DataStream<Committable> |
CdcUnawareBucketSink.doWrite(org.apache.flink.streaming.api.datastream.DataStream<CdcRecord> input,
String initialCommitUser,
Integer parallelism) |
Modifier and Type | Method and Description |
---|---|
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
CdcRecordStoreWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
CdcDynamicBucketWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
<T extends org.apache.flink.streaming.api.operators.StreamOperator<Committable>> |
CdcUnawareBucketWriteOperator.Factory.createStreamOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters) |
Constructor and Description |
---|
CdcRecordStoreWriteOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<Committable> parameters,
FileStoreTable table,
StoreSinkWrite.Provider storeSinkWriteProvider,
String initialCommitUser) |
Modifier and Type | Method and Description |
---|---|
protected org.apache.flink.streaming.api.operators.OneInputStreamOperatorFactory<org.apache.flink.api.java.tuple.Tuple2<InternalRow,Integer>,Committable> |
GlobalDynamicBucketSink.createWriteOperatorFactory(StoreSinkWrite.Provider writeProvider,
String commitUser) |
Copyright © 2023–2025 The Apache Software Foundation. All rights reserved.