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

scwhittle pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new d1ebefb6f63 [Spanner Change Streams] Claim last processed timestamp 
instead of artificial query end timestamp for unbounded queries (#40209)
d1ebefb6f63 is described below

commit d1ebefb6f6301f6b99358b63dafa1c206cce5849
Author: Jiang Zhu <[email protected]>
AuthorDate: Wed Sep 23 06:08:50 2026 +0000

    [Spanner Change Streams] Claim last processed timestamp instead of 
artificial query end timestamp for unbounded queries (#40209)
---
 .../spanner/changestreams/action/QueryChangeStreamAction.java  | 10 +++++-----
 .../changestreams/action/QueryChangeStreamActionTest.java      |  7 ++++---
 2 files changed, 9 insertions(+), 8 deletions(-)

diff --git 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
index ac4acbb6282..7a5780144ea 100644
--- 
a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
+++ 
b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java
@@ -358,11 +358,11 @@ public class QueryChangeStreamAction {
         "[{}] change stream completed successfully up to {}", token, 
changeStreamQueryEndTimestamp);
 
     if (!stopAfterQuerySucceeds) {
-      // Records stopped being returned for the query due to our artificial 
query end timestamp but
-      // we want to continue processing the partition, resuming from 
changeStreamQueryEndTimestamp.
-      if (!tracker.tryClaim(changeStreamQueryEndTimestamp)) {
-        return ProcessContinuation.stop();
-      }
+      // Leave the tracker at the last claimed position (record or heartbeat)
+      // instead of advancing to the query end timestamp.
+      // This works around spanner backend issue where some child partition 
records
+      // were not sent. In other cases since heartbeating is regular the last
+      // received timestamp will not be far behind the query end timestamp.
       bundleFinalizer.afterBundleCommit(
           Instant.now().plus(BUNDLE_FINALIZER_TIMEOUT),
           updateWatermarkCallback(
diff --git 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
index 13648e1e676..d95569b5e23 100644
--- 
a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
+++ 
b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java
@@ -834,7 +834,6 @@ public class QueryChangeStreamActionTest {
         .thenReturn(changeStreamResultSet);
     when(changeStreamResultSet.next()).thenReturn(false);
     when(watermarkEstimator.currentWatermark()).thenReturn(WATERMARK);
-    when(restrictionTracker.tryClaim(any(Timestamp.class))).thenReturn(true);
 
     final ProcessContinuation result =
         action.run(
@@ -842,7 +841,7 @@ public class QueryChangeStreamActionTest {
     assertEquals(ProcessContinuation.resume(), result);
     assertNotEquals(MAX_INCLUSIVE_END_AT, timestampCaptor.getValue());
 
-    verify(restrictionTracker).tryClaim(timestampCaptor.getValue());
+    verify(restrictionTracker, never()).tryClaim(timestampCaptor.getValue());
     verify(partitionMetadataDao).updateWatermark(PARTITION_TOKEN, 
WATERMARK_TIMESTAMP);
     verify(partitionMetadataDao, never()).updateToFinished(PARTITION_TOKEN);
     verify(metrics, never()).decActivePartitionReadCounter();
@@ -1081,8 +1080,10 @@ public class QueryChangeStreamActionTest {
     long diff = timestampCaptor.getValue().getSeconds() - now.getSeconds();
     assertTrue("Query should be capped at approx 2 minutes (120s)", 
Math.abs(diff - 120) < 10);
 
-    // Crucial: Should RESUME to process the rest later
+    // Crucial: Should RESUME to process the rest later without claiming the 
capped query end
+    // timestamp.
     assertEquals(ProcessContinuation.resume(), result);
+    verify(restrictionTracker, never()).tryClaim(timestampCaptor.getValue());
   }
 
   @Test

Reply via email to