This is an automated email from the ASF dual-hosted git repository.
jyothsnakonisa pushed a commit to branch trunk
in repository https://gitbox.apache.org/repos/asf/cassandra-analytics.git
The following commit(s) were added to refs/heads/trunk by this push:
new 839020f6 CASSANALYTICS-182 : CDC batch-write mixing a CDC-enabled and
CDC-disabled table drops the CDC table's mutation (#224)
839020f6 is described below
commit 839020f69311864c7f0a1407898da8bd4aabc239
Author: Jyothsna konisa <[email protected]>
AuthorDate: Thu Jul 23 15:12:46 2026 -0700
CASSANALYTICS-182 : CDC batch-write mixing a CDC-enabled and CDC-disabled
table drops the CDC table's mutation (#224)
Patch by Jyothsna Konisa; Reviewed by Josh McKenzie for CASSANALYTICS-182
---
CHANGES.txt | 1 +
.../cdc/schemastore/CachingSchemaStore.java | 22 ++-
.../apache/cassandra/cdc/sidecar/SidecarCdc.java | 12 +-
.../cassandra/cdc/sidecar/SidecarCdcBuilder.java | 8 +
.../main/java/org/apache/cassandra/cdc/Cdc.java | 31 +++-
.../apache/cassandra/cdc/api/SchemaSupplier.java | 25 ++-
.../java/org/apache/cassandra/cdc/CdcTests.java | 184 ++++++++++++++++++++-
.../org/apache/cassandra/spark/data/CqlTable.java | 33 +++-
.../cassandra/spark/KryoSerializationTests.java | 30 ++++
.../org/apache/cassandra/bridge/CdcBridge.java | 20 +++
.../cassandra/spark/utils/test/TestSchema.java | 3 +-
.../cassandra/bridge/CassandraSchemaTests.java | 2 +
.../bridge/AbstractCdcBridgeImplementation.java | 6 +
.../cassandra/bridge/CassandraSchemaTests.java | 49 ++++++
.../apache/cassandra/bridge/CassandraSchema.java | 165 ++++++++++++++----
.../cassandra/spark/reader/SchemaBuilder.java | 7 +-
16 files changed, 535 insertions(+), 63 deletions(-)
diff --git a/CHANGES.txt b/CHANGES.txt
index 811417a7..afd381b8 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
0.5.0
-----
+ * CDC batch-write mixing a CDC-enabled and CDC-disabled table drops the CDC
table's mutation (CASSANALYTICS-182)
* CdcState.ReplicaCountSerializer map-size overflow corrupts persisted CDC
state (CASSANALYTICS-184)
* SSTable-version-based bridge determination (CASSANALYTICS-24)
* Upgrade sidecar version to 0.4.0 (CASSANALYTICS-176)
diff --git
a/cassandra-analytics-cdc-codec/src/main/java/org/apache/cassandra/cdc/schemastore/CachingSchemaStore.java
b/cassandra-analytics-cdc-codec/src/main/java/org/apache/cassandra/cdc/schemastore/CachingSchemaStore.java
index 94981908..35251619 100644
---
a/cassandra-analytics-cdc-codec/src/main/java/org/apache/cassandra/cdc/schemastore/CachingSchemaStore.java
+++
b/cassandra-analytics-cdc-codec/src/main/java/org/apache/cassandra/cdc/schemastore/CachingSchemaStore.java
@@ -83,8 +83,8 @@ public class CachingSchemaStore implements SchemaStore
public void initialize()
{
LOGGER.info("Initializing CachingSchemaStore");
- schemaSupplier.getCdcEnabledTables()
- .thenAccept(refreshedCdcTables -> {
+ schemaSupplier.getTables()
+ .thenAccept(ignored -> {
loadPublisher();
publishSchemas();
LOGGER.info("CachingSchemaStore initialized");
@@ -106,7 +106,7 @@ public class CachingSchemaStore implements SchemaStore
*/
public void onSchemaChange()
{
- schemaSupplier.getCdcEnabledTables().thenAccept(refreshedCdcTables -> {
+ schemaSupplier.getCDCEnabledTables().thenAccept(refreshedCdcTables -> {
for (CqlTable cqlTable : refreshedCdcTables)
{
TableIdentifier tableIdentifier =
TableIdentifier.of(cqlTable.keyspace(), cqlTable.table());
@@ -118,11 +118,19 @@ public class CachingSchemaStore implements SchemaStore
}
return value;
});
+ }
+ // We call publishSchemas() out here, after the loop, because it
re-fetches and
+ // republishes every CDC table's schema in one pass — calling it
inside the loop
+ // would redundantly re-fetch and republish every table N times
over.
+ if (!refreshedCdcTables.isEmpty())
+ {
publishSchemas();
}
- // Remove any old schema entries for deleted tables, this
operation can be done in the end as this is
- // only for removing stale entries and no one is going to use
these entries once the table is removed.
- // This doesn't have to be an atomic operation.
+ else
+ {
+ LOGGER.warn("No CDC-enabled tables found; no schemas will be
published until CDC is enabled on a table");
+ }
+ // Remove any old schema entries for deleted tables
List<TableIdentifier> refreshedTableIds = refreshedCdcTables
.stream()
.map(cqlTable ->
TableIdentifier.of(cqlTable.keyspace(), cqlTable.table()))
@@ -160,7 +168,7 @@ public class CachingSchemaStore implements SchemaStore
private void publishSchemas()
{
schemaSupplier
- .getCdcEnabledTables()
+ .getCDCEnabledTables()
.thenAccept(refreshedCdcTables -> {
for (CqlTable cqlTable : refreshedCdcTables)
{
diff --git
a/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdc.java
b/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdc.java
index f7e0513c..14bff329 100644
---
a/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdc.java
+++
b/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdc.java
@@ -62,7 +62,7 @@ public class SidecarCdc extends Cdc
* @param cdcOptions CDC processing options
* @param clusterConfigProvider provider for cluster configuration (e.g.
datacenter, hosts)
* @param eventConsumer consumer that receives CDC change events
- * @param schemaSupplier supplier for CDC-enabled table schemas
+ * @param schemaSupplier supplier for all table schemas
(CDC-enabled and disabled); see {@link SchemaSupplier}
* @param tokenRangeSupplier supplier for the token ranges assigned to
this partition
* @param sidecarCdcClient externally managed Sidecar HTTP client;
<em>not</em> closed by
* {@code SidecarCdc} or {@code
SidecarCdcBuilder}
@@ -92,11 +92,11 @@ public class SidecarCdc extends Cdc
public void initSchema()
{
- Set<CqlTable> tables =
FutureUtils.get(schemaSupplier.getCdcEnabledTables());
- Optional<ReplicationFactor> rfOp = tables.stream()
-
.map(CqlTable::replicationFactor)
- .filter(rf ->
rf.getOptions().containsKey(dc()))
-
.max(Comparator.comparingInt(rf -> rf.getOptions().get(dc())));
+ Set<CqlTable> cdcEnabledTables =
FutureUtils.get(schemaSupplier.getCDCEnabledTables());
+ Optional<ReplicationFactor> rfOp = cdcEnabledTables.stream()
+
.map(CqlTable::replicationFactor)
+ .filter(rf ->
rf.getOptions().containsKey(dc()))
+
.max(Comparator.comparingInt(rf -> rf.getOptions().get(dc())));
if (!rfOp.isPresent())
{
diff --git
a/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
b/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
index 4027db2d..559edc2f 100644
---
a/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
+++
b/cassandra-analytics-cdc-sidecar/src/main/java/org/apache/cassandra/cdc/sidecar/SidecarCdcBuilder.java
@@ -25,6 +25,7 @@ import org.apache.cassandra.cdc.CdcBuilder;
import org.apache.cassandra.cdc.api.CdcOptions;
import org.apache.cassandra.cdc.api.EventConsumer;
import org.apache.cassandra.cdc.api.SchemaSupplier;
+import org.apache.cassandra.cdc.api.TableIdLookup;
import org.apache.cassandra.cdc.api.TokenRangeSupplier;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.spark.utils.AsyncExecutor;
@@ -102,6 +103,13 @@ public class SidecarCdcBuilder extends CdcBuilder
return withSidecarCdcCassandraClient(cassandraClient); // rebuild
SidecarStatePersister with new AsyncExecutor
}
+ @Override
+ public SidecarCdcBuilder withTableIdLookup(@NotNull TableIdLookup
tableIdLookup)
+ {
+ super.withTableIdLookup(tableIdLookup);
+ return this;
+ }
+
public SidecarCdcBuilder
withSidecarCdcCassandraClient(SidecarCdcCassandraClient cassandraClient)
{
this.cassandraClient = cassandraClient;
diff --git
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java
index b3270c9f..57484324 100644
--- a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java
+++ b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/Cdc.java
@@ -415,22 +415,39 @@ public class Cdc implements Closeable
try
{
schemaSupplier
- .getCdcEnabledTables()
- .handle((tables, throwable) -> {
+ .getTables()
+ .handle((allTables, throwable) -> {
if (throwable != null)
{
LOGGER.warn("Error refreshing schema", throwable);
return null;
}
- this.cdcEnabledTables = tables;
- if (tables == null || tables.isEmpty())
+ if (allTables == null || allTables.isEmpty())
{
- LOGGER.warn("No CQL enabled tables");
+ // A brand-new cluster with no user tables yet is expected,
+ // not a failure — refreshSchema() retries on its next
scheduled run.
+ LOGGER.warn("No tables returned from schema supplier; will
retry on next scheduled schema refresh");
return null;
}
- // update Schema instance with latest schema
- cdcBridge().updateCdcSchema(tables, cdcOptions.partitioner(),
tableIdLookup);
+ // Filter CDC-enabled tables for publishing decisions
+ Set<CqlTable> cdcTables = allTables.stream()
+ .filter(CqlTable::cdc)
+
.collect(Collectors.toSet());
+ this.cdcEnabledTables = cdcTables;
+ if (cdcTables.isEmpty())
+ {
+ // Schema.instance is still updated below with allTables
(which is non-empty
+ // here), so this doesn't fast-fail — it just means no CDC
updates will be
+ // processed until CDC is enabled on at least one table's
schema.
+ LOGGER.warn("No CDC-enabled tables found; no cdc updates
will be processed "
+ + "until CDC is enabled on a table's schema");
+ }
+
+ // Update Schema.instance with ALL tables so deserialization
never throws
+ // UnknownTableException for non-CDC tables co-located in
batch mutations.
+ // Each table's CDC flag is set correctly from CqlTable.cdc().
+ cdcBridge().updateCdcSchema(allTables,
cdcOptions.partitioner(), tableIdLookup);
return null;
})
.whenComplete((aVoid, throwable) -> {
diff --git
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/api/SchemaSupplier.java
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/api/SchemaSupplier.java
index 0bd55a2f..0d86adea 100644
---
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/api/SchemaSupplier.java
+++
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/api/SchemaSupplier.java
@@ -21,13 +21,34 @@ package org.apache.cassandra.cdc.api;
import java.util.Set;
import java.util.concurrent.CompletableFuture;
+import java.util.stream.Collectors;
import org.apache.cassandra.spark.data.CqlTable;
/**
- * Supplies all CDC enabled tables
+ * Supplies schema for tables relevant to CDC processing.
*/
public interface SchemaSupplier
{
- CompletableFuture<Set<CqlTable>> getCdcEnabledTables();
+ /**
+ * @return ALL tables (CDC-enabled and CDC-disabled) so that the bridge's
Schema.instance
+ * is complete enough to deserialize any commit log mutation without
UnknownTableException.
+ * Callers use {@link org.apache.cassandra.spark.data.CqlTable#cdc()} to
filter for publishing.
+ */
+ CompletableFuture<Set<CqlTable>> getTables();
+
+ /**
+ * @return the subset of {@link #getTables()} that are CDC-enabled — i.e.
what to actually
+ * publish/process, as opposed to the full set needed for schema
completeness. A default
+ * method (rather than requiring implementations to filter themselves) so
every caller that
+ * only cares about CDC-enabled tables shares one implementation of the
+ * {@code getTables().filter(CqlTable::cdc)} pattern, instead of repeating
it at each call
+ * site.
+ */
+ default CompletableFuture<Set<CqlTable>> getCDCEnabledTables()
+ {
+ return getTables().thenApply(tables -> tables.stream()
+ .filter(CqlTable::cdc)
+
.collect(Collectors.toSet()));
+ }
}
diff --git
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
index b67fc2fd..b4490cb0 100644
---
a/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
+++
b/cassandra-analytics-cdc/src/test/java/org/apache/cassandra/cdc/CdcTests.java
@@ -25,6 +25,7 @@ import java.nio.file.Files;
import java.nio.file.Path;
import java.util.ArrayList;
import java.util.Arrays;
+import java.util.Collection;
import java.util.Collections;
import java.util.HashMap;
import java.util.HashSet;
@@ -56,17 +57,26 @@ import org.slf4j.LoggerFactory;
import org.apache.cassandra.bridge.CassandraVersion;
import org.apache.cassandra.bridge.CdcBridge;
import org.apache.cassandra.bridge.TokenRange;
+import org.apache.cassandra.cdc.api.CassandraSource;
import org.apache.cassandra.cdc.api.CommitLog;
+import org.apache.cassandra.cdc.api.CommitLogInstance;
+import org.apache.cassandra.cdc.api.CommitLogMarkers;
import org.apache.cassandra.cdc.api.CommitLogProvider;
+import org.apache.cassandra.cdc.api.CommitLogReader;
import org.apache.cassandra.cdc.api.EventConsumer;
import org.apache.cassandra.cdc.api.Marker;
+import org.apache.cassandra.cdc.api.Row;
import org.apache.cassandra.cdc.api.SchemaSupplier;
import org.apache.cassandra.cdc.api.StatePersister;
+import org.apache.cassandra.cdc.api.TableIdLookup;
import org.apache.cassandra.cdc.msg.CdcEvent;
import org.apache.cassandra.cdc.msg.Value;
+import org.apache.cassandra.cdc.scanner.CdcStreamScanner;
import org.apache.cassandra.cdc.state.CdcState;
+import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.cdc.test.CdcTestBase;
import org.apache.cassandra.cdc.test.CdcTester;
+import org.apache.cassandra.db.commitlog.PartitionUpdateWrapper;
import org.apache.cassandra.db.marshal.ByteBufferAccessor;
import org.apache.cassandra.serializers.CollectionSerializer;
import org.apache.cassandra.spark.data.CqlField;
@@ -77,6 +87,8 @@ import
org.apache.cassandra.spark.data.partitioner.Partitioner;
import org.apache.cassandra.spark.utils.AsyncExecutor;
import org.apache.cassandra.spark.utils.ByteBufferUtils;
import org.apache.cassandra.spark.utils.IOUtils;
+import org.apache.cassandra.spark.utils.TableIdentifier;
+import org.apache.cassandra.spark.utils.TimeProvider;
import org.apache.cassandra.spark.utils.TimeUtils;
import org.apache.cassandra.spark.utils.test.TestSchema;
import org.apache.cassandra.transport.ProtocolVersion;
@@ -166,7 +178,12 @@ public class CdcTests extends CdcTestBase
final int maxRows = 5000;
final int batchSize = 500;
final int numBatches = maxRows / batchSize;
- TestSchema testSchema = TestSchema.basicBuilder(bridge).build();
+ // withCdc(true) is required so table.cdc() is true —
SchemaSupplier.getTables() now
+ // returns all tables, and Cdc.refreshSchema() filters to
CDC-enabled ones via
+ // CqlTable#cdc(), not via any pre-filtering by the supplier
itself. Without this,
+ // cdcEnabledTables ends up empty, keyspaceSupplier() watches no
keyspaces, and the
+ // CDC scanner never sees the mutations this test writes.
+ TestSchema testSchema =
TestSchema.basicBuilder(bridge).withCdc(true).build();
CqlTable table = testSchema.buildTable();
bridge.buildSchema(table.createStatement(),
@@ -266,6 +283,171 @@ public class CdcTests extends CdcTestBase
}
}
+ @ParameterizedTest
+
@MethodSource("org.apache.cassandra.cdc.test.TestVersionSupplier#testVersions")
+ public void
testRefreshSchemaUpdatesBridgeWithAllTablesButFiltersCdcEnabledTables(CassandraVersion
version)
+ {
+ // Regression test for the CDC batch-write bug: refreshSchema() must
register EVERY
+ // table returned by the schema supplier (CDC-enabled and
CDC-disabled) with the
+ // bridge's Schema.instance — not just the CDC-enabled ones — so that
a commit log
+ // Mutation spanning both a CDC-enabled and a CDC-disabled table in
the same keyspace
+ // (e.g. a BEGIN BATCH statement) can still be deserialized.
Separately, only the
+ // CDC-enabled tables should end up in cdcEnabledTables, since that
set drives
+ // publishing/replica-check decisions, not schema completeness.
+ TestSchema cdcSchema = TestSchema.builder(bridge)
+ .withPartitionKey("pk", bridge.uuid())
+ .withColumn("val", bridge.text())
+ .withCdc(true)
+ .build();
+ CqlTable cdcTable = cdcSchema.buildTable();
+
+ TestSchema nonCdcSchema = TestSchema.builder(bridge)
+ .withKeyspace(cdcSchema.keyspace)
+ .withPartitionKey("pk",
bridge.uuid())
+ .withColumn("val", bridge.text())
+ .withCdc(false)
+ .build();
+ CqlTable nonCdcTable = nonCdcSchema.buildTable();
+
+ AtomicReference<Set<CqlTable>> tablesPassedToBridge = new
AtomicReference<>();
+ SchemaSupplier schemaSupplier = () ->
CompletableFuture.completedFuture(ImmutableSet.of(cdcTable, nonCdcTable));
+
+ try (RecordingCdc cdc = new
RecordingCdc(Cdc.builder("test-refresh-schema", 0, event -> { }, schemaSupplier)
+
.withExecutor(CdcTests.ASYNC_EXECUTOR)
+
.withTableIdLookup(cdcBridge.internalTableIdLookup())
+
.withCommitLogProvider(tokenRange -> Stream.empty())
+
.withCdcOptions(cdcOptions),
+ tablesPassedToBridge))
+ {
+ // start() flips isRunning to true and calls refreshSchema()
synchronously (the
+ // schemaSupplier future is already completed, so the
.handle()/.whenComplete()
+ // chain runs inline). scheduleRun()/scheduleMonitorSchema() are
overridden to
+ // no-ops in RecordingCdc, so this only exercises the schema
refresh path.
+ cdc.start();
+
+ assertThat(tablesPassedToBridge.get())
+ .as("Schema.instance must be updated with ALL tables (CDC and
non-CDC) so a "
+ + "batch mutation spanning both can still deserialize")
+ .containsExactlyInAnyOrder(cdcTable, nonCdcTable);
+
+ assertThat(cdc.cdcEnabledTables)
+ .as("cdcEnabledTables (used for publishing/replica-check
decisions) must contain "
+ + "only the CDC-enabled table, even though Schema.instance
was given both")
+ .containsExactly(cdcTable);
+ }
+ }
+
+ /**
+ * Test-only {@link Cdc} subclass that intercepts {@link #cdcBridge()} to
record exactly
+ * what table set gets passed to {@code CdcBridge.updateCdcSchema(...)}
during
+ * {@code refreshSchema()}, while still delegating to the real bridge so
Schema.instance is
+ * genuinely updated (not just observed).
scheduleRun()/scheduleMonitorSchema() are
+ * overridden to no-ops so start() only exercises the schema refresh path,
without kicking
+ * off an unrelated micro-batch read or a recurring schema-refresh loop.
+ */
+ private static final class RecordingCdc extends Cdc
+ {
+ private final AtomicReference<Set<CqlTable>> tablesPassedToBridge;
+
+ RecordingCdc(CdcBuilder builder, AtomicReference<Set<CqlTable>>
tablesPassedToBridge)
+ {
+ super(builder);
+ this.tablesPassedToBridge = tablesPassedToBridge;
+ }
+
+ @Override
+ protected void scheduleRun(long delayMillis)
+ {
+ // no-op — this test only exercises schema refresh, not the
micro-batch read loop
+ }
+
+ @Override
+ public void scheduleMonitorSchema()
+ {
+ // no-op — avoid recursively rescheduling refreshSchema() after
start() calls it
+ }
+
+ @Override
+ protected CdcBridge cdcBridge()
+ {
+ CdcBridge real = super.cdcBridge();
+ return new RecordingCdcBridge(real, tablesPassedToBridge);
+ }
+ }
+
+ /**
+ * Delegates every call to the real {@link CdcBridge}, except {@code
updateCdcSchema}, whose
+ * argument is captured before delegating.
+ */
+ private static final class RecordingCdcBridge extends CdcBridge
+ {
+ private final CdcBridge delegate;
+ private final AtomicReference<Set<CqlTable>> tablesPassedToBridge;
+
+ RecordingCdcBridge(CdcBridge delegate, AtomicReference<Set<CqlTable>>
tablesPassedToBridge)
+ {
+ this.delegate = delegate;
+ this.tablesPassedToBridge = tablesPassedToBridge;
+ }
+
+ @Override
+ public void updateCdcSchema(@NotNull Set<CqlTable> cdcTables,
+ @NotNull Partitioner partitioner,
+ @NotNull TableIdLookup tableIdLookup)
+ {
+ tablesPassedToBridge.set(cdcTables);
+ delegate.updateCdcSchema(cdcTables, partitioner, tableIdLookup);
+ }
+
+ @Override
+ public void unregisterNonCdcTables(@NotNull Set<TableIdentifier>
tables)
+ {
+ delegate.unregisterNonCdcTables(tables);
+ }
+
+ @Override
+ public CommitLogInstance createCommitLogInstance(Path path)
+ {
+ return delegate.createCommitLogInstance(path);
+ }
+
+ @Override
+ public TableIdLookup internalTableIdLookup()
+ {
+ return delegate.internalTableIdLookup();
+ }
+
+ @Override
+ public CommitLogReader.Result readLog(@NotNull CommitLog log,
+ @Nullable TokenRange tokenRange,
+ @NotNull CommitLogMarkers
markers,
+ int partitionId,
+ @NotNull ICdcStats stats,
+ @Nullable AsyncExecutor executor,
+ @Nullable Consumer<Marker>
listener,
+ @Nullable Long
startTimestampMicros,
+ boolean readCommitLogHeader)
+ {
+ return delegate.readLog(log, tokenRange, markers, partitionId,
stats, executor, listener, startTimestampMicros, readCommitLogHeader);
+ }
+
+ @Override
+ public CdcStreamScanner
openCdcStreamScanner(Collection<PartitionUpdateWrapper> updates,
+ @NotNull CdcState
endState,
+ Random random,
+ CassandraSource
cassandraSource,
+ double traceSampleRate)
+ {
+ return delegate.openCdcStreamScanner(updates, endState, random,
cassandraSource, traceSampleRate);
+ }
+
+ @Override
+ public void log(TimeProvider timeProvider, CqlTable cqlTable,
CommitLogInstance log, Row row, long timestamp)
+ {
+ delegate.log(cqlTable, log, row, timestamp);
+ }
+ }
+
@ParameterizedTest
@MethodSource("org.apache.cassandra.cdc.test.TestVersionSupplier#testVersions")
public void testSinglePartitionKey(CassandraVersion version)
diff --git
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlTable.java
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlTable.java
index d71e471c..73ed6b22 100644
---
a/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlTable.java
+++
b/cassandra-analytics-common/src/main/java/org/apache/cassandra/spark/data/CqlTable.java
@@ -49,6 +49,7 @@ public class CqlTable implements Serializable
private final String createStatement;
private final List<CqlField> fields;
private final Set<CqlField.CqlUdt> udts;
+ private final boolean cdc;
private final Map<String, CqlField> fieldsMap;
private final Set<String> columnsWithUdts;
@@ -75,11 +76,24 @@ public class CqlTable implements Serializable
@NotNull List<CqlField> fields,
@NotNull Set<CqlField.CqlUdt> udts,
int indexCount)
+ {
+ this(keyspace, table, createStatement, replicationFactor, fields,
udts, indexCount, false);
+ }
+
+ public CqlTable(@NotNull String keyspace,
+ @NotNull String table,
+ @NotNull String createStatement,
+ @NotNull ReplicationFactor replicationFactor,
+ @NotNull List<CqlField> fields,
+ @NotNull Set<CqlField.CqlUdt> udts,
+ int indexCount,
+ boolean cdc)
{
this.keyspace = keyspace;
this.table = table;
this.createStatement = createStatement;
this.replicationFactor = replicationFactor;
+ this.cdc = cdc;
this.fields = fields.stream().sorted().collect(Collectors.toList());
this.fieldsMap =
this.fields.stream().collect(Collectors.toMap(CqlField::name,
Function.identity()));
this.partitionKeys =
this.fields.stream().filter(CqlField::isPartitionKey).sorted().collect(Collectors.toList());
@@ -241,6 +255,16 @@ public class CqlTable implements Serializable
return createStatement;
}
+ /**
+ * Returns true if CDC is enabled on this table.
+ * This flag is explicitly set at construction time (e.g. from the CREATE
TABLE statement)
+ * and stored as a field — it is not derived from createStatement at
runtime.
+ */
+ public boolean cdc()
+ {
+ return cdc;
+ }
+
public int indexCount()
{
return indexCount;
@@ -327,7 +351,7 @@ public class CqlTable implements Serializable
@Override
public int hashCode()
{
- return Objects.hash(keyspace, table, createStatement, fields, udts);
+ return Objects.hash(keyspace, table, createStatement, fields, udts,
cdc);
}
@Override
@@ -351,7 +375,8 @@ public class CqlTable implements Serializable
&& Objects.equals(this.table, that.table)
&& Objects.equals(this.createStatement, that.createStatement)
&& Objects.equals(this.fields, that.fields)
- && Objects.equals(this.udts, that.udts);
+ && Objects.equals(this.udts, that.udts)
+ && this.cdc == that.cdc;
}
public static class Serializer extends
com.esotericsoftware.kryo.Serializer<CqlTable>
@@ -383,7 +408,8 @@ public class CqlTable implements Serializable
udts.add((CqlField.CqlUdt) CqlField.CqlType.read(input,
cassandraTypes));
}
int indexCount = input.readInt();
- return new CqlTable(keyspace, table, createStatement,
replicationFactor, fields, udts, indexCount);
+ boolean cdc = input.readBoolean();
+ return new CqlTable(keyspace, table, createStatement,
replicationFactor, fields, udts, indexCount, cdc);
}
@Override
@@ -406,6 +432,7 @@ public class CqlTable implements Serializable
udt.write(output);
}
output.writeInt(table.indexCount());
+ output.writeBoolean(table.cdc());
}
}
}
diff --git
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
index 62e76774..23fd7874 100644
---
a/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
+++
b/cassandra-analytics-core/src/test/java/org/apache/cassandra/spark/KryoSerializationTests.java
@@ -21,6 +21,7 @@ package org.apache.cassandra.spark;
import java.math.BigInteger;
import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
import java.util.Map;
@@ -269,6 +270,35 @@ public class KryoSerializationTests
assertThat(deserialized).isEqualTo(table);
}
+ @ParameterizedTest
+ @MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
+ public void testCqlTableCdcFlagSurvivesSerialization(CassandraBridge
bridge)
+ {
+ List<CqlField> fields = ImmutableList.of(new CqlField(true, false,
false, "a", bridge.bigint(), 0));
+ ReplicationFactor replicationFactor = new
ReplicationFactor(ReplicationFactor.ReplicationStrategy.NetworkTopologyStrategy,
+
ImmutableMap.of("DC1", 3, "DC2", 3));
+
+ for (boolean cdc : new boolean[]{true, false})
+ {
+ CqlTable table = new CqlTable("test_keyspace",
+ "test_table",
+ "create table
test_keyspace.test_table (a bigint, primary key(a));",
+ replicationFactor,
+ fields,
+ Collections.emptySet(),
+ 0,
+ cdc);
+
+ Output out = serialize(bridge.getVersion(), table);
+ CqlTable deserialized = deserialize(bridge.getVersion(), out,
CqlTable.class);
+ assertThat(deserialized).isNotNull();
+ assertThat(deserialized).isEqualTo(table);
+ assertThat(deserialized.cdc())
+ .as("cdc=%s must survive Kryo serialization round-trip", cdc)
+ .isEqualTo(cdc);
+ }
+ }
+
@ParameterizedTest
@MethodSource("org.apache.cassandra.bridge.VersionRunner#bridges")
public void testCassandraInstance(CassandraBridge bridge)
diff --git
a/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CdcBridge.java
b/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CdcBridge.java
index d233808b..b52e247f 100644
--- a/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CdcBridge.java
+++ b/cassandra-bridge/src/main/java/org/apache/cassandra/bridge/CdcBridge.java
@@ -42,6 +42,7 @@ import
org.apache.cassandra.db.commitlog.PartitionUpdateWrapper;
import org.apache.cassandra.spark.data.CqlTable;
import org.apache.cassandra.spark.data.partitioner.Partitioner;
import org.apache.cassandra.spark.utils.AsyncExecutor;
+import org.apache.cassandra.spark.utils.TableIdentifier;
import org.apache.cassandra.spark.utils.TimeProvider;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
@@ -74,6 +75,25 @@ public abstract class CdcBridge
@NotNull Partitioner partitioner,
@NotNull TableIdLookup tableIdLookup);
+ /**
+ * Removes tables from the bridge's {@code Schema.instance} that were
previously registered
+ * (via {@link #updateCdcSchema}) but are no longer needed — e.g. a
non-CDC table that used
+ * to share partition-key structure with a CDC-enabled table in the same
keyspace, where a
+ * later schema change (the CDC table dropped, or either table's partition
key altered) means
+ * it can no longer be co-located with a CDC-enabled table's update in the
same commit-log
+ * {@code Mutation}.
+ *
+ * <p>Implementations must be idempotent (a no-op for a table that isn't
currently
+ * registered) and must never remove a table that is currently CDC-enabled
— the caller is
+ * responsible for only passing tables it has determined are safe to
unregister, but this is
+ * a last-line defense against silently dropping schema for a table CDC
still needs. (Named
+ * to make that CDC-enabled-table refusal unsurprising at the call site,
rather than
+ * implying unconditional removal of whatever is passed in.)
+ *
+ * @param tables the tables to unregister
+ */
+ public abstract void unregisterNonCdcTables(@NotNull Set<TableIdentifier>
tables);
+
public abstract CommitLogReader.Result readLog(@NotNull CommitLog log,
@Nullable TokenRange
tokenRange,
@NotNull CommitLogMarkers
markers,
diff --git
a/cassandra-bridge/src/testFixtures/java/org/apache/cassandra/spark/utils/test/TestSchema.java
b/cassandra-bridge/src/testFixtures/java/org/apache/cassandra/spark/utils/test/TestSchema.java
index e9684749..505ada88 100644
---
a/cassandra-bridge/src/testFixtures/java/org/apache/cassandra/spark/utils/test/TestSchema.java
+++
b/cassandra-bridge/src/testFixtures/java/org/apache/cassandra/spark/utils/test/TestSchema.java
@@ -498,7 +498,8 @@ public final class TestSchema
rf,
allFields,
udts,
- 0);
+ 0,
+ withCdc);
}
public void writeSSTable(TemporaryDirectory directory,
diff --git
a/cassandra-five-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
b/cassandra-five-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
index 7b3c2db8..0a5f1350 100644
---
a/cassandra-five-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
+++
b/cassandra-five-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
@@ -45,6 +45,7 @@ public class CassandraSchemaTests
.withPartitionKey("a",
BRIDGE.bigint())
.withClusteringKey("b",
BRIDGE.text())
.withColumn("c",
BRIDGE.timeuuid())
+ .withCdc(true)
.build();
final CqlTable cqlTable1 = testSchema1.buildTable();
@@ -52,6 +53,7 @@ public class CassandraSchemaTests
.withPartitionKey("pk",
BRIDGE.uuid())
.withClusteringKey("ck",
BRIDGE.aInt())
.withColumn("val",
BRIDGE.blob())
+ .withCdc(true)
.build();
final CqlTable cqlTable2 = testSchema2.buildTable();
diff --git
a/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/bridge/AbstractCdcBridgeImplementation.java
b/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/bridge/AbstractCdcBridgeImplementation.java
index 77f8f881..67e25aae 100644
---
a/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/bridge/AbstractCdcBridgeImplementation.java
+++
b/cassandra-four-zero-bridge/src/main/java/org/apache/cassandra/bridge/AbstractCdcBridgeImplementation.java
@@ -70,6 +70,7 @@ import org.apache.cassandra.spark.data.CqlType;
import org.apache.cassandra.spark.data.partitioner.Partitioner;
import org.apache.cassandra.spark.utils.AsyncExecutor;
import org.apache.cassandra.spark.utils.ByteBufferUtils;
+import org.apache.cassandra.spark.utils.TableIdentifier;
import org.apache.cassandra.spark.utils.TimeProvider;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
@@ -108,6 +109,11 @@ public abstract class AbstractCdcBridgeImplementation
extends CdcBridge
CassandraSchema.updateCdcSchema(cdcTables, partitioner, tableIdLookup);
}
+ public void unregisterNonCdcTables(@NotNull Set<TableIdentifier> tables)
+ {
+ CassandraSchema.unregisterNonCdcTables(tables);
+ }
+
public CommitLogReader.Result readLog(@NotNull CommitLog log,
@Nullable TokenRange tokenRange,
@NotNull CommitLogMarkers markers,
diff --git
a/cassandra-four-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
b/cassandra-four-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
index 7b3c2db8..7d39071e 100644
---
a/cassandra-four-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
+++
b/cassandra-four-zero-bridge/src/test/java/org/apache/cassandra/bridge/CassandraSchemaTests.java
@@ -26,6 +26,7 @@ import org.junit.jupiter.api.Test;
import org.apache.cassandra.schema.Schema;
import org.apache.cassandra.spark.data.CqlTable;
+import org.apache.cassandra.spark.utils.TableIdentifier;
import org.apache.cassandra.spark.utils.test.TestSchema;
import org.apache.cassandra.spark.data.partitioner.Partitioner;
@@ -45,6 +46,7 @@ public class CassandraSchemaTests
.withPartitionKey("a",
BRIDGE.bigint())
.withClusteringKey("b",
BRIDGE.text())
.withColumn("c",
BRIDGE.timeuuid())
+ .withCdc(true)
.build();
final CqlTable cqlTable1 = testSchema1.buildTable();
@@ -52,6 +54,7 @@ public class CassandraSchemaTests
.withPartitionKey("pk",
BRIDGE.uuid())
.withClusteringKey("ck",
BRIDGE.aInt())
.withColumn("val",
BRIDGE.blob())
+ .withCdc(true)
.build();
final CqlTable cqlTable2 = testSchema2.buildTable();
@@ -90,4 +93,50 @@ public class CassandraSchemaTests
assertThat(CassandraSchema.isCdcEnabled(schema, cqlTable1)).isFalse();
assertThat(CassandraSchema.isCdcEnabled(schema, cqlTable2)).isFalse();
}
+
+ @Test
+ public void testUnregisterNonCdcTables()
+ {
+ Schema schema = Schema.instance;
+
+ TestSchema nonCdcSchema = TestSchema.builder(BRIDGE)
+ .withPartitionKey("a",
BRIDGE.uuid())
+ .withColumn("b", BRIDGE.text())
+ .build();
+ CqlTable nonCdcTable = nonCdcSchema.buildTable();
+ TableIdentifier nonCdcId = TableIdentifier.of(nonCdcTable.keyspace(),
nonCdcTable.table());
+
+ TestSchema cdcSchema = TestSchema.builder(BRIDGE)
+ .withKeyspace(nonCdcTable.keyspace())
+ .withPartitionKey("a", BRIDGE.uuid())
+ .withColumn("b", BRIDGE.text())
+ .withCdc(true)
+ .build();
+ CqlTable cdcTable = cdcSchema.buildTable();
+ TableIdentifier cdcId = TableIdentifier.of(cdcTable.keyspace(),
cdcTable.table());
+
+ // register both tables (as if they'd been found to share
partition-key structure)
+ CassandraSchema.updateCdcSchema(schema, ImmutableSet.of(nonCdcTable,
cdcTable), Partitioner.Murmur3Partitioner, (keyspace, table) -> null);
+ assertThat(CassandraSchema.has(schema, nonCdcTable.keyspace(),
nonCdcTable.table())).isTrue();
+ assertThat(CassandraSchema.has(schema, cdcTable.keyspace(),
cdcTable.table())).isTrue();
+
+ // a later refresh determines nonCdcTable is no longer at risk —
unregister it
+ CassandraSchema.unregisterNonCdcTables(schema,
ImmutableSet.of(nonCdcId));
+ assertThat(CassandraSchema.has(schema, nonCdcTable.keyspace(),
nonCdcTable.table())).isFalse();
+ // the CDC-enabled table must be completely unaffected
+ assertThat(CassandraSchema.has(schema, cdcTable.keyspace(),
cdcTable.table())).isTrue();
+ assertThat(CassandraSchema.isCdcEnabled(schema, cdcTable)).isTrue();
+
+ // idempotent: unregistering an already-unregistered table is a no-op,
not an error
+ CassandraSchema.unregisterNonCdcTables(schema,
ImmutableSet.of(nonCdcId));
+ assertThat(CassandraSchema.has(schema, nonCdcTable.keyspace(),
nonCdcTable.table())).isFalse();
+
+ // refuses to unregister a table that is currently CDC-enabled
+ CassandraSchema.unregisterNonCdcTables(schema, ImmutableSet.of(cdcId));
+ assertThat(CassandraSchema.has(schema, cdcTable.keyspace(),
cdcTable.table())).isTrue();
+ assertThat(CassandraSchema.isCdcEnabled(schema, cdcTable)).isTrue();
+
+ // unregistering an unknown table (never registered) is a no-op, not
an error
+ CassandraSchema.unregisterNonCdcTables(schema,
ImmutableSet.of(TableIdentifier.of("unknown_ks", "unknown_table")));
+ }
}
diff --git
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraSchema.java
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraSchema.java
index bd4f2c78..bf90bbf0 100644
---
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraSchema.java
+++
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/bridge/CassandraSchema.java
@@ -20,6 +20,7 @@
package org.apache.cassandra.bridge;
import java.util.ArrayList;
+import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Objects;
@@ -48,6 +49,7 @@ import org.apache.cassandra.schema.Types;
import org.apache.cassandra.spark.data.CqlTable;
import org.apache.cassandra.spark.data.partitioner.Partitioner;
import org.apache.cassandra.spark.reader.SchemaBuilder;
+import org.apache.cassandra.spark.utils.TableIdentifier;
import org.jetbrains.annotations.NotNull;
import org.jetbrains.annotations.Nullable;
@@ -278,7 +280,7 @@ public final class CassandraSchema
return;
}
- LOGGER.info("Schema change detected, updating new table schema
keyspace={} table={}", keyspace, cqlTable.table());
+ LOGGER.info("Schema change detected, updating table schema
keyspace={} table={} cdc={}", keyspace, cqlTable.table(), enableCdc);
SchemaUpdater.updateTable(s, ks.get(), updatedTable);
});
}
@@ -296,52 +298,147 @@ public final class CassandraSchema
.collect(Collectors.joining(",")));
}
- Map<String, Set<String>> cdcEnabledTables =
CassandraSchema.cdcEnabledTables(schema);
+ Set<TableIdentifier> currentlyCdcEnabled =
currentlyCdcEnabledTables(schema);
+
for (CqlTable table : cdcTables)
{
table.udts().forEach(udt ->
CassandraTypesImplementation.INSTANCE.updateUDTs(table.keyspace(), udt));
UUID tableId = tableIdLookup.lookup(table.keyspace(),
table.table());
- if (cdcEnabledTables.containsKey(table.keyspace()) &&
cdcEnabledTables.get(table.keyspace()).contains(table.table()))
+ boolean previouslyCdcEnabled =
currentlyCdcEnabled.contains(TableIdentifier.of(table.keyspace(),
table.table()));
+ if (previouslyCdcEnabled)
{
- // table has cdc enabled already, update schema if it has
changed
- LOGGER.debug("CDC already enabled keyspace={} table={}",
table.keyspace(), table.table());
- cdcEnabledTables.get(table.keyspace()).remove(table.table());
- CassandraSchema.maybeUpdateSchema(schema, partitioner, table,
tableId, true);
-
Preconditions.checkArgument(CassandraSchema.isCdcEnabled(schema, table),
- "CDC not enabled for table: " +
table.keyspace() + "." + table.table());
- continue;
+ // maybeUpdateSchema logs on its own when it actually performs
an update.
+ CassandraSchema.maybeUpdateSchema(schema, partitioner, table,
tableId, table.cdc());
}
-
- if (CassandraSchema.has(schema, table))
+ else if (CassandraSchema.has(schema, table))
+ {
+ // table exists but wasn't tracked as cdc-enabled (e.g. a
non-CDC table the
+ // caller included in cdcTables anyway) — update if schema
changed.
+ CassandraSchema.maybeUpdateSchema(schema, partitioner, table,
tableId, table.cdc());
+ }
+ else
{
- // update schema if changed for existing table
- LOGGER.info("Enabling CDC on existing table keyspace={}
table={}", table.keyspace(), table.table());
- CassandraSchema.maybeUpdateSchema(schema, partitioner, table,
tableId, true);
-
Preconditions.checkArgument(CassandraSchema.isCdcEnabled(schema, table),
- "CDC not enabled for table: " +
table.keyspace() + "." + table.table());
- continue;
+ // new table — register with the CDC flag from the create
statement
+ LOGGER.info("Registering new table keyspace={} table={}
cdc={}", table.keyspace(), table.table(), table.cdc());
+ new SchemaBuilder(table, partitioner, tableId, table.cdc());
+ if (tableId != null && table.cdc())
+ {
+ // verify CDC-enabled tables are correctly initialized
+ TableId tableIdAfter = TableId.fromUUID(tableId);
+
Preconditions.checkNotNull(schema.getTableMetadata(tableIdAfter), "Table not
initialized in the schema");
+
Preconditions.checkArgument(Objects.requireNonNull(schema.getKeyspaceInstance(table.keyspace())).hasColumnFamilyStore(tableIdAfter),
+ "ColumnFamilyStore not
initialized in the schema");
+
Preconditions.checkArgument(CassandraSchema.isCdcEnabled(schema, table),
+ "CDC not enabled for table: "
+ table.keyspace() + "." + table.table());
+ }
}
+ }
+ disableCdcOnStaleTables(schema, currentlyCdcEnabled, cdcTables);
+ }
- // new table so initialize table with cdc = true
- LOGGER.info("Adding new CDC enabled table keyspace={} table={}",
table.keyspace(), table.table());
- new SchemaBuilder(table, partitioner, tableId, true);
- if (tableId != null)
+ private static Set<TableIdentifier> currentlyCdcEnabledTables(Schema
schema)
+ {
+ return CassandraSchema.cdcEnabledTables(schema)
+ .entrySet()
+ .stream()
+ .flatMap(e -> e.getValue().stream().map(table ->
TableIdentifier.of(e.getKey(), table)))
+ .collect(Collectors.toSet());
+ }
+
+ /**
+ * Disables CDC on every table in {@code currentlyCdcEnabled} that is not
CDC-enabled in
+ * {@code cdcTables} (dropped, or CDC disabled in its CREATE TABLE).
+ */
+ private static void disableCdcOnStaleTables(Schema schema,
Set<TableIdentifier> currentlyCdcEnabled, Set<CqlTable> cdcTables)
+ {
+ Set<TableIdentifier> stillCdcEnabled = cdcTables.stream()
+ .filter(CqlTable::cdc)
+ .map(t ->
TableIdentifier.of(t.keyspace(), t.table()))
+
.collect(Collectors.toSet());
+ Set<TableIdentifier> stale = new HashSet<>(currentlyCdcEnabled);
+ stale.removeAll(stillCdcEnabled);
+
+ stale.forEach(id -> {
+ LOGGER.warn("Disabling CDC on table keyspace={} table={}",
id.keyspace(), id.table());
+ CassandraSchema.disableCdc(schema, id.keyspace(), id.table());
+ });
+ }
+
+ /**
+ * Removes tables from {@code Schema.instance} that were previously
registered via
+ * {@link #updateCdcSchema} but are no longer needed — e.g. a non-CDC
table that no longer
+ * shares partition-key structure with any CDC-enabled table in its
keyspace after a schema
+ * change. See {@code CdcBridge#unregisterNonCdcTables} for the full
rationale.
+ *
+ * <p>Idempotent: a table not currently registered is silently skipped.
Refuses (skips, with
+ * a warning) to unregister any table that is currently CDC-enabled — the
caller is
+ * responsible for only requesting removal of tables it has determined are
safe, but this is
+ * a last-line defense against silently dropping schema CDC still needs.
+ *
+ * @param tables the tables to unregister
+ */
+ public static void unregisterNonCdcTables(@NotNull Set<TableIdentifier>
tables)
+ {
+ unregisterNonCdcTables(Schema.instance, tables);
+ }
+
+ public static void unregisterNonCdcTables(@NotNull Schema schema, @NotNull
Set<TableIdentifier> tables)
+ {
+ for (TableIdentifier id : tables)
+ {
+ String keyspace = id.keyspace();
+ String table = id.table();
+ try
{
- // verify TableMetadata and ColumnFamilyStore initialized in
Schema
- TableId tableIdAfter = TableId.fromUUID(tableId);
-
Preconditions.checkNotNull(schema.getTableMetadata(tableIdAfter), "Table not
initialized in the schema");
-
Preconditions.checkArgument(Objects.requireNonNull(schema.getKeyspaceInstance(table.keyspace())).hasColumnFamilyStore(tableIdAfter),
- "ColumnFamilyStore not initialized
in the schema");
-
Preconditions.checkArgument(CassandraSchema.isCdcEnabled(schema, table),
- "CDC not enabled for table: " +
table.keyspace() + "." + table.table());
+ unregisterNonCdcTable(schema, keyspace, table);
+ }
+ catch (RuntimeException e)
+ {
+ // Don't let one bad table abort unregistration of the rest of
the batch.
+ LOGGER.warn("Failed to unregister table keyspace={} table={}",
keyspace, table, e);
}
}
- // existing table no longer with cdc = true, so disable
- cdcEnabledTables.forEach((ks, tables) -> tables.forEach(table -> {
- LOGGER.warn("Disabling CDC on table keyspace={} table={}", ks,
table);
- CassandraSchema.disableCdc(schema, ks, table);
- }));
+ }
+
+ private static void unregisterNonCdcTable(@NotNull Schema schema, @NotNull
String keyspace, @NotNull String table)
+ {
+ Optional<TableMetadata> tableMetadata = getTable(schema, keyspace,
table);
+ if (!tableMetadata.isPresent())
+ {
+ // already unregistered (or never was) — nothing to do
+ return;
+ }
+
+ if (tableMetadata.get().params.cdc)
+ {
+ LOGGER.warn("Refusing to unregister CDC-enabled table keyspace={}
table={}", keyspace, table);
+ return;
+ }
+
+ update(s -> {
+ Optional<KeyspaceMetadata> ks = getKeyspaceMetadata(s, keyspace);
+ Optional<TableMetadata> tableOpt = getTable(s, keyspace, table);
+ if (!ks.isPresent() || !tableOpt.isPresent())
+ {
+ // unregistered by a concurrent call, or keyspace itself is
gone
+ return;
+ }
+ if (tableOpt.get().params.cdc)
+ {
+ // became CDC-enabled since the check above (race with a
concurrent
+ // updateCdcSchema) — do not remove it
+ return;
+ }
+
+ LOGGER.info("Unregistering table no longer at risk of a batch with
a CDC-enabled table keyspace={} table={}", keyspace, table);
+ // Only remove the table's schema metadata (so it goes back to
throwing
+ // UnknownTableException on deserialization) — this bridge never
performs real
+ // writes/compactions on the tables it mirrors, so a full
Keyspace.dropCf() would
+ // pull in unrelated production machinery (e.g. lazily initializing
+ // CompactionManager's thread pools) with no corresponding benefit
here.
+ SchemaUpdater.load(s,
ks.get().withSwapped(ks.get().tables.without(table)));
+ });
}
public static void enableCdc(Schema schema, CqlTable cqlTable)
diff --git
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
index 5d0493a1..b565565b 100644
---
a/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
+++
b/cassandra-four-zero-types/src/main/java/org/apache/cassandra/spark/reader/SchemaBuilder.java
@@ -85,6 +85,7 @@ public class SchemaBuilder
private final ReplicationFactor replicationFactor;
private final CassandraTypes cassandraTypes;
private final int indexCount;
+ private final boolean enableCdc;
public SchemaBuilder(CqlTable table, Partitioner partitioner, boolean
enableCdc)
{
@@ -93,7 +94,7 @@ public class SchemaBuilder
public SchemaBuilder(CqlTable table, Partitioner partitioner)
{
- this(table, partitioner, null, false);
+ this(table, partitioner, null, table.cdc());
}
public SchemaBuilder(CqlTable table, Partitioner partitioner, UUID
tableId, boolean enableCdc)
@@ -137,6 +138,7 @@ public class SchemaBuilder
this.replicationFactor = replicationFactor;
this.cassandraTypes = new CassandraTypesImplementation();
this.indexCount = indexCount;
+ this.enableCdc = enableCdc;
Pair<KeyspaceMetadata, TableMetadata> updated =
CassandraSchema.apply(schema ->
updateSchema(schema,
@@ -477,7 +479,8 @@ public class SchemaBuilder
replicationFactor,
fields,
new HashSet<>(udts.values()),
- indexCount);
+ indexCount,
+ enableCdc);
}
private Map<String, CqlField.CqlUdt> buildsUdts(KeyspaceMetadata
keyspaceMetadata)
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]