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 dd4d2563 CASSANALYTICS-191: CDC reader stats silently dropped in
SidecarCdcBuilder (#232)
dd4d2563 is described below
commit dd4d25637a50d9b20138e657972174b4bfb6fab3
Author: Jyothsna konisa <[email protected]>
AuthorDate: Mon Aug 10 12:19:39 2026 -0700
CASSANALYTICS-191: CDC reader stats silently dropped in SidecarCdcBuilder
(#232)
Patch by Jyothsna Konisa; Reviewed by Saranya Krishnakumar for
CASSANALYTICS-191
---
CHANGES.txt | 1 +
.../cassandra/cdc/sidecar/SidecarCdcBuilder.java | 1 +
.../cassandra/cdc/sidecar/SidecarCdcTest.java | 56 ++++++++++++++++++++++
.../main/java/org/apache/cassandra/cdc/Cdc.java | 12 +++++
.../java/org/apache/cassandra/cdc/CdcBuilder.java | 3 +-
5 files changed, 71 insertions(+), 2 deletions(-)
diff --git a/CHANGES.txt b/CHANGES.txt
index 4d74b003..c3cda0a5 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
0.5.0
-----
+ * CDC reader stats silently dropped in SidecarCdcBuilder (CASSANALYTICS-191)
* Add CapturePublishedSchema metric to SidecarCdcStats (CASSANALYTICS-189)
* Expand list of architecture that supports unaligned access in
FastByteOperations (CASSANALYTICS-188)
* Fix FastByteOperations Silently Falling Back to Pure-Java Comparator
(CASSANALYTICS-187)
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 559edc2f..26c57279 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
@@ -55,6 +55,7 @@ public class SidecarCdcBuilder extends CdcBuilder
super(jobId, partitionId, eventConsumer, schemaSupplier);
this.clusterConfigProvider = clusterConfigProvider;
this.sidecarCdcClient = sidecarCdcClient;
+ withStats(cdcStats);
withCdcOptions(cdcOptions);
withTokenRangeSupplier(tokenRangeSupplier);
}
diff --git
a/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
b/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
index 27c41ee1..e2988296 100644
---
a/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
+++
b/cassandra-analytics-cdc-sidecar/src/test/java/org/apache/cassandra/cdc/sidecar/SidecarCdcTest.java
@@ -20,6 +20,8 @@
package org.apache.cassandra.cdc.sidecar;
import java.util.Map;
+import java.util.Set;
+import java.util.concurrent.CompletableFuture;
import org.junit.jupiter.api.BeforeEach;
import org.junit.jupiter.api.Test;
@@ -31,10 +33,14 @@ import org.apache.cassandra.cdc.api.EventConsumer;
import org.apache.cassandra.cdc.api.SchemaSupplier;
import org.apache.cassandra.cdc.api.TokenRangeSupplier;
import org.apache.cassandra.cdc.stats.ICdcStats;
+import org.apache.cassandra.spark.data.CqlTable;
+import org.apache.cassandra.spark.data.ReplicationFactor;
import org.apache.cassandra.spark.data.partitioner.CassandraInstance;
+import org.apache.cassandra.spark.utils.AsyncExecutor;
import static org.assertj.core.api.Assertions.assertThat;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.when;
/**
* Unit tests for SidecarCdc class
@@ -81,6 +87,56 @@ public class SidecarCdcTest
assertThat(builder.sidecarCdcClient).isEqualTo(mockSidecarCdcClient);
}
+ /**
+ * Regression test for a bug where {@link SidecarCdcBuilder}'s constructor
accepted an
+ * {@link ICdcStats} parameter but never wired it into the builder (via
{@link SidecarCdcBuilder#withStats}),
+ * so every {@link SidecarCdc} built through it silently used {@link
ICdcStats#STUB} instead of the real,
+ * caller-supplied stats implementation — with no exception anywhere to
reveal it. This mirrors the exact
+ * call shape a consuming project (e.g. cassandra-sidecar's CdcManager)
uses in production:
+ * {@code
SidecarCdc.builder(...).withExecutor(...).withReplicationFactorSupplier(...).withSidecarStatePersister(...).build()}.
+ */
+ @Test
+ public void testBuiltSidecarCdcUsesSuppliedStatsNotStub() throws Exception
+ {
+ String jobId = "test-job-123";
+ int partitionId = 0;
+ CdcOptions cdcOptions = mock(CdcOptions.class);
+ ClusterConfigProvider clusterConfigProvider =
mock(ClusterConfigProvider.class);
+ when(clusterConfigProvider.dc()).thenReturn("DC1");
+ EventConsumer eventConsumer = mock(EventConsumer.class);
+ TokenRangeSupplier tokenRangeSupplier = mock(TokenRangeSupplier.class);
+ SidecarCdcClient mockSidecarCdcClient = mock(SidecarCdcClient.class);
+ AsyncExecutor asyncExecutor = mock(AsyncExecutor.class);
+
+ // Just enough of a CDC-enabled table (with a replication factor for
"DC1") to satisfy
+ // SidecarCdc.initSchema(), which runs synchronously inside the
constructor.
+ ReplicationFactor rf = new
ReplicationFactor(ReplicationFactor.ReplicationStrategy.NetworkTopologyStrategy,
+ Map.of("DC1", 3));
+ CqlTable cqlTable = mock(CqlTable.class);
+ when(cqlTable.replicationFactor()).thenReturn(rf);
+ SchemaSupplier schemaSupplier = mock(SchemaSupplier.class);
+
when(schemaSupplier.getCDCEnabledTables()).thenReturn(CompletableFuture.completedFuture(Set.of(cqlTable)));
+
+ SidecarCdc consumer = SidecarCdc.builder(jobId,
+ partitionId,
+ cdcOptions,
+ clusterConfigProvider,
+ eventConsumer,
+ schemaSupplier,
+ tokenRangeSupplier,
+ mockSidecarCdcClient,
+ cdcStats)
+ .withExecutor(asyncExecutor)
+ .build();
+
+ assertThat(consumer.stats())
+ .as("SidecarCdc.builder(...)'s cdcStats argument must reach
Cdc.stats — if it doesn't, every "
+ + "ICdcStats call (changeProduced, insufficientReplicas,
mutationsReadCount, etc.) silently "
+ + "no-ops against ICdcStats.STUB instead of the real
implementation, with no exception to reveal it.")
+ .isSameAs(cdcStats)
+ .isNotSameAs(ICdcStats.STUB);
+ }
+
@Test
public void testPerInstancePortResolution()
{
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 57484324..630a1f3c 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
@@ -36,6 +36,7 @@ import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import com.esotericsoftware.kryo.io.Output;
+import com.google.common.annotations.VisibleForTesting;
import org.apache.cassandra.bridge.CassandraBridge;
import org.apache.cassandra.bridge.CdcBridge;
import org.apache.cassandra.bridge.CdcBridgeFactory;
@@ -110,6 +111,17 @@ public class Cdc implements Closeable
return new CdcBuilder(jobId, partitionId, eventConsumer,
schemaSupplier);
}
+ /**
+ * @return the {@link ICdcStats} this {@link Cdc} instance was built with.
Exposed so tests can assert
+ * the stats implementation supplied to the builder is actually the
instance in use, and not the
+ * {@link ICdcStats#STUB} default silently falling through.
+ */
+ @VisibleForTesting
+ public ICdcStats stats()
+ {
+ return stats;
+ }
+
public String jobId()
{
return jobId;
diff --git
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
index cd8b3ea9..f7c57966 100644
---
a/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
+++
b/cassandra-analytics-cdc/src/main/java/org/apache/cassandra/cdc/CdcBuilder.java
@@ -31,7 +31,6 @@ 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.api.TokenRangeSupplier;
-import org.apache.cassandra.cdc.stats.CdcStats;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.spark.utils.AsyncExecutor;
import org.jetbrains.annotations.NotNull;
@@ -131,7 +130,7 @@ public class CdcBuilder
return this;
}
- public CdcBuilder withStats(@NotNull CdcStats stats)
+ public CdcBuilder withStats(@NotNull ICdcStats stats)
{
this.stats = stats;
return this;
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]