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