Skip to content

Commit d1ebefb

Browse files
authored
[Spanner Change Streams] Claim last processed timestamp instead of artificial query end timestamp for unbounded queries (#40209)
1 parent 2906a89 commit d1ebefb

2 files changed

Lines changed: 9 additions & 8 deletions

File tree

‎sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamAction.java‎

Lines changed: 5 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -358,11 +358,11 @@ public ProcessContinuation run(
358358
"[{}] change stream completed successfully up to {}", token, changeStreamQueryEndTimestamp);
359359

360360
if (!stopAfterQuerySucceeds) {
361-
// Records stopped being returned for the query due to our artificial query end timestamp but
362-
// we want to continue processing the partition, resuming from changeStreamQueryEndTimestamp.
363-
if (!tracker.tryClaim(changeStreamQueryEndTimestamp)) {
364-
return ProcessContinuation.stop();
365-
}
361+
// Leave the tracker at the last claimed position (record or heartbeat)
362+
// instead of advancing to the query end timestamp.
363+
// This works around spanner backend issue where some child partition records
364+
// were not sent. In other cases since heartbeating is regular the last
365+
// received timestamp will not be far behind the query end timestamp.
366366
bundleFinalizer.afterBundleCommit(
367367
Instant.now().plus(BUNDLE_FINALIZER_TIMEOUT),
368368
updateWatermarkCallback(

‎sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/QueryChangeStreamActionTest.java‎

Lines changed: 4 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -834,15 +834,14 @@ public void testQueryChangeStreamFinishedWithResume() {
834834
.thenReturn(changeStreamResultSet);
835835
when(changeStreamResultSet.next()).thenReturn(false);
836836
when(watermarkEstimator.currentWatermark()).thenReturn(WATERMARK);
837-
when(restrictionTracker.tryClaim(any(Timestamp.class))).thenReturn(true);
838837

839838
final ProcessContinuation result =
840839
action.run(
841840
partition, restrictionTracker, outputReceiver, watermarkEstimator, bundleFinalizer);
842841
assertEquals(ProcessContinuation.resume(), result);
843842
assertNotEquals(MAX_INCLUSIVE_END_AT, timestampCaptor.getValue());
844843

845-
verify(restrictionTracker).tryClaim(timestampCaptor.getValue());
844+
verify(restrictionTracker, never()).tryClaim(timestampCaptor.getValue());
846845
verify(partitionMetadataDao).updateWatermark(PARTITION_TOKEN, WATERMARK_TIMESTAMP);
847846
verify(partitionMetadataDao, never()).updateToFinished(PARTITION_TOKEN);
848847
verify(metrics, never()).decActivePartitionReadCounter();
@@ -1081,8 +1080,10 @@ public void testQueryChangeStreamWithMutableChangeStreamCappedEndTimestamp() {
10811080
long diff = timestampCaptor.getValue().getSeconds() - now.getSeconds();
10821081
assertTrue("Query should be capped at approx 2 minutes (120s)", Math.abs(diff - 120) < 10);
10831082

1084-
// Crucial: Should RESUME to process the rest later
1083+
// Crucial: Should RESUME to process the rest later without claiming the capped query end
1084+
// timestamp.
10851085
assertEquals(ProcessContinuation.resume(), result);
1086+
verify(restrictionTracker, never()).tryClaim(timestampCaptor.getValue());
10861087
}
10871088

10881089
@Test

0 commit comments

Comments
 (0)