Skip to content

Commit 86a9ef0

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

1 file changed

Lines changed: 3 additions & 0 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: 3 additions & 0 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.KeySet;
3536
import com.google.cloud.spanner.Mutation;
3637
import com.google.cloud.spanner.Options;
3738
import com.google.cloud.spanner.ResultSet;
@@ -53,6 +54,8 @@
5354
import org.apache.beam.sdk.io.gcp.spanner.changestreams.model.PartitionMetadata.State;
5455
import org.apache.beam.vendor.guava.v32_1_2_jre.com.google.common.collect.ImmutableList;
5556
import org.slf4j.Logger;
57+
58+
5659
import org.slf4j.LoggerFactory;
5760

5861
/** Data access object for the Connector metadata tables. */

0 commit comments

Comments
 (0)