Skip to content

Commit 9583fe9

Browse files
authored
[Cloud Spanner Change Streams] Fix inverted evaluation of cancelQueryOnHeartbeat (#38695)
This meant that low latency mode for heartbeats was enabled by default and disabled in low latency mode instead of the desired opposite.
1 parent d553b78 commit 9583fe9

2 files changed

Lines changed: 3 additions & 3 deletions

File tree

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

Lines changed: 1 addition & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -104,6 +104,6 @@ public Optional<ProcessContinuation> run(
104104
return Optional.empty();
105105
}
106106
// no new data, finish reading data
107-
return cancelQueryOnHeartbeat ? Optional.empty() : Optional.of(ProcessContinuation.resume());
107+
return cancelQueryOnHeartbeat ? Optional.of(ProcessContinuation.resume()) : Optional.empty();
108108
}
109109
}

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

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -232,7 +232,7 @@ public void testEndTimestampNotReachedOnCancellingAction() {
232232
watermarkEstimator,
233233
endTimestamp);
234234

235-
assertEquals(Optional.empty(), maybeContinuation);
235+
assertEquals(Optional.of(ProcessContinuation.resume()), maybeContinuation);
236236
verify(watermarkEstimator).setWatermark(new Instant(timestamp.toSqlTimestamp().getTime()));
237237
}
238238

@@ -254,7 +254,7 @@ public void testEndTimestampNotReachedOnAction() {
254254
watermarkEstimator,
255255
endTimestamp);
256256

257-
assertEquals(Optional.of(ProcessContinuation.resume()), maybeContinuation);
257+
assertEquals(Optional.empty(), maybeContinuation);
258258
verify(watermarkEstimator).setWatermark(new Instant(timestamp.toSqlTimestamp().getTime()));
259259
}
260260
}

0 commit comments

Comments
 (0)