From e274da41468d7d35840ed161243a453a3e9e19e2 Mon Sep 17 00:00:00 2001 From: Andrew Crites Date: Thu, 3 Sep 2026 10:44:23 -0700 Subject: [PATCH 1/2] Adds drain support to Spanner CDC source. Overriding truncateRestriction on SDFs allows terminating immediately and changing to DROP TABLE IF EXISTS makes deletion idempotent. --- .../changestreams/dao/PartitionMetadataAdminDao.java | 4 ++-- .../changestreams/dofn/DetectNewPartitionsDoFn.java | 7 +++++++ .../changestreams/dofn/ReadChangeStreamPartitionDoFn.java | 7 +++++++ 3 files changed, 16 insertions(+), 2 deletions(-) diff --git a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java index 80bf178f49a9..5bf4436de6bd 100644 --- a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java +++ b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java @@ -259,10 +259,10 @@ public void deletePartitionMetadataTable(List indexes) { List ddl = new ArrayList<>(); if (this.isPostgres()) { indexes.forEach(index -> ddl.add("DROP INDEX \"" + index + "\"")); - ddl.add("DROP TABLE \"" + names.getTableName() + "\""); + ddl.add("DROP TABLE IF EXISTS \"" + names.getTableName() + "\""); } else { indexes.forEach(index -> ddl.add("DROP INDEX " + index)); - ddl.add("DROP TABLE " + names.getTableName()); + ddl.add("DROP TABLE IF EXISTS " + names.getTableName()); } OperationFuture op = databaseAdminClient.updateDatabaseDdl(instanceId, databaseId, ddl, null); diff --git a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java index 4c46307aa661..169526b3624a 100644 --- a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java +++ b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/DetectNewPartitionsDoFn.java @@ -37,7 +37,9 @@ import org.apache.beam.sdk.transforms.DoFn.UnboundedPerElement; import org.apache.beam.sdk.transforms.splittabledofn.ManualWatermarkEstimator; import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker; +import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker.TruncateResult; import org.apache.beam.sdk.transforms.splittabledofn.WatermarkEstimators.Manual; +import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Duration; import org.joda.time.Instant; import org.slf4j.Logger; @@ -121,6 +123,11 @@ public TimestampRange initialRestriction(@Element PartitionMetadata partition) { TimestampUtils.previous(createdAt), com.google.cloud.Timestamp.MAX_VALUE); } + @TruncateRestriction + public @Nullable TruncateResult truncateRestriction() { + return null; + } + @GetSize public double getSize(@Restriction TimestampRange restriction) { if (!averagePartitionBytesSizeSet) { diff --git a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java index b37d1ab8b7da..3d436445c1f4 100644 --- a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java +++ b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dofn/ReadChangeStreamPartitionDoFn.java @@ -48,7 +48,9 @@ import org.apache.beam.sdk.transforms.DoFn.UnboundedPerElement; import org.apache.beam.sdk.transforms.splittabledofn.ManualWatermarkEstimator; import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker; +import org.apache.beam.sdk.transforms.splittabledofn.RestrictionTracker.TruncateResult; import org.apache.beam.sdk.transforms.splittabledofn.WatermarkEstimators.Manual; +import org.checkerframework.checker.nullness.qual.Nullable; import org.joda.time.Duration; import org.joda.time.Instant; import org.slf4j.Logger; @@ -169,6 +171,11 @@ public TimestampRange initialRestriction(@Element PartitionMetadata partition) { return TimestampRange.of(startTimestamp, endTimestamp); } + @TruncateRestriction + public @Nullable TruncateResult truncateRestriction() { + return null; + } + @GetSize public double getSize(@Element PartitionMetadata partition, @Restriction TimestampRange range) throws Exception { From 12f2b7b9e16db1ec398da09afaefa3210adb0cee Mon Sep 17 00:00:00 2001 From: Andrew Crites Date: Thu, 3 Sep 2026 14:16:51 -0700 Subject: [PATCH 2/2] Adds IF EXISTS to DROP INDEX calls and updates tests. --- .../dao/PartitionMetadataAdminDao.java | 4 ++-- .../dao/PartitionMetadataAdminDaoTest.java | 16 ++++++++-------- 2 files changed, 10 insertions(+), 10 deletions(-) diff --git a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java index 5bf4436de6bd..609d7394deb4 100644 --- a/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java +++ b/sdks/java/io/google-cloud-platform/src/main/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDao.java @@ -258,10 +258,10 @@ public void createPartitionMetadataTable() { public void deletePartitionMetadataTable(List indexes) { List ddl = new ArrayList<>(); if (this.isPostgres()) { - indexes.forEach(index -> ddl.add("DROP INDEX \"" + index + "\"")); + indexes.forEach(index -> ddl.add("DROP INDEX IF EXISTS \"" + index + "\"")); ddl.add("DROP TABLE IF EXISTS \"" + names.getTableName() + "\""); } else { - indexes.forEach(index -> ddl.add("DROP INDEX " + index)); + indexes.forEach(index -> ddl.add("DROP INDEX IF EXISTS " + index)); ddl.add("DROP TABLE IF EXISTS " + names.getTableName()); } OperationFuture op = diff --git a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDaoTest.java b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDaoTest.java index 02b9d111583b..61904ae92b1e 100644 --- a/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDaoTest.java +++ b/sdks/java/io/google-cloud-platform/src/test/java/org/apache/beam/sdk/io/gcp/spanner/changestreams/dao/PartitionMetadataAdminDaoTest.java @@ -145,9 +145,9 @@ public void testDeletePartitionMetadataTable() throws Exception { .updateDatabaseDdl(eq(INSTANCE_ID), eq(DATABASE_ID), statements.capture(), isNull()); assertEquals(3, ((Collection) statements.getValue()).size()); Iterator it = statements.getValue().iterator(); - assertTrue(it.next().contains("DROP INDEX")); - assertTrue(it.next().contains("DROP INDEX")); - assertTrue(it.next().contains("DROP TABLE")); + assertTrue(it.next().contains("DROP INDEX IF EXISTS")); + assertTrue(it.next().contains("DROP INDEX IF EXISTS")); + assertTrue(it.next().contains("DROP TABLE IF EXISTS")); } @Test @@ -158,7 +158,7 @@ public void testDeletePartitionMetadataTableWithNoIndexes() throws Exception { .updateDatabaseDdl(eq(INSTANCE_ID), eq(DATABASE_ID), statements.capture(), isNull()); assertEquals(1, ((Collection) statements.getValue()).size()); Iterator it = statements.getValue().iterator(); - assertTrue(it.next().contains("DROP TABLE")); + assertTrue(it.next().contains("DROP TABLE IF EXISTS")); } @Test @@ -170,9 +170,9 @@ public void testDeletePartitionMetadataTablePostgres() throws Exception { .updateDatabaseDdl(eq(INSTANCE_ID), eq(DATABASE_ID), statements.capture(), isNull()); assertEquals(3, ((Collection) statements.getValue()).size()); Iterator it = statements.getValue().iterator(); - assertTrue(it.next().contains("DROP INDEX \"")); - assertTrue(it.next().contains("DROP INDEX \"")); - assertTrue(it.next().contains("DROP TABLE \"")); + assertTrue(it.next().contains("DROP INDEX IF EXISTS \"")); + assertTrue(it.next().contains("DROP INDEX IF EXISTS \"")); + assertTrue(it.next().contains("DROP TABLE IF EXISTS \"")); } @Test @@ -183,7 +183,7 @@ public void testDeletePartitionMetadataTablePostgresWithNoIndexes() throws Excep .updateDatabaseDdl(eq(INSTANCE_ID), eq(DATABASE_ID), statements.capture(), isNull()); assertEquals(1, ((Collection) statements.getValue()).size()); Iterator it = statements.getValue().iterator(); - assertTrue(it.next().contains("DROP TABLE \"")); + assertTrue(it.next().contains("DROP TABLE IF EXISTS \"")); } @Test