Package | Description |
---|---|
org.apache.flink.table.store.connector.sink |
Modifier and Type | Method and Description |
---|---|
Committable |
CommittableSerializer.deserialize(int committableVersion,
byte[] bytes) |
Modifier and Type | Method and Description |
---|---|
org.apache.flink.api.common.typeutils.TypeSerializer<Committable> |
CommittableTypeInfo.createSerializer(org.apache.flink.api.common.ExecutionConfig config) |
protected abstract org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.table.data.RowData,Committable> |
FlinkSink.createWriteOperator(StoreSinkWrite.Provider writeProvider,
boolean isStreaming) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.table.data.RowData,Committable> |
FileStoreSink.createWriteOperator(StoreSinkWrite.Provider writeProvider,
boolean isStreaming) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.table.data.RowData,Committable> |
CompactorSink.createWriteOperator(StoreSinkWrite.Provider writeProvider,
boolean isStreaming) |
Class<Committable> |
CommittableTypeInfo.getTypeClass() |
protected List<Committable> |
StoreWriteOperator.prepareCommit(boolean doCompaction,
long checkpointId) |
protected List<Committable> |
StoreCompactOperator.prepareCommit(boolean doCompaction,
long checkpointId) |
protected abstract List<Committable> |
PrepareCommitOperator.prepareCommit(boolean doCompaction,
long checkpointId) |
List<Committable> |
StoreSinkWriteImpl.prepareCommit(boolean doCompaction,
long checkpointId) |
List<Committable> |
FullChangelogStoreSinkWrite.prepareCommit(boolean doCompaction,
long checkpointId) |
Modifier and Type | Method and Description |
---|---|
byte[] |
CommittableSerializer.serialize(Committable committable) |
Modifier and Type | Method and Description |
---|---|
ManifestCommittable |
Committer.combine(long checkpointId,
List<Committable> committables)
Compute an aggregated committable from a list of committables.
|
ManifestCommittable |
StoreCommitter.combine(long checkpointId,
List<Committable> committables) |
void |
CommitterOperator.processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<Committable> element) |
void |
StoreWriteOperator.setup(org.apache.flink.streaming.runtime.tasks.StreamTask<?,?> containingTask,
org.apache.flink.streaming.api.graph.StreamConfig config,
org.apache.flink.streaming.api.operators.Output<org.apache.flink.streaming.runtime.streamrecord.StreamRecord<Committable>> output) |
Copyright © 2019–2023 The Apache Software Foundation. All rights reserved.