MartijnVisser commented on code in PR #253:
URL:
https://github.com/apache/flink-connector-jdbc/pull/253#discussion_r4230195334
##########
flink-connector-jdbc-core/src/main/java/org/apache/flink/connector/jdbc/core/datastream/source/enumerator/JdbcSourceEnumerator.java:
##########
@@ -65,6 +79,9 @@ public JdbcSourceEnumerator(
@Override
public void start() {
splitterEnumerator.start(connectionProvider);
+ // No split has been handled yet: this is the state a concurrent
checkpoint has to fall
+ // back to while the first batch is still being enumerated on a worker
thread.
+ lastHandledSplitterState = splitterEnumerator.serializableState();
Review Comment:
For a `JdbcSlideTimingParameterProvider` this is `null`, and restoring
`null` after a global failover keeps the advanced provider, so the first batch
is still lost. Can you fall back to `getLatestOptionalState()` in
`TemplateSqlSplitEnumeratorProvider.create()` when no state is set, with a test?
##########
flink-connector-jdbc-core/src/test/java/org/apache/flink/connector/jdbc/core/datastream/source/enumerator/JdbcSourceEnumeratorCheckpointRaceTest.java:
##########
@@ -0,0 +1,210 @@
+/*
+ * 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.flink.connector.jdbc.core.datastream.source.enumerator;
+
+import org.apache.flink.api.connector.source.Boundedness;
+import
org.apache.flink.connector.jdbc.core.datastream.source.enumerator.splitter.SplitterEnumerator;
+import
org.apache.flink.connector.jdbc.core.datastream.source.split.CheckpointedOffset;
+import
org.apache.flink.connector.jdbc.core.datastream.source.split.JdbcSourceSplit;
+import
org.apache.flink.connector.jdbc.datasource.connections.JdbcConnectionProvider;
+import
org.apache.flink.connector.testutils.source.reader.TestingSplitEnumeratorContext;
+
+import org.junit.jupiter.api.Test;
+
+import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.Collections;
+import java.util.List;
+import java.util.concurrent.Callable;
+import java.util.function.BiConsumer;
+
+import static org.assertj.core.api.Assertions.assertThat;
+
+/** Regression test for FLINK-40779. */
+class JdbcSourceEnumeratorCheckpointRaceTest {
+
+ /**
+ * Reproduces the window that FLINK-40779 describes.
+ *
+ * <p>{@code JdbcSourceEnumerator} enumerates splits on a worker thread
and handles the result
+ * on the coordinator thread. A checkpoint taken in between used to record
the splitter state
+ * that had already advanced past the splits still in flight, so those
splits were lost on
+ * restore.
+ */
+ @Test
+ void testCheckpointBetweenEnumerationAndHandlingDoesNotLoseSplits() throws
Exception {
+ final ControllableAsyncContext context = new
ControllableAsyncContext(1);
+ final CountingSplitterEnumerator splitter = new
CountingSplitterEnumerator();
+ final JdbcSourceEnumerator enumerator =
+ new JdbcSourceEnumerator(context, splitter, null, new
ArrayList<>());
+
+ // Schedules the first enumeration, which has not run yet.
+ enumerator.start();
+
+ // The enumeration runs "on the worker thread": the splitter advances
to state 1 and
+ // produces split-000, but the handler that would buffer split-000 has
not run yet.
+ context.runNextCallable();
+ assertThat(splitter.serializableState()).isEqualTo(1);
+
+ // A checkpoint taken in this window must not claim the advanced
state, because split-000 is
+ // not part of it yet.
+ final JdbcSourceEnumeratorState checkpoint =
enumerator.snapshotState(1L);
+ assertThat(checkpoint.getRemainingSplits()).isEmpty();
+
assertThat(checkpoint.getOptionalUserDefinedSplitEnumeratorState()).isEqualTo(0);
+
+ // Restoring from that checkpoint must produce split-000 again rather
than dropping it.
+ final ControllableAsyncContext restoredContext = new
ControllableAsyncContext(1);
+ final JdbcSourceEnumerator restored =
+ new JdbcSourceEnumerator(
+ restoredContext,
+ splitter.restoreState(
+
checkpoint.getOptionalUserDefinedSplitEnumeratorState()),
+ null,
+ new ArrayList<>(checkpoint.getRemainingSplits()));
+ restored.start();
+ restoredContext.runNextCallable();
+ restoredContext.runNextHandler();
+
+ assertThat(restored.snapshotState(2L).getRemainingSplits())
+ .extracting(JdbcSourceSplit::splitId)
+ .containsExactly("split-000");
+ }
+
+ @Test
+ void testStateIsCommittedOnceSplitsHaveBeenHandled() throws Exception {
+ final ControllableAsyncContext context = new
ControllableAsyncContext(1);
+ final CountingSplitterEnumerator splitter = new
CountingSplitterEnumerator();
+ final JdbcSourceEnumerator enumerator =
+ new JdbcSourceEnumerator(context, splitter, null, new
ArrayList<>());
+
+ enumerator.start();
+ context.runNextCallable();
+ context.runNextHandler();
+
+ // split-000 is now owned by the enumerator, so the state that
produced it is checkpointed
+ // together with it.
+ final JdbcSourceEnumeratorState checkpoint =
enumerator.snapshotState(1L);
+ assertThat(checkpoint.getRemainingSplits())
+ .extracting(JdbcSourceSplit::splitId)
+ .containsExactly("split-000");
+
assertThat(checkpoint.getOptionalUserDefinedSplitEnumeratorState()).isEqualTo(1);
+
+ // The next enumeration is already scheduled but has not run: its
state must not leak into
+ // the checkpoint, otherwise the split it produces would be skipped
after a restore.
+ final JdbcSourceEnumeratorState inFlight =
enumerator.snapshotState(2L);
Review Comment:
This half passes on main too: the next enumeration hasn't run, so the
splitter is still at 1. Running the next callable before this snapshot makes it
fail on main.
##########
flink-connector-jdbc-core/src/main/java/org/apache/flink/connector/jdbc/core/datastream/source/enumerator/JdbcSourceEnumerator.java:
##########
@@ -138,22 +155,40 @@ private void preDiscoverSplits() {
while (asyncCallsPending.get() < targetParallelism
&& !splitterEnumerator.isAllSplitsFinished()) {
asyncCallsPending.incrementAndGet();
- context.callAsync(() -> splitterEnumerator.enumerateSplits(),
this::onSplitsDiscovered);
+ context.callAsync(this::enumerateSplitsWithState,
this::onSplitsDiscovered);
}
signalNoMoreSplitsIfDone();
}
- private void onSplitsDiscovered(List<JdbcSourceSplit> splits, Throwable
error) {
+ /**
+ * Enumerates splits and captures the splitter state that corresponds to
exactly those splits.
+ *
+ * <p>Both are captured on the same worker thread invocation, so the
returned state and splits
+ * belong together. Splits must never be separated from the state that
produced them, otherwise
+ * a checkpoint taken in between would either drop splits or duplicate
them.
+ */
+ private EnumerationResult enumerateSplitsWithState() {
+ final List<JdbcSourceSplit> splits =
splitterEnumerator.enumerateSplits();
+ return new EnumerationResult(splits,
splitterEnumerator.serializableState());
+ }
+
+ private void onSplitsDiscovered(EnumerationResult result, Throwable error)
{
asyncCallsPending.decrementAndGet();
if (error != null) {
LOG.error("Failed to discover splits.", error);
preDiscoverSplits();
return;
}
- if (splits != null && !splits.isEmpty()) {
- assignOrBuffer(splits);
+ // The splits are now owned by this enumerator (buffered) or by a
reader, so the state that
+ // produced them may be included in a checkpoint from here on.
Handlers run on the
+ // coordinator thread in completion order, matching the order in which
the enumerations
+ // advanced the splitter state, so committing unconditionally keeps
the two in sync.
Review Comment:
`callAsync` allows concurrent callables, so this order only holds with
today's single worker thread, not because of `synchronized` as the commit
message says. I would prefer capping enumerations in flight at one
(`asyncCallsPending.get() < 1`), which removes that dependency.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]