Package | Description |
---|---|
org.apache.paimon.flink.sink | |
org.apache.paimon.flink.sink.cdc |
Modifier and Type | Method and Description |
---|---|
MultiTableCommittable |
MultiTableCommittableSerializer.deserialize(int committableVersion,
byte[] bytes) |
static MultiTableCommittable |
MultiTableCommittable.fromCommittable(Identifier id,
Committable committable) |
Modifier and Type | Method and Description |
---|---|
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.table.data.RowData,MultiTableCommittable> |
CombinedTableCompactorSink.combinedMultiComacptionWriteOperator(org.apache.flink.streaming.api.environment.CheckpointConfig checkpointConfig,
boolean isStreaming,
String commitUser) |
protected Committer.Factory<MultiTableCommittable,WrappedManifestCommittable> |
CombinedTableCompactorSink.createCommitterFactory(boolean isStreaming) |
org.apache.flink.api.common.typeutils.TypeSerializer<MultiTableCommittable> |
MultiTableCommittableTypeInfo.createSerializer(org.apache.flink.api.common.ExecutionConfig config) |
org.apache.flink.streaming.api.datastream.DataStream<MultiTableCommittable> |
CombinedTableCompactorSink.doWrite(org.apache.flink.streaming.api.datastream.DataStream<org.apache.flink.table.data.RowData> awareBucketTableSource,
org.apache.flink.streaming.api.datastream.DataStream<MultiTableUnawareAppendCompactionTask> unawareBucketTableSource,
String commitUser) |
Class<MultiTableCommittable> |
MultiTableCommittableTypeInfo.getTypeClass() |
Map<Long,List<MultiTableCommittable>> |
StoreMultiCommitter.groupByCheckpoint(Collection<MultiTableCommittable> committables) |
protected List<MultiTableCommittable> |
MultiTablesStoreCompactOperator.prepareCommit(boolean waitCompaction,
long checkpointId) |
protected List<MultiTableCommittable> |
AppendOnlyMultiTableCompactionWorkerOperator.prepareCommit(boolean waitCompaction,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
int |
MultiTableCommittableChannelComputer.channel(MultiTableCommittable multiTableCommittable) |
byte[] |
MultiTableCommittableSerializer.serialize(MultiTableCommittable committable) |
byte[] |
MultiTableCommittableSerializer.serializeCommittable(MultiTableCommittable committable) |
Modifier and Type | Method and Description |
---|---|
WrappedManifestCommittable |
StoreMultiCommitter.combine(long checkpointId,
long watermark,
List<MultiTableCommittable> committables) |
WrappedManifestCommittable |
StoreMultiCommitter.combine(long checkpointId,
long watermark,
WrappedManifestCommittable wrappedManifestCommittable,
List<MultiTableCommittable> committables) |
protected org.apache.flink.streaming.api.datastream.DataStreamSink<?> |
CombinedTableCompactorSink.doCommit(org.apache.flink.streaming.api.datastream.DataStream<MultiTableCommittable> written,
String commitUser) |
Map<Long,List<MultiTableCommittable>> |
StoreMultiCommitter.groupByCheckpoint(Collection<MultiTableCommittable> committables) |
Modifier and Type | Method and Description |
---|---|
protected Committer.Factory<MultiTableCommittable,WrappedManifestCommittable> |
FlinkCdcMultiTableSink.createCommitterFactory() |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<CdcMultiplexRecord,MultiTableCommittable> |
FlinkCdcMultiTableSink.createWriteOperator(StoreSinkWrite.WithWriteBufferProvider writeProvider,
String commitUser) |
protected List<MultiTableCommittable> |
CdcRecordStoreMultiWriteOperator.prepareCommit(boolean waitCompaction,
long checkpointId) |
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.