diff --git a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/sink/cdc/FlinkCdcMultiTableSink.java b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/sink/cdc/FlinkCdcMultiTableSink.java index cdf619725d85..c4dee4dddbd2 100644 --- a/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/sink/cdc/FlinkCdcMultiTableSink.java +++ b/paimon-flink/paimon-flink-cdc/src/main/java/org/apache/paimon/flink/sink/cdc/FlinkCdcMultiTableSink.java @@ -177,6 +177,8 @@ public DataStreamSink sinkFrom( protected CommittableStateManager createCommittableStateManager() { return new RestoreAndFailCommittableStateManager<>( - WrappedManifestCommittableSerializer::new, true); + WrappedManifestCommittableSerializer::new, + true, + StoreMultiCommitter.END_INPUT_HANDLER); } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CombinedTableCompactorSink.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CombinedTableCompactorSink.java index 1f71adc096c0..873f06412a7b 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CombinedTableCompactorSink.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CombinedTableCompactorSink.java @@ -208,6 +208,7 @@ protected DataStreamSink doCommit( protected CommittableStateManager createCommittableStateManager() { return new RestoreAndFailCommittableStateManager<>( WrappedManifestCommittableSerializer::new, - options.get(PARTITION_MARK_DONE_RECOVER_FROM_STATE)); + options.get(PARTITION_MARK_DONE_RECOVER_FROM_STATE), + StoreMultiCommitter.END_INPUT_HANDLER); } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommittableStateManager.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommittableStateManager.java index a39c57986730..bd9f48dab369 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommittableStateManager.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommittableStateManager.java @@ -27,8 +27,13 @@ */ public interface CommittableStateManager extends Serializable { - void initializeState(Committer.Context context, Committer committer) - throws Exception; + /** + * Initializes the state and returns restored committables which must remain pending in the + * operator. + */ + List initializeState( + Committer.Context context, Committer committer) throws Exception; - void snapshotState(List committables) throws Exception; + /** Snapshots pending committables together with the complete end-input state. */ + void snapshotState(List committables, boolean completeEndInput) throws Exception; } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommitterOperator.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommitterOperator.java index 9056440cf6f9..e62b9d55b83c 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommitterOperator.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/CommitterOperator.java @@ -47,7 +47,7 @@ public class CommitterOperator extends AbstractStreamOpe implements OneInputStreamOperator, BoundedOneInput { private static final long serialVersionUID = 1L; - private static final long END_INPUT_CHECKPOINT_ID = Long.MAX_VALUE; + static final long END_INPUT_CHECKPOINT_ID = Long.MAX_VALUE; /** Record all the inputs until commit. */ private final Deque inputs = new ArrayDeque<>(); @@ -85,7 +85,7 @@ public class CommitterOperator extends AbstractStreamOpe private transient long currentWatermark; - private transient boolean endInput; + private transient boolean completeEndInput; private transient String commitUser; @@ -123,7 +123,7 @@ public void initializeState(StateInitializationContext context) throws Exception "Committer Operator parallelism in paimon MUST be one."); this.currentWatermark = Long.MIN_VALUE; - this.endInput = false; + this.completeEndInput = false; // each job can only have one user name and this name must be consistent across restarts // we cannot use job id as commit user name here because user may change job id by creating // a savepoint, stop the job and then resume from savepoint @@ -149,7 +149,14 @@ public void initializeState(StateInitializationContext context) throws Exception .getSpillingDirectoriesPaths()); committer = committerFactory.create(committerContext); - committableStateManager.initializeState(committerContext, committer); + List pendingEndInputCommittables = + committableStateManager.initializeState(committerContext, committer); + for (GlobalCommitT committable : pendingEndInputCommittables) { + Preconditions.checkState( + !committablesPerCheckpoint.containsKey(END_INPUT_CHECKPOINT_ID), + "State manager returned multiple pending end-input committables."); + committablesPerCheckpoint.put(END_INPUT_CHECKPOINT_ID, committable); + } } @Override @@ -170,7 +177,8 @@ public void snapshotState(StateSnapshotContext context) throws Exception { super.snapshotState(context); pollInputs(); committer.snapshotState(); - committableStateManager.snapshotState(committables(committablesPerCheckpoint)); + committableStateManager.snapshotState( + committables(committablesPerCheckpoint), completeEndInput); } private List committables(NavigableMap map) { @@ -179,23 +187,24 @@ private List committables(NavigableMap map) @Override public void endInput() throws Exception { - endInput = true; if (endInputWatermark != null) { currentWatermark = endInputWatermark; } + pollInputs(); + completeEndInput = true; + if (streamingCheckpointEnabled) { return; } - pollInputs(); commitUpToCheckpoint(END_INPUT_CHECKPOINT_ID); } @Override public void notifyCheckpointComplete(long checkpointId) throws Exception { super.notifyCheckpointComplete(checkpointId); - commitUpToCheckpoint(endInput ? END_INPUT_CHECKPOINT_ID : checkpointId); + commitUpToCheckpoint(completeEndInput ? END_INPUT_CHECKPOINT_ID : checkpointId); } private void commitUpToCheckpoint(long checkpointId) throws Exception { diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/EndInputCommittableHandler.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/EndInputCommittableHandler.java new file mode 100644 index 000000000000..f9c225b71f0b --- /dev/null +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/EndInputCommittableHandler.java @@ -0,0 +1,34 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.paimon.flink.sink; + +import java.io.Serializable; + +/** Operations required to keep incomplete end-input committables pending during recovery. */ +public interface EndInputCommittableHandler extends Serializable { + + /** Returns whether the committable belongs to end input. */ + boolean isEndInput(GlobalCommitT committable); + + /** + * Merges two restored end-input committables without rebuilding them or refreshing watermark + * metadata. + */ + GlobalCommitT merge(GlobalCommitT target, GlobalCommitT source); +} diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkWriteSink.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkWriteSink.java index 7b7cb18bb819..bd3476afffd3 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkWriteSink.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/FlinkWriteSink.java @@ -66,7 +66,8 @@ protected CommittableStateManager createCommittableStateMan Options options = table.coreOptions().toConfiguration(); return new RestoreAndFailCommittableStateManager<>( ManifestCommittableSerializer::new, - options.get(PARTITION_MARK_DONE_RECOVER_FROM_STATE)); + options.get(PARTITION_MARK_DONE_RECOVER_FROM_STATE), + StoreCommitter.END_INPUT_HANDLER); } protected static OneInputStreamOperatorFactory @@ -89,6 +90,7 @@ public StreamOperator createStreamOperator(StreamOperatorParameters parameters) Options options = table.coreOptions().toConfiguration(); return new RestoreCommittableStateManager<>( ManifestCommittableSerializer::new, - options.get(PARTITION_MARK_DONE_RECOVER_FROM_STATE)); + options.get(PARTITION_MARK_DONE_RECOVER_FROM_STATE), + StoreCommitter.END_INPUT_HANDLER); } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/NoopCommittableStateManager.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/NoopCommittableStateManager.java index 6e694464c771..84c21d135599 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/NoopCommittableStateManager.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/NoopCommittableStateManager.java @@ -20,6 +20,7 @@ import org.apache.paimon.manifest.ManifestCommittable; +import java.util.Collections; import java.util.List; /** @@ -32,14 +33,15 @@ public class NoopCommittableStateManager implements CommittableStateManager { @Override - public void initializeState( + public List initializeState( Committer.Context context, Committer committer) throws Exception { - // nothing to do + return Collections.emptyList(); } @Override - public void snapshotState(List committables) throws Exception { + public void snapshotState(List committables, boolean completeEndInput) + throws Exception { // nothing to do } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreAndFailCommittableStateManager.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreAndFailCommittableStateManager.java index 8556cb6e7765..04da6c7074fb 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreAndFailCommittableStateManager.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreAndFailCommittableStateManager.java @@ -41,8 +41,9 @@ public class RestoreAndFailCommittableStateManager public RestoreAndFailCommittableStateManager( SerializableSupplier> committableSerializer, - boolean partitionMarkDoneRecoverFromState) { - super(committableSerializer, partitionMarkDoneRecoverFromState); + boolean partitionMarkDoneRecoverFromState, + EndInputCommittableHandler endInputHandler) { + super(committableSerializer, partitionMarkDoneRecoverFromState, endInputHandler); } @Override diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreCommittableStateManager.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreCommittableStateManager.java index 9e5a34ecebf3..e47f3a7bf3b3 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreCommittableStateManager.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/RestoreCommittableStateManager.java @@ -25,16 +25,19 @@ import org.apache.flink.api.common.state.ListState; import org.apache.flink.api.common.state.ListStateDescriptor; +import org.apache.flink.api.common.typeutils.base.BooleanSerializer; import org.apache.flink.api.common.typeutils.base.array.BytePrimitiveArraySerializer; import org.apache.flink.streaming.api.operators.util.SimpleVersionedListState; import java.util.ArrayList; +import java.util.Collections; import java.util.List; /** * A {@link CommittableStateManager} which stores uncommitted {@link ManifestCommittable}s in state. * - *

When the job restarts, these {@link ManifestCommittable}s will be restored and committed. + *

When the job restarts, regular checkpoint committables are restored and committed. An + * incomplete END_INPUT committable remains pending until the operator receives complete end input. */ public class RestoreCommittableStateManager implements CommittableStateManager { @@ -46,19 +49,41 @@ public class RestoreCommittableStateManager private final boolean partitionMarkDoneRecoverFromState; + private final EndInputCommittableHandler endInputHandler; + /** GlobalCommitT state of this job. Used to filter out previous successful commits. */ private ListState streamingCommitterState; + /** Whether every committer operator completed end input before the restored checkpoint. */ + private ListState completeEndInputState; + public RestoreCommittableStateManager( SerializableSupplier> committableSerializer, - boolean partitionMarkDoneRecoverFromState) { + boolean partitionMarkDoneRecoverFromState, + EndInputCommittableHandler endInputHandler) { this.committableSerializer = committableSerializer; this.partitionMarkDoneRecoverFromState = partitionMarkDoneRecoverFromState; + this.endInputHandler = endInputHandler; } @Override - public void initializeState(Committer.Context context, Committer committer) - throws Exception { + public List initializeState( + Committer.Context context, Committer committer) throws Exception { + completeEndInputState = + context.stateStore() + .getUnionListState( + new ListStateDescriptor<>( + "streaming_committer_complete_end_input_state", + BooleanSerializer.INSTANCE)); + boolean hasCompleteEndInputState = false; + boolean restoredCompleteEndInput = true; + for (Boolean value : completeEndInputState.get()) { + hasCompleteEndInputState = true; + restoredCompleteEndInput &= Boolean.TRUE.equals(value); + } + restoredCompleteEndInput &= hasCompleteEndInputState; + completeEndInputState.clear(); + streamingCommitterState = new SimpleVersionedListState<>( context.stateStore() @@ -70,7 +95,26 @@ public void initializeState(Committer.Context context, Committer restored = new ArrayList<>(); streamingCommitterState.get().forEach(restored::add); streamingCommitterState.clear(); + + List pendingEndInput = new ArrayList<>(); + if (!restoredCompleteEndInput) { + restored.removeIf( + committable -> { + if (endInputHandler.isEndInput(committable)) { + if (pendingEndInput.isEmpty()) { + pendingEndInput.add(committable); + } else { + pendingEndInput.set( + 0, + endInputHandler.merge(pendingEndInput.get(0), committable)); + } + return true; + } + return false; + }); + } recover(restored, committer); + return pendingEndInput; } protected int recover(List committables, Committer committer) @@ -79,7 +123,9 @@ protected int recover(List committables, Committer committables) throws Exception { + public void snapshotState(List committables, boolean completeEndInput) + throws Exception { streamingCommitterState.update(committables); + completeEndInputState.update(Collections.singletonList(completeEndInput)); } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreCommitter.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreCommitter.java index 4c353517c3f0..34d75e4c0d46 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreCommitter.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreCommitter.java @@ -41,9 +41,32 @@ import java.util.List; import java.util.Map; +import static org.apache.paimon.utils.Preconditions.checkArgument; + /** {@link Committer} for dynamic store. */ public class StoreCommitter implements Committer { + public static final EndInputCommittableHandler END_INPUT_HANDLER = + new EndInputCommittableHandler() { + private static final long serialVersionUID = 1L; + + @Override + public boolean isEndInput(ManifestCommittable committable) { + return committable.identifier() == CommitterOperator.END_INPUT_CHECKPOINT_ID; + } + + @Override + public ManifestCommittable merge( + ManifestCommittable target, ManifestCommittable source) { + checkArgument( + target.identifier() == source.identifier(), + "Cannot merge committables from different checkpoints %s and %s.", + target.identifier(), + source.identifier()); + return mergeManifestCommittables(target, source); + } + }; + private final TableCommitImpl commit; @Nullable private final CommitterMetrics committerMetrics; private final CommitListeners commitListeners; @@ -106,6 +129,15 @@ public ManifestCommittable combine( return manifestCommittable; } + static ManifestCommittable mergeManifestCommittables( + ManifestCommittable target, ManifestCommittable source) { + for (CommitMessage commitMessage : source.fileCommittables()) { + target.addFileCommittable(commitMessage); + } + source.properties().forEach(target::addProperty); + return target; + } + @Override public void commit(List committables) throws IOException, InterruptedException { diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreMultiCommitter.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreMultiCommitter.java index 855c06f65459..accd19f7e300 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreMultiCommitter.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/StoreMultiCommitter.java @@ -39,6 +39,8 @@ import java.util.Map; import java.util.stream.Collectors; +import static org.apache.paimon.utils.Preconditions.checkArgument; + /** * {@link StoreMultiCommitter} for multiple dynamic store. During the commit process, it will group * the WrappedManifestCommittables by their table identifier and use different committers to commit @@ -47,6 +49,37 @@ public class StoreMultiCommitter implements Committer { + public static final EndInputCommittableHandler END_INPUT_HANDLER = + new EndInputCommittableHandler() { + private static final long serialVersionUID = 1L; + + @Override + public boolean isEndInput(WrappedManifestCommittable committable) { + return committable.checkpointId() == CommitterOperator.END_INPUT_CHECKPOINT_ID; + } + + @Override + public WrappedManifestCommittable merge( + WrappedManifestCommittable target, WrappedManifestCommittable source) { + checkArgument( + target.checkpointId() == source.checkpointId(), + "Cannot merge committables from different checkpoints %s and %s.", + target.checkpointId(), + source.checkpointId()); + for (Map.Entry entry : + source.manifestCommittables().entrySet()) { + ManifestCommittable previous = + target.manifestCommittables().get(entry.getKey()); + if (previous == null) { + target.putManifestCommittable(entry.getKey(), entry.getValue()); + } else { + StoreCommitter.mergeManifestCommittables(previous, entry.getValue()); + } + } + return target; + } + }; + private final Catalog catalog; private final Context context; diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/OperatorBackendStateStore.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/OperatorBackendStateStore.java index e57f5e8a449a..8c8459abd5f1 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/OperatorBackendStateStore.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/OperatorBackendStateStore.java @@ -35,4 +35,9 @@ public OperatorBackendStateStore(OperatorStateStore delegate) { public ListState getListState(ListStateDescriptor descriptor) throws Exception { return delegate.getListState(descriptor); } + + @Override + public ListState getUnionListState(ListStateDescriptor descriptor) throws Exception { + return delegate.getUnionListState(descriptor); + } } diff --git a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/StateStore.java b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/StateStore.java index 22401eb8b90c..c83e2737854d 100644 --- a/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/StateStore.java +++ b/paimon-flink/paimon-flink-common/src/main/java/org/apache/paimon/flink/sink/state/StateStore.java @@ -37,4 +37,12 @@ public interface StateStore { * (subtask or coordinator) and is checkpointed together with that component. */ ListState getListState(ListStateDescriptor descriptor) throws Exception; + + /** + * Returns union list state. Coordinator-side implementations do not rescale and may use regular + * list-state semantics. + */ + default ListState getUnionListState(ListStateDescriptor descriptor) throws Exception { + return getListState(descriptor); + } } diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/BatchWriteGeneratorTagOperatorTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/BatchWriteGeneratorTagOperatorTest.java index 1bc68477a323..d55d7bebe7a5 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/BatchWriteGeneratorTagOperatorTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/BatchWriteGeneratorTagOperatorTest.java @@ -67,7 +67,9 @@ public void testBatchWriteGeneratorTag() throws Exception { table, initialCommitUser, new RestoreAndFailCommittableStateManager<>( - ManifestCommittableSerializer::new, true)); + ManifestCommittableSerializer::new, + true, + StoreCommitter.END_INPUT_HANDLER)); OneInputStreamOperator committerOperator = committerOperatorFactory.createStreamOperator( @@ -143,7 +145,9 @@ public void testBatchWriteGeneratorCustomizedTag() throws Exception { table, initialCommitUser, new RestoreAndFailCommittableStateManager<>( - ManifestCommittableSerializer::new, true)); + ManifestCommittableSerializer::new, + true, + StoreCommitter.END_INPUT_HANDLER)); OneInputStreamOperator committerOperator = committerOperatorFactory.createStreamOperator( diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommitterOperatorTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommitterOperatorTest.java index c533e175fa69..2acbd685ad0b 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommitterOperatorTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/CommitterOperatorTest.java @@ -133,6 +133,59 @@ public void testFailIntentionallyAfterRestore() throws Exception { testHarness.close(); } + @Test + public void testPartialEndInputMergedBeforeOldCheckpointNotification() throws Exception { + FileStoreTable table = createFileStoreTable(); + long timestamp = 1; + + // Only one writer has reached endInput. Its END_INPUT committable must remain pending + // until the commit operator itself receives endInput. + OneInputStreamOperatorTestHarness testHarness = + createRecoverableTestHarness(table); + testHarness.open(); + StreamTableWrite writer1 = + table.newStreamWriteBuilder().withCommitUser(initialCommitUser).newWrite(); + writer1.write(GenericRow.of(1, 10L)); + for (CommitMessage message : writer1.prepareCommit(true, Long.MAX_VALUE)) { + testHarness.processElement(new Committable(Long.MAX_VALUE, message), timestamp++); + } + + OperatorSubtaskState checkpoint = testHarness.snapshot(2, timestamp++); + testHarness.notifyOfCompletedCheckpoint(2); + assertResults(table); + testHarness.close(); + + // Recovery must keep the partial END_INPUT committable pending. There is nothing to commit, + // so RestoreAndFailCommittableStateManager must not trigger an intentional failure. + testHarness = createRecoverableTestHarness(table); + testHarness.initializeState(checkpoint); + testHarness.open(); + assertResults(table); + + // A second writer emits a different committable with the same identifier. It must be + // merged with the restored END_INPUT committable before the final commit. + StreamTableWrite writer2 = + table.newStreamWriteBuilder().withCommitUser(initialCommitUser).newWrite(); + writer2.write(GenericRow.of(2, 20L)); + for (CommitMessage message : writer2.prepareCommit(true, Long.MAX_VALUE)) { + testHarness.processElement(new Committable(Long.MAX_VALUE, message), timestamp++); + } + + testHarness.endInput(); + + // A completion notification for the restored checkpoint may arrive before another + // snapshot. endInput must have drained inputs first, otherwise this notification commits + // only the restored partial MAX_VALUE committable and permanently filters writer2's data. + testHarness.notifyOfCompletedCheckpoint(2); + assertResults(table, "1, 10", "2, 20"); + + testHarness.snapshot(3, timestamp++); + testHarness.notifyOfCompletedCheckpoint(3); + testHarness.close(); + + assertResults(table, "1, 10", "2, 20"); + } + @Test public void testCheckpointAbort() throws Exception { FileStoreTable table = createFileStoreTable(); @@ -467,6 +520,23 @@ private static OperatorSubtaskState writeAndSnapshot( return snapshot; } + private static OperatorSubtaskState writeEndInputAndSnapshot( + FileStoreTable table, + String commitUser, + long timestamp, + long checkpoint, + OneInputStreamOperatorTestHarness testHarness) + throws Exception { + StreamTableWrite write = + table.newStreamWriteBuilder().withCommitUser(commitUser).newWrite(); + write.write(GenericRow.of(1, 10L)); + for (CommitMessage committable : write.prepareCommit(true, checkpoint)) { + testHarness.processElement(new Committable(checkpoint, committable), ++timestamp); + } + testHarness.endInput(); + return testHarness.snapshot(checkpoint, ++timestamp); + } + @Test public void testWatermarkCommit() throws Exception { FileStoreTable table = createFileStoreTable(); @@ -549,7 +619,7 @@ public void testNotTriggerPartitionMarkDownWhenRecoverFromState() throws Excepti createRecoverableTestHarness(table, false); testHarness.open(); OperatorSubtaskState snapshotState = - writeAndSnapshot(table, "commitUser", 1, Long.MAX_VALUE, testHarness); + writeEndInputAndSnapshot(table, "commitUser", 1, Long.MAX_VALUE, testHarness); testHarness.close(); testHarness = createRecoverableTestHarness(table, false); @@ -567,7 +637,6 @@ public void testNotTriggerPartitionMarkDownWhenRecoverFromState() throws Excepti + "writers can start writing based on these new commits."); } - testHarness.notifyOfCompletedCheckpoint(Long.MAX_VALUE); Snapshot snapshot = table.snapshotManager().latestSnapshot(); assertThat(snapshot).isNotNull(); @@ -575,6 +644,41 @@ public void testNotTriggerPartitionMarkDownWhenRecoverFromState() throws Excepti assertThat(table.fileIO().exists(successFile)).isEqualTo(false); } + @Test + public void testTriggerPartitionMarkDownWhenRecoverFromCompleteEndInputState() + throws Exception { + FileStoreTable table = + createFileStoreTable( + options -> { + options.set(CoreOptions.COMMIT_FORCE_CREATE_SNAPSHOT.key(), "true"); + options.set( + CoreOptions.PARTITION_MARK_DONE_WHEN_END_INPUT.key(), "true"); + options.set( + FlinkConnectorOptions.PARTITION_IDLE_TIME_TO_DONE.key(), "1h"); + }, + Collections.singletonList("b")); + + OneInputStreamOperatorTestHarness testHarness = + createRecoverableTestHarness(table, true); + testHarness.open(); + OperatorSubtaskState snapshotState = + writeEndInputAndSnapshot(table, "commitUser", 1, Long.MAX_VALUE, testHarness); + testHarness.close(); + + testHarness = createRecoverableTestHarness(table, true); + try { + testHarness.initializeState(snapshotState); + testHarness.open(); + fail("Expecting intentional exception"); + } catch (Exception e) { + assertThat(e).hasMessageContaining("This exception is intentionally thrown"); + } + + assertThat(table.snapshotManager().latestSnapshot()).isNotNull(); + Path successFile = new Path(table.location(), "b=10/_SUCCESS"); + assertThat(table.fileIO().exists(successFile)).isTrue(); + } + @Test public void testEmptyCommitWithProcessTimeTag() throws Exception { FileStoreTable table = @@ -678,7 +782,9 @@ public void testCommitMetrics() throws Exception { table, null, new RestoreAndFailCommittableStateManager<>( - ManifestCommittableSerializer::new, true)); + ManifestCommittableSerializer::new, + true, + StoreCommitter.END_INPUT_HANDLER)); OneInputStreamOperatorTestHarness testHarness = createTestHarness(operatorFactory); testHarness.open(); @@ -781,7 +887,8 @@ public void testParallelism() throws Exception { null, new RestoreAndFailCommittableStateManager<>( ManifestCommittableSerializer::new, - partitionMarkDownRecoverFromState)); + partitionMarkDownRecoverFromState, + StoreCommitter.END_INPUT_HANDLER)); return createTestHarness(operatorFactory); } diff --git a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/StoreMultiCommitterTest.java b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/StoreMultiCommitterTest.java index 39d0e899ac38..fe2f50b3153c 100644 --- a/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/StoreMultiCommitterTest.java +++ b/paimon-flink/paimon-flink-common/src/test/java/org/apache/paimon/flink/sink/StoreMultiCommitterTest.java @@ -31,6 +31,7 @@ import org.apache.paimon.flink.utils.TestingMetricUtils; import org.apache.paimon.fs.Path; import org.apache.paimon.fs.local.LocalFileIO; +import org.apache.paimon.manifest.ManifestCommittable; import org.apache.paimon.manifest.WrappedManifestCommittable; import org.apache.paimon.options.CatalogOptions; import org.apache.paimon.options.Options; @@ -159,6 +160,43 @@ public void after() throws Exception { // Recoverable operator tests // ------------------------------------------------------------------------ + @Test + public void testMergeEndInputCommittables() throws Exception { + StoreMultiCommitter committer = + new StoreMultiCommitter( + catalogLoader, + Committer.createContext(initialCommitUser, null, true, true, null, 1, 0)); + try { + long checkpointId = Long.MAX_VALUE; + WrappedManifestCommittable target = + new WrappedManifestCommittable(checkpointId, Long.MIN_VALUE); + ManifestCommittable targetTableCommittable = + new ManifestCommittable(checkpointId, Long.MIN_VALUE); + targetTableCommittable.addProperty("target", "value"); + target.putManifestCommittable(firstTable, targetTableCommittable); + + WrappedManifestCommittable source = + new WrappedManifestCommittable(checkpointId, Long.MIN_VALUE); + ManifestCommittable sourceFirstTableCommittable = + new ManifestCommittable(checkpointId, Long.MIN_VALUE); + sourceFirstTableCommittable.addProperty("source", "value"); + source.putManifestCommittable(firstTable, sourceFirstTableCommittable); + source.putManifestCommittable( + secondTable, new ManifestCommittable(checkpointId, Long.MIN_VALUE)); + + WrappedManifestCommittable merged = + StoreMultiCommitter.END_INPUT_HANDLER.merge(target, source); + + assertThat(merged).isSameAs(target); + assertThat(merged.manifestCommittables()).containsOnlyKeys(firstTable, secondTable); + assertThat(merged.manifestCommittables().get(firstTable).properties()) + .containsEntry("target", "value") + .containsEntry("source", "value"); + } finally { + committer.close(); + } + } + @SuppressWarnings("CatchMayIgnoreException") @Test public void testFailIntentionallyAfterRestore() throws Exception { @@ -617,7 +655,9 @@ firstTable, new Committable(cpId, write1.prepareCommit(true, cpId).get(0))), initialCommitUser, context -> new StoreMultiCommitter(catalogLoader, context), new RestoreAndFailCommittableStateManager<>( - WrappedManifestCommittableSerializer::new, true)); + WrappedManifestCommittableSerializer::new, + true, + StoreMultiCommitter.END_INPUT_HANDLER)); return createTestHarness(operator); } @@ -631,13 +671,16 @@ firstTable, new Committable(cpId, write1.prepareCommit(true, cpId).get(0))), context -> new StoreMultiCommitter(catalogLoader, context), new CommittableStateManager() { @Override - public void initializeState( + public List initializeState( Committer.Context context, - Committer committer) {} + Committer committer) { + return Collections.emptyList(); + } @Override public void snapshotState( - List committables) {} + List committables, + boolean completeEndInput) {} }); return createTestHarness(operator); }