From 5d31039789fadcd00a642eeb648d2ce81a094454 Mon Sep 17 00:00:00 2001 From: Sam Whittle Date: Tue, 26 May 2026 15:30:34 +0200 Subject: [PATCH] [Cloud Spanner Change Streams] Fix inverted evaluation of cancelQueryOnHeartbeat --- .../spanner/changestreams/action/HeartbeatRecordAction.java | 2 +- .../changestreams/action/HeartbeatRecordActionTest.java | 4 ++-- 2 files changed, 3 insertions(+), 3 deletions(-) diff --git a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordAction.java b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordAction.java index 1b66a548b3d2..773f54a15d24 100644 --- a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordAction.java +++ b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordAction.java @@ -104,6 +104,6 @@ public Optional run( return Optional.empty(); } // no new data, finish reading data - return cancelQueryOnHeartbeat ? Optional.empty() : Optional.of(ProcessContinuation.resume()); + return cancelQueryOnHeartbeat ? Optional.of(ProcessContinuation.resume()) : Optional.empty(); } } diff --git a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordActionTest.java b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordActionTest.java index adfc4ea35d48..48fd7c30a1a8 100644 --- a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordActionTest.java +++ b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/action/HeartbeatRecordActionTest.java @@ -232,7 +232,7 @@ public void testEndTimestampNotReachedOnCancellingAction() { watermarkEstimator, endTimestamp); - assertEquals(Optional.empty(), maybeContinuation); + assertEquals(Optional.of(ProcessContinuation.resume()), maybeContinuation); verify(watermarkEstimator).setWatermark(new Instant(timestamp.toSqlTimestamp().getTime())); } @@ -254,7 +254,7 @@ public void testEndTimestampNotReachedOnAction() { watermarkEstimator, endTimestamp); - assertEquals(Optional.of(ProcessContinuation.resume()), maybeContinuation); + assertEquals(Optional.empty(), maybeContinuation); verify(watermarkEstimator).setWatermark(new Instant(timestamp.toSqlTimestamp().getTime())); } }