Skip to content

Commit 294bb44

Browse files
committed
Only update the partition watermark iff the partition watermark in metadata table is smaller than the given watermark.
1 parent 08c96f2 commit 294bb44

1 file changed

Lines changed: 13 additions & 3 deletions

File tree

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

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

Lines changed: 13 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -528,14 +528,24 @@ public Void updateToFinished(String partitionToken) {
528528
}
529529

530530
/**
531-
* Update the partition watermark to the given timestamp.
531+
* Update the partition watermark to the given timestamp iff the partition
532+
* watermark in metadata table is smaller than the given watermark.
532533
*
533534
* @param partitionToken the partition unique identifier
534-
* @param watermark the new partition watermark
535+
* @param watermark the new partition watermark
535536
* @return the commit timestamp of the read / write transaction
536537
*/
537538
public Void updateWatermark(String partitionToken, Timestamp watermark) {
538-
transaction.buffer(createUpdateMetadataWatermarkMutationFrom(partitionToken, watermark));
539+
Struct row = transaction.readRow(metadataTableName, KeySet.singleKey(partitionToken),
540+
Collections.singleton(COLUMN_WATERMARK));
541+
if (row == null) {
542+
LOG.error("[{}] Failed to read Watermark column", partitionToken);
543+
return null;
544+
}
545+
Timestamp partitionWatermark = row.getTimestamp(COLUMN_WATERMARK);
546+
if (partitionWatermark.compareTo(watermark) < 0) {
547+
transaction.buffer(createUpdateMetadataWatermarkMutationFrom(partitionToken, watermark));
548+
}
539549
return null;
540550
}
541551

0 commit comments

Comments
 (0)