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..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,11 +258,11 @@ public void createPartitionMetadataTable() { 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() + "\""); + 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)); - ddl.add("DROP TABLE " + names.getTableName()); + indexes.forEach(index -> ddl.add("DROP INDEX IF EXISTS " + index)); + 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 { 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