This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-12412-a5efad47c59b13606792f4ca3f8357d6dfd3e829
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit c8c497376030b2254830a734fe88e2bc6576ac69
Author: Gökhan Elbistan <[email protected]>
AuthorDate: Fri Oct 2 07:42:39 2026 +0000

    [Bug][Connector-V2][Http] Release the checkpoint lock while waiting out 
poll_interval_millis (#12412)
    
    Co-authored-by: Claude Opus 5.5 <[email protected]>
---
 .../source/reader/GraphQLSourceHttpReader.java     | 11 ++--
 .../seatunnel/http/source/HttpSourceReader.java    | 12 +++--
 .../http/HttpSourceReaderInternalPollNextTest.java | 60 ++++++++++++++++++++++
 .../prometheus/source/PrometheusSourceReader.java  | 11 ++--
 4 files changed, 82 insertions(+), 12 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-graphql/src/main/java/org/apache/seatunnel/connectors/seatunnel/graphql/source/reader/GraphQLSourceHttpReader.java
 
b/seatunnel-connectors-v2/connector-graphql/src/main/java/org/apache/seatunnel/connectors/seatunnel/graphql/source/reader/GraphQLSourceHttpReader.java
index 9ff0021231..69c98e2ca4 100644
--- 
a/seatunnel-connectors-v2/connector-graphql/src/main/java/org/apache/seatunnel/connectors/seatunnel/graphql/source/reader/GraphQLSourceHttpReader.java
+++ 
b/seatunnel-connectors-v2/connector-graphql/src/main/java/org/apache/seatunnel/connectors/seatunnel/graphql/source/reader/GraphQLSourceHttpReader.java
@@ -83,6 +83,13 @@ public class GraphQLSourceHttpReader extends 
AbstractSingleSplitReader<SeaTunnel
         synchronized (output.getCheckpointLock()) {
             internalPollNext(output);
         }
+        // Wait out the poll interval only after the checkpoint lock is
+        // released. Sleeping while holding it kept the checkpoint barrier
+        // out for the whole interval, so a streaming job's checkpoint expired.
+        if (!Boundedness.BOUNDED.equals(context.getBoundedness())
+                && httpParameter.getPollIntervalMillis() > 0) {
+            Thread.sleep(httpParameter.getPollIntervalMillis());
+        }
     }
 
     @Override
@@ -94,10 +101,6 @@ public class GraphQLSourceHttpReader extends 
AbstractSingleSplitReader<SeaTunnel
                 // signal to the source that we have reached the end of the 
data.
                 log.info("Closed the bounded http source");
                 context.signalNoMoreElement();
-            } else {
-                if (httpParameter.getPollIntervalMillis() > 0) {
-                    Thread.sleep(httpParameter.getPollIntervalMillis());
-                }
             }
         }
     }
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/source/HttpSourceReader.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/source/HttpSourceReader.java
index 3ed214ce67..7efd9c6a24 100644
--- 
a/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/source/HttpSourceReader.java
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/http/source/HttpSourceReader.java
@@ -405,6 +405,14 @@ public class HttpSourceReader extends 
AbstractSingleSplitReader<SeaTunnelRow> {
         synchronized (output.getCheckpointLock()) {
             internalPollNext(output);
         }
+        // Wait out the poll interval only after the checkpoint lock is
+        // released. Sleeping while holding it kept the checkpoint barrier
+        // out for the whole interval, so a streaming job's checkpoint expired.
+        boolean finished =
+                Boundedness.BOUNDED.equals(context.getBoundedness()) && 
noMoreElementFlag;
+        if (!finished && httpParameter.getPollIntervalMillis() > 0) {
+            Thread.sleep(httpParameter.getPollIntervalMillis());
+        }
     }
 
     @Override
@@ -441,10 +449,6 @@ public class HttpSourceReader extends 
AbstractSingleSplitReader<SeaTunnelRow> {
                 // signal to the source that we have reached the end of the 
data.
                 log.info("Closed the bounded http source");
                 context.signalNoMoreElement();
-            } else {
-                if (httpParameter.getPollIntervalMillis() > 0) {
-                    Thread.sleep(httpParameter.getPollIntervalMillis());
-                }
             }
         }
     }
diff --git 
a/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/HttpSourceReaderInternalPollNextTest.java
 
b/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/HttpSourceReaderInternalPollNextTest.java
index 13ac4dc8ac..5e4bda1ab1 100644
--- 
a/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/HttpSourceReaderInternalPollNextTest.java
+++ 
b/seatunnel-connectors-v2/connector-http/connector-http-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/http/HttpSourceReaderInternalPollNextTest.java
@@ -43,6 +43,11 @@ import org.mockito.MockitoAnnotations;
 
 import java.util.HashMap;
 import java.util.Map;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 
 import static org.mockito.ArgumentMatchers.any;
@@ -183,6 +188,61 @@ public class HttpSourceReaderInternalPollNextTest {
         httpSourceReader.close();
     }
 
+    @Test
+    public void testPollIntervalDoesNotHoldCheckpointLock() throws Exception {
+        httpParameter.setPollIntervalMillis(60_000);
+        when(context.getBoundedness()).thenReturn(Boundedness.UNBOUNDED);
+        CountDownLatch polled = new CountDownLatch(1);
+        when(httpClientProvider.execute(
+                        anyString(), anyString(), any(), any(), any(), 
anyBoolean()))
+                .thenAnswer(
+                        invocation -> {
+                            polled.countDown();
+                            return new HttpResponse(200, "[]");
+                        });
+        Object checkpointLock = new Object();
+        Collector<SeaTunnelRow> lockingCollector =
+                new Collector<SeaTunnelRow>() {
+                    @Override
+                    public void collect(SeaTunnelRow record) {}
+
+                    @Override
+                    public Object getCheckpointLock() {
+                        return checkpointLock;
+                    }
+                };
+
+        httpSourceReader =
+                new HttpSourceReader(
+                        httpParameter, context, deserializationSchema, 
jsonField, null);
+        httpSourceReader.open();
+        httpSourceReader.setHttpClient(httpClientProvider);
+
+        ExecutorService executor = Executors.newFixedThreadPool(2);
+        try {
+            executor.submit(
+                    () -> {
+                        httpSourceReader.pollNext(lockingCollector);
+                        return null;
+                    });
+            // The request runs while pollNext holds the checkpoint lock
+            Assertions.assertTrue(polled.await(10, TimeUnit.SECONDS));
+            // A checkpoint barrier must get the lock during the 60 s poll
+            // interval, not only after it
+            Future<Boolean> barrier =
+                    executor.submit(
+                            () -> {
+                                synchronized (checkpointLock) {
+                                    return true;
+                                }
+                            });
+            Assertions.assertTrue(barrier.get(10, TimeUnit.SECONDS));
+        } finally {
+            executor.shutdownNow();
+            httpSourceReader.close();
+        }
+    }
+
     @AfterEach
     public void tearDown() throws Exception {
         mock.close();
diff --git 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/source/PrometheusSourceReader.java
 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/source/PrometheusSourceReader.java
index f15accdf7c..1f8043ef6c 100644
--- 
a/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/source/PrometheusSourceReader.java
+++ 
b/seatunnel-connectors-v2/connector-prometheus/src/main/java/org/apache/seatunnel/connectors/seatunnel/prometheus/source/PrometheusSourceReader.java
@@ -94,6 +94,13 @@ public class PrometheusSourceReader extends 
AbstractSingleSplitReader<SeaTunnelR
         synchronized (output.getCheckpointLock()) {
             internalPollNext(output);
         }
+        // Wait out the poll interval only after the checkpoint lock is
+        // released. Sleeping while holding it kept the checkpoint barrier
+        // out for the whole interval, so a streaming job's checkpoint expired.
+        if (!Boundedness.BOUNDED.equals(context.getBoundedness())
+                && httpParameter.getPollIntervalMillis() > 0) {
+            Thread.sleep(httpParameter.getPollIntervalMillis());
+        }
     }
 
     @Override
@@ -105,10 +112,6 @@ public class PrometheusSourceReader extends 
AbstractSingleSplitReader<SeaTunnelR
                 // signal to the source that we have reached the end of the 
data.
                 log.info("Closed the bounded http source");
                 context.signalNoMoreElement();
-            } else {
-                if (httpParameter.getPollIntervalMillis() > 0) {
-                    Thread.sleep(httpParameter.getPollIntervalMillis());
-                }
             }
         }
     }

Reply via email to