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]

Reply via email to