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()); - } } } }
