Skip to content

Commit b9f07c9

Browse files
authored
[Spanner Change Streams] Ensure the partition watermark is monotonic by reading within the transaction (#36463)
Bundle finalizations which are used to update the watermark are not ordered, so we must guard against stale watermark updates to ensure the watermark is correct.
1 parent 65ee225 commit b9f07c9

2 files changed

Lines changed: 44 additions & 5 deletions

File tree

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

Lines changed: 14 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@
3232
import com.google.cloud.Timestamp;
3333
import com.google.cloud.spanner.DatabaseClient;
3434
import com.google.cloud.spanner.Dialect;
35+
import com.google.cloud.spanner.Key;
3536
import com.google.cloud.spanner.Mutation;
3637
import com.google.cloud.spanner.Options;
3738
import com.google.cloud.spanner.ResultSet;
@@ -528,14 +529,25 @@ public Void updateToFinished(String partitionToken) {
528529
}
529530

530531
/**
531-
* Update the partition watermark to the given timestamp.
532+
* Update the partition watermark to the given timestamp iff the partition watermark in metadata
533+
* table is smaller than the given watermark.
532534
*
533535
* @param partitionToken the partition unique identifier
534536
* @param watermark the new partition watermark
535537
* @return the commit timestamp of the read / write transaction
536538
*/
537539
public Void updateWatermark(String partitionToken, Timestamp watermark) {
538-
transaction.buffer(createUpdateMetadataWatermarkMutationFrom(partitionToken, watermark));
540+
Struct row =
541+
transaction.readRow(
542+
metadataTableName, Key.of(partitionToken), Collections.singleton(COLUMN_WATERMARK));
543+
if (row == null) {
544+
LOG.error("[{}] Failed to read Watermark column", partitionToken);
545+
return null;
546+
}
547+
Timestamp partitionWatermark = row.getTimestamp(COLUMN_WATERMARK);
548+
if (partitionWatermark.compareTo(watermark) < 0) {
549+
transaction.buffer(createUpdateMetadataWatermarkMutationFrom(partitionToken, watermark));
550+
}
539551
return null;
540552
}
541553

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

Lines changed: 30 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -36,6 +36,8 @@
3636
import com.google.cloud.spanner.TransactionContext;
3737
import com.google.cloud.spanner.TransactionRunner;
3838
import com.google.cloud.spanner.Value;
39+
import java.time.Duration;
40+
import java.time.Instant;
3941
import java.util.Collections;
4042
import java.util.Map;
4143
import org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata;
@@ -238,14 +240,39 @@ public void testInTransactionContextUpdateToFinished() {
238240
@Test
239241
public void testInTransactionContextUpdateWatermark() {
240242
ArgumentCaptor<Mutation> mutation = ArgumentCaptor.forClass(Mutation.class);
241-
doNothing().when(transaction).buffer(mutation.capture());
242-
assertNull(inTransactionContext.updateWatermark(PARTITION_TOKEN, WATERMARK));
243+
when(transaction.readRow(any(), any(), any()))
244+
.thenReturn(
245+
Struct.newBuilder()
246+
.set(PartitionMetadataAdminDao.COLUMN_WATERMARK)
247+
.to(WATERMARK)
248+
.build());
249+
Instant largerWatermark = WATERMARK.toSqlTimestamp().toInstant().plus(Duration.ofSeconds(1));
250+
assertNull(
251+
inTransactionContext.updateWatermark(
252+
PARTITION_TOKEN,
253+
Timestamp.ofTimeSecondsAndNanos(
254+
largerWatermark.getEpochSecond(), largerWatermark.getNano())));
255+
verify(transaction).buffer(mutation.capture());
243256
Map<String, Value> mutationValueMap = mutation.getValue().asMap();
244257
assertEquals(
245258
PARTITION_TOKEN,
246259
mutationValueMap.get(PartitionMetadataAdminDao.COLUMN_PARTITION_TOKEN).getString());
247260
assertEquals(
248-
WATERMARK, mutationValueMap.get(PartitionMetadataAdminDao.COLUMN_WATERMARK).getTimestamp());
261+
Timestamp.ofTimeSecondsAndNanos(
262+
largerWatermark.getEpochSecond(), largerWatermark.getNano()),
263+
mutationValueMap.get(PartitionMetadataAdminDao.COLUMN_WATERMARK).getTimestamp());
264+
}
265+
266+
@Test
267+
public void testInTransactionContextDoNotUpdateWatermark() {
268+
when(transaction.readRow(any(), any(), any()))
269+
.thenReturn(
270+
Struct.newBuilder()
271+
.set(PartitionMetadataAdminDao.COLUMN_WATERMARK)
272+
.to(WATERMARK)
273+
.build());
274+
assertNull(inTransactionContext.updateWatermark(PARTITION_TOKEN, WATERMARK));
275+
verify(transaction, times(0)).buffer(any(Mutation.class));
249276
}
250277

251278
@Test

0 commit comments

Comments
 (0)