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-sidecar.git
The following commit(s) were added to refs/heads/trunk by this push:
new bb735d44 CASSSIDECAR-489: Fix dead/dropped CDC lifecycle metrics and
duplicate SidecarCdcStats interface (#375)
bb735d44 is described below
commit bb735d447fea2548b420c3140609ae0e0afeacad
Author: Jyothsna konisa <[email protected]>
AuthorDate: Wed Aug 5 12:44:51 2026 -0700
CASSSIDECAR-489: Fix dead/dropped CDC lifecycle metrics and duplicate
SidecarCdcStats interface (#375)
Patch by Jyothsna Konisa; Reviewed by Josh McKenzie for CASSSIDECAR-489
---
CHANGES.txt | 1 +
.../sidecar/testing/TestCdcPublisher.java | 2 +-
...SharedClusterCdcSidecarIntegrationTestBase.java | 2 +-
.../cassandra/sidecar/cdc/CachingSchemaStore.java | 7 +-
.../cassandra/sidecar/cdc/CdcConsumerEntry.java | 7 +-
.../cassandra/sidecar/cdc/CdcEventConsumer.java | 7 +-
.../apache/cassandra/sidecar/cdc/CdcManager.java | 7 +-
.../apache/cassandra/sidecar/cdc/CdcPublisher.java | 9 +-
.../cassandra/sidecar/cdc/SidecarCdcStats.java | 277 ---------------------
.../sidecar/metrics/server/CdcMetrics.java | 5 +
.../cassandra/sidecar/modules/CdcModule.java | 2 +-
.../sidecar/tasks/CdcRawDirectorySpaceCleaner.java | 3 +-
.../sidecar/cdc/CachingSchemaStoreTest.java | 1 +
.../sidecar/cdc/CdcConsumerEntryTest.java | 20 +-
.../cassandra/sidecar/cdc/CdcManagerTest.java | 20 +-
.../cassandra/sidecar/cdc/CdcPublisherTests.java | 44 ++++
.../tasks/CdcRawDirectorySpaceCleanerTest.java | 5 +
17 files changed, 117 insertions(+), 302 deletions(-)
diff --git a/CHANGES.txt b/CHANGES.txt
index d9caa34b..71c6953d 100644
--- a/CHANGES.txt
+++ b/CHANGES.txt
@@ -1,5 +1,6 @@
0.5.0
-----
+ * Fix dead/dropped CDC lifecycle metrics and duplicate SidecarCdcStats
interface (CASSSIDECAR-489)
* Wire CDC configs in configs table to
SidecarCdcOptions/SidecarStatePersister (CASSSIDECAR-483)
* Implement durable operational job tracker (CASSSIDECAR-374)
* Remove filesystem path from Http response (CASSSIDECAR-477)
diff --git
a/integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/TestCdcPublisher.java
b/integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/TestCdcPublisher.java
index 8399a892..1aca8527 100644
---
a/integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/TestCdcPublisher.java
+++
b/integration-framework/src/main/java/org/apache/cassandra/sidecar/testing/TestCdcPublisher.java
@@ -25,12 +25,12 @@ import org.apache.cassandra.cdc.api.SchemaSupplier;
import org.apache.cassandra.cdc.kafka.KafkaProducerFactory;
import org.apache.cassandra.cdc.sidecar.ClusterConfigProvider;
import org.apache.cassandra.cdc.sidecar.SidecarCdcClient;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.sidecar.bridge.CassandraBridgeFactory;
import org.apache.cassandra.sidecar.cdc.CachingSchemaStore;
import org.apache.cassandra.sidecar.cdc.CdcConfig;
import org.apache.cassandra.sidecar.cdc.CdcPublisher;
-import org.apache.cassandra.sidecar.cdc.SidecarCdcStats;
import org.apache.cassandra.sidecar.concurrent.ExecutorPools;
import org.apache.cassandra.sidecar.coordination.RangeManager;
import org.apache.cassandra.sidecar.db.CdcDatabaseAccessor;
diff --git
a/integration-tests/src/integrationTest/org/apache/cassandra/sidecar/testing/SharedClusterCdcSidecarIntegrationTestBase.java
b/integration-tests/src/integrationTest/org/apache/cassandra/sidecar/testing/SharedClusterCdcSidecarIntegrationTestBase.java
index 70ff33b9..111220a7 100644
---
a/integration-tests/src/integrationTest/org/apache/cassandra/sidecar/testing/SharedClusterCdcSidecarIntegrationTestBase.java
+++
b/integration-tests/src/integrationTest/org/apache/cassandra/sidecar/testing/SharedClusterCdcSidecarIntegrationTestBase.java
@@ -32,13 +32,13 @@ import org.apache.cassandra.cdc.api.CdcOptions;
import org.apache.cassandra.cdc.api.SchemaSupplier;
import org.apache.cassandra.cdc.sidecar.ClusterConfigProvider;
import org.apache.cassandra.cdc.sidecar.SidecarCdcClient;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.distributed.api.ICluster;
import org.apache.cassandra.distributed.api.IInstance;
import org.apache.cassandra.sidecar.bridge.CassandraBridgeFactory;
import org.apache.cassandra.sidecar.cdc.CdcConfig;
import org.apache.cassandra.sidecar.cdc.CdcPublisher;
-import org.apache.cassandra.sidecar.cdc.SidecarCdcStats;
import org.apache.cassandra.sidecar.concurrent.ExecutorPools;
import org.apache.cassandra.sidecar.config.ServiceConfiguration;
import org.apache.cassandra.sidecar.config.SidecarClientConfiguration;
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStore.java
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStore.java
index 47db6ca2..2c812a34 100644
---
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStore.java
+++
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStore.java
@@ -45,6 +45,7 @@ import org.apache.cassandra.cdc.kafka.KafkaOptions;
import org.apache.cassandra.cdc.schemastore.SchemaStore;
import org.apache.cassandra.cdc.schemastore.SchemaStorePublisherFactory;
import org.apache.cassandra.cdc.schemastore.TableSchemaPublisher;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.sidecar.db.TableHistoryDatabaseAccessor;
import org.apache.cassandra.sidecar.db.schema.SidecarSchema;
import org.apache.cassandra.sidecar.tasks.CassandraClusterSchemaMonitor;
@@ -179,7 +180,11 @@ public class CachingSchemaStore implements SchemaStore
metadata.put(METADATA_NAME_KEY, cqlTable.table());
metadata.put(METADATA_NAMESPACE_KEY, cqlTable.keyspace());
publisher.publishSchema(mergedSchema.toString(false),
metadata);
- sidecarCdcStats.capturePublishedSchema();
+ // TODO: capturePublishedSchema() now lives on the
canonical
+ // org.apache.cassandra.cdc.sidecar.SidecarCdcStats
(cassandra-analytics-cdc-sidecar),
+ // which cassandra-sidecar consumes as a published
artifact. Re-wire
+ // sidecarCdcStats.capturePublishedSchema() here once that
artifact is released
+ // with this method included.
}
return new SchemaCacheEntry(cqlTable, payloadSchema);
});
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntry.java
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntry.java
index 2e33982f..c8a470a1 100644
---
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntry.java
+++
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntry.java
@@ -19,6 +19,7 @@
package org.apache.cassandra.sidecar.cdc;
import org.apache.cassandra.cdc.sidecar.SidecarCdc;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.sidecar.SidecarStatePersister;
/**
@@ -29,11 +30,13 @@ class CdcConsumerEntry
{
private final SidecarCdc consumer;
private final SidecarStatePersister persister;
+ private final SidecarCdcStats sidecarCdcStats;
- CdcConsumerEntry(SidecarCdc consumer, SidecarStatePersister persister)
+ CdcConsumerEntry(SidecarCdc consumer, SidecarStatePersister persister,
SidecarCdcStats sidecarCdcStats)
{
this.consumer = consumer;
this.persister = persister;
+ this.sidecarCdcStats = sidecarCdcStats;
}
SidecarCdc consumer()
@@ -51,11 +54,13 @@ class CdcConsumerEntry
persister.start();
consumer.initSchema();
consumer.start();
+ sidecarCdcStats.captureCdcConsumerStarted();
}
void stop()
{
consumer.stop(); // blocking — waits for any active run() to
complete
persister.stop(true); // flush buffered state to Cassandra, then
cancel timer
+ sidecarCdcStats.captureCdcConsumerStopped();
}
}
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcEventConsumer.java
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcEventConsumer.java
index 92db0077..1bf3461d 100644
---
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcEventConsumer.java
+++
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcEventConsumer.java
@@ -23,6 +23,7 @@ import java.util.function.Consumer;
import org.apache.cassandra.cdc.api.EventConsumer;
import org.apache.cassandra.cdc.kafka.KafkaPublisher;
import org.apache.cassandra.cdc.msg.CdcEvent;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.jetbrains.annotations.NotNull;
/**
@@ -31,10 +32,12 @@ import org.jetbrains.annotations.NotNull;
public class CdcEventConsumer implements EventConsumer
{
private final transient KafkaPublisher kafka;
+ private final SidecarCdcStats sidecarCdcStats;
- public CdcEventConsumer(KafkaPublisher kafka)
+ public CdcEventConsumer(KafkaPublisher kafka, SidecarCdcStats
sidecarCdcStats)
{
this.kafka = kafka;
+ this.sidecarCdcStats = sidecarCdcStats;
}
public void accept(CdcEvent cdcEvent)
@@ -45,7 +48,9 @@ public class CdcEventConsumer implements EventConsumer
@Override
public void flush() throws InterruptedException
{
+ long startNanos = System.nanoTime();
kafka.flush();
+ sidecarCdcStats.captureKafkaFlushTime(System.nanoTime() - startNanos);
}
public @NotNull Consumer<CdcEvent> andThen(@NotNull Consumer<? super
CdcEvent> after)
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java
index 8c97b1fd..6bcdb871 100644
--- a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java
+++ b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcManager.java
@@ -84,6 +84,7 @@ public class CdcManager
private final ClusterConfigProvider clusterConfigProvider;
private final SidecarCdcClient sidecarCdcClient;
private final ICdcStats cdcStats;
+ private final SidecarCdcStats sidecarCdcStats;
private List<CdcConsumerEntry> entries = new ArrayList<>();
private final ReplicationFactorSupplier rfSupplier;
private final CdcOptions cdcOptions;
@@ -99,6 +100,7 @@ public class CdcManager
ClusterConfigProvider clusterConfigProvider,
SidecarCdcClient sidecarCdcClient,
ICdcStats cdcStats,
+ SidecarCdcStats sidecarCdcStats,
TaskExecutorPool taskExecutorPool,
CdcDatabaseAccessor cdcDatabaseAccessor,
CdcOptions cdcOptions)
@@ -111,6 +113,7 @@ public class CdcManager
this.clusterConfigProvider = clusterConfigProvider;
this.sidecarCdcClient = sidecarCdcClient;
this.cdcStats = cdcStats;
+ this.sidecarCdcStats = sidecarCdcStats;
this.cdcOptions = cdcOptions;
this.asyncExecutor = new ExecutorPoolsExecutor(taskExecutorPool);
this.cassandraClient = new
StateSidecarCdcCassandraClient(cdcDatabaseAccessor);
@@ -216,14 +219,14 @@ public class CdcManager
.withReplicationFactorSupplier(rfSupplier)
.withSidecarStatePersister(persister)
.build();
- return new CdcConsumerEntry(consumer, persister);
+ return new CdcConsumerEntry(consumer, persister, sidecarCdcStats);
}
private @NotNull SidecarStatePersister getSidecarStatePersister()
{
return new SidecarStatePersister(new
ConfigBackedPersisterOptions(conf),
cdcOptions,
- SidecarCdcStats.STUB,
+ sidecarCdcStats,
cassandraClient,
asyncExecutor);
}
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcPublisher.java
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcPublisher.java
index 2fc6ac20..5a110540 100644
--- a/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcPublisher.java
+++ b/server/src/main/java/org/apache/cassandra/sidecar/cdc/CdcPublisher.java
@@ -38,6 +38,7 @@ import org.apache.cassandra.cdc.kafka.KafkaPublisher;
import org.apache.cassandra.cdc.kafka.TopicSupplier;
import org.apache.cassandra.cdc.sidecar.ClusterConfigProvider;
import org.apache.cassandra.cdc.sidecar.SidecarCdcClient;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.sidecar.bridge.CassandraBridgeFactory;
import org.apache.cassandra.sidecar.common.server.utils.DurationSpec;
@@ -120,6 +121,7 @@ public class CdcPublisher implements
Handler<Message<Object>>, PeriodicTask
if (conf.cdcEnabled())
{
+ sidecarCdcStats.captureCdcEnabled();
vertx.eventBus().localConsumer(RangeManager.RangeManagerEvents.ON_TOKEN_RANGE_CHANGED.address(),
this);
vertx.eventBus().localConsumer(RangeManager.LeadershipEvents.ON_TOKEN_RANGE_GAINED.address(),
this);
vertx.eventBus().localConsumer(RangeManager.LeadershipEvents.ON_TOKEN_RANGE_LOST.address(),
this);
@@ -127,6 +129,10 @@ public class CdcPublisher implements
Handler<Message<Object>>, PeriodicTask
vertx.eventBus().localConsumer(ON_CDC_CACHE_WARMED_UP.address(),
this);
vertx.eventBus().localConsumer(ON_CDC_CONFIGURATION_CHANGED.address(), new
ConfigChangedHandler());
}
+ else
+ {
+ sidecarCdcStats.captureCdcDisabled();
+ }
}
public EventConsumer eventConsumer(CdcConfig conf)
@@ -146,7 +152,7 @@ public class CdcPublisher implements
Handler<Message<Object>>, PeriodicTask
conf.failOnRecordTooLargeError(),
conf.failOnKafkaError(),
CdcLogMode.FULL);
- return new CdcEventConsumer(kafkaPublisher);
+ return new CdcEventConsumer(kafkaPublisher, sidecarCdcStats);
}
/**
@@ -220,6 +226,7 @@ public class CdcPublisher implements
Handler<Message<Object>>, PeriodicTask
clusterConfigProvider,
sidecarCdcClientProvider.get(),
cdcStats,
+ sidecarCdcStats,
this.executorPools,
databaseAccessor,
cdcOptions);
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/cdc/SidecarCdcStats.java
b/server/src/main/java/org/apache/cassandra/sidecar/cdc/SidecarCdcStats.java
deleted file mode 100644
index 93588a43..00000000
--- a/server/src/main/java/org/apache/cassandra/sidecar/cdc/SidecarCdcStats.java
+++ /dev/null
@@ -1,277 +0,0 @@
-/*
- * Licensed to the Apache Software Foundation (ASF) under one
- * or more contributor license agreements. See the NOTICE file
- * distributed with this work for additional information
- * regarding copyright ownership. The ASF licenses this file
- * to you under the Apache License, Version 2.0 (the
- * "License"); you may not use this file except in compliance
- * with the License. You may obtain a copy of the License at
- *
- * http://www.apache.org/licenses/LICENSE-2.0
- *
- * Unless required by applicable law or agreed to in writing,
- * software distributed under the License is distributed on an
- * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
- * KIND, either express or implied. See the License for the
- * specific language governing permissions and limitations
- * under the License.
- */
-
-package org.apache.cassandra.sidecar.cdc;
-
-/**
- * Interface for capturing and reporting CDC (Change Data Capture) related
statistics and metrics
- * in the Cassandra Sidecar application.
- */
-public interface SidecarCdcStats
-{
- /**
- * Cdc is enabled
- */
- default void captureCdcEnabled()
- {
- }
-
- /**
- * Cdc disabled and CdcPublisher will not be started.
- */
- default void captureCdcDisabled()
- {
- }
-
- // cdc lifecycle stats
-
- /**
- * CdcPublisher started successfully.
- *
- * @param numConsumers number of Cdc consumers started
- */
- default void captureCdcStarted(int numConsumers)
- {
- }
-
- /**
- * CdcPublisher failed to start.
- *
- * @param t throwable
- */
- default void captureCdcStartFailure(Throwable t)
- {
- }
-
- /**
- * CdcConsumer started.
- */
- default void captureCdcConsumerStarted()
- {
- }
-
- /**
- * CdcConsumer stopped.
- */
- default void captureCdcConsumerStopped()
- {
- }
-
- /**
- * Cdc or Kafka config change resulting in Cdc restart.
- */
- default void captureCdcConfigChange()
- {
- }
-
- /**
- * Cassandra cluster topology changed resulting in Cdc restart.
- */
- default void captureCdcClusterTopologyChange()
- {
- }
-
- /**
- * Cdc gained ownership for a token range.
- */
- default void captureCdcTokenRangeGained()
- {
- }
-
- /**
- * Cdc lost ownership for a token range.
- */
- default void captureCdcTokenRangeLost()
- {
- }
-
- /**
- * CdcPublisher restarted.
- */
- default void captureCdcRestart()
- {
- }
-
- /**
- * CdcPublisher stopped.
- */
- default void captureCdcStopped()
- {
- }
-
- /**
- * CdcPublisher failed to stop gracefully.
- *
- * @param t throwable
- */
- default void captureCdcStopFailed(Throwable t)
- {
- }
-
- /**
- * CdcConsumer read from state.
- *
- * @param count number of state objects.
- * @param len byte length of state.
- */
- default void captureCdcConsumerReadFromState(int count, int len)
- {
- }
-
- default void captureCdcStateDeserializationFailure()
- {
- }
-
- /**
- * New blank CdcConsumer initialized.
- */
- default void captureNewBlankCdcConsumer()
- {
- }
-
- // cdc consumer stats
-
- /**
- * Cdc consumer completed
- *
- * @param epoch epoch number, monotonically increasing 64-bit
signed integer.
- * @param numDownInstances number of down instances
- * @param numTables number of cached CDC enabled tables
- * @param runtimeNanos runtime of the previous epoch in nanoseconds
- */
- default void captureCdcNextEpoch(long epoch, int numDownInstances, int
numTables, long runtimeNanos)
- {
-
- }
-
- /**
- * Single Cdc event processed.
- */
- default void captureCdcEventProcessed()
- {
- }
-
- /**
- * Recoverable Error detected in the CdcConsumer.
- */
- default void captureRecoverableCdcError(Throwable t)
- {
- }
-
- /**
- * Unrecoverable Error detected in the CdcConsumer causing it to stop.
- */
- default void captureUnrecoverableCdcError(Throwable t)
- {
- }
-
- // state persist stats
-
- /**
- * Kafka queued events flushed.
- *
- * @param timeNanos time taken to flush in nanoseconds.
- */
- default void captureKafkaFlushTime(long timeNanos)
- {
- }
-
- /**
- * State persister backed up and not keeping up with persist cadence.
- *
- * @param numTasks number of active flush tasks queued.
- */
- default void capturePersistBackedUp(int numTasks)
- {
- }
-
- /**
- * Persist state request to storage.
- *
- * @param len byte length of state
- */
- default void capturePersistingCdcStateLength(int len)
- {
- }
-
- /**
- * Persist request succeeded.
- *
- * @param timeNanos time in nanos taken to persist.
- */
- default void capturePersistSucceeded(long timeNanos)
- {
- }
-
- /**
- * Persist failed with error.
- *
- * @param t throwable
- */
- default void capturePersistFailed(Throwable t)
- {
- }
-
- /**
- * Schema has been published.
- */
- default void capturePublishedSchema()
- {
- }
-
- /**
- * Cdc state has been written through http api.
- *
- * @param len length of the request body
- */
- default void captureCdcStatePutRequest(int len)
- {
- }
-
- /**
- * PutCdcStateHandler failed to write cdc state to store.
- */
- default void captureCdcStatePutRequestFailure()
- {
- }
-
- /**
- * Cdc state has been read through the http api.
- *
- * @param count number of Cdc state objects returned
- * @param len length of the response body
- */
- default void captureCdcStateGetRequest(int count, int len)
- {
- }
-
- /**
- * A Cdc enabled table is not enabled in the Schema.instance singleton,
this is a critical alert as the table will skipped when reading the commit log.
- */
- default void captureCdcTableNotEnabled()
- {
- }
-
- /**
- * cdc_on_repair is enabled in the yaml file and should be disabled.
- */
- default void captureCdcOnRepairEnabled()
- {
- }
-}
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/metrics/server/CdcMetrics.java
b/server/src/main/java/org/apache/cassandra/sidecar/metrics/server/CdcMetrics.java
index 1f0765e6..cf25f0bb 100644
---
a/server/src/main/java/org/apache/cassandra/sidecar/metrics/server/CdcMetrics.java
+++
b/server/src/main/java/org/apache/cassandra/sidecar/metrics/server/CdcMetrics.java
@@ -41,6 +41,7 @@ public class CdcMetrics
public final NamedMetric<DefaultSettableGauge<Integer>> oldestSegmentAge;
public final NamedMetric<DeltaGauge> totalConsumedCdcBytes;
public final NamedMetric<DefaultSettableGauge<Long>> totalCdcSpaceUsed;
+ public final NamedMetric<DefaultSettableGauge<Long>> maxCdcSpaceConfigured;
public final NamedMetric<DeltaGauge> deletedSegment;
public final NamedMetric<DeltaGauge> lowCdcRawSpace;
public final NamedMetric<DeltaGauge> criticalCdcRawSpace;
@@ -51,6 +52,10 @@ public class CdcMetrics
this.cdcRawCleanerFailed = createMetric("CleanerFailed", name ->
metricRegistry.gauge(name, DeltaGauge::new));
this.totalConsumedCdcBytes = createMetric("TotalConsumedBytes", name
-> metricRegistry.gauge(name, DeltaGauge::new));
this.totalCdcSpaceUsed = createMetric("TotalSpaceUsed", name ->
metricRegistry.gauge(name, () -> new DefaultSettableGauge<>(0L)));
+ // Configured cdc_total_space limit
(CdcRawDirectorySpaceCleaner.maxUsageBytes()), so a
+ // dashboard can show capacity alongside TotalSpaceUsed. Peak usage,
if ever needed, can be
+ // derived from TotalSpaceUsed with stat=max instead of a dedicated
gauge.
+ this.maxCdcSpaceConfigured = createMetric("MaxSpaceConfigured", name
-> metricRegistry.gauge(name, () -> new DefaultSettableGauge<>(0L)));
this.orphanedIdx = createMetric("OrphanedIdxFile", name ->
metricRegistry.gauge(name, DeltaGauge::new));
this.deletedSegment = createMetric("DeletedSegment", name ->
metricRegistry.gauge(name, DeltaGauge::new));
this.oldestSegmentAge = createMetric("OldestSegmentAgeSeconds", name
-> metricRegistry.gauge(name, () -> new DefaultSettableGauge<>(0)));
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java
b/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java
index 9d1811d4..5edcd976 100644
--- a/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java
+++ b/server/src/main/java/org/apache/cassandra/sidecar/modules/CdcModule.java
@@ -44,6 +44,7 @@ import
org.apache.cassandra.cdc.schemastore.SchemaStorePublisherFactory;
import org.apache.cassandra.cdc.sidecar.CdcSidecarInstancesProvider;
import org.apache.cassandra.cdc.sidecar.ClusterConfigProvider;
import org.apache.cassandra.cdc.sidecar.SidecarCdcClient;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.stats.CdcStats;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.secrets.SecretsProvider;
@@ -56,7 +57,6 @@ import org.apache.cassandra.sidecar.cdc.CdcLogCache;
import org.apache.cassandra.sidecar.cdc.CdcPublisher;
import org.apache.cassandra.sidecar.cdc.CdcSchemaSupplier;
import org.apache.cassandra.sidecar.cdc.SidecarCdcOptions;
-import org.apache.cassandra.sidecar.cdc.SidecarCdcStats;
import org.apache.cassandra.sidecar.cdc.SidecarClientSecretsProvider;
import org.apache.cassandra.sidecar.cdc.SidecarClusterConfigProvider;
import org.apache.cassandra.sidecar.cdc.SidecarCqlToAvroSchemaConverter;
diff --git
a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java
b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java
index 03ed40a1..7d10b0d7 100644
---
a/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java
+++
b/server/src/main/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleaner.java
@@ -218,6 +218,8 @@ public class CdcRawDirectorySpaceCleaner implements
PeriodicTask
)
.orElseGet(List::of);
publishCdcStats(segmentFiles);
+ long maxUsageBytes = maxUsageBytes();
+ cdcMetrics.maxCdcSpaceConfigured.metric.setValue(maxUsageBytes);
if (segmentFiles.size() < 2)
{
LOGGER.debug("Skipping cdc data cleaner routine cleanup: No cdc
data or only one single cdc segment is found.");
@@ -225,7 +227,6 @@ public class CdcRawDirectorySpaceCleaner implements
PeriodicTask
}
long directorySizeBytes =
FileUtils.directorySizeBytes(cdcRawDirectory);
- long maxUsageBytes = maxUsageBytes();
long upperLimitBytes = (long) (maxUsageBytes *
cdcConfiguration.cdcRawDirectoryMaxPercentUsage());
// Sort the files by segmentId to delete commit log segments in write
order
// The latest file is the current active segment, but it could be
created before the retention duration, e.g. slow data ingress
diff --git
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStoreTest.java
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStoreTest.java
index b8c8c085..77746603 100644
---
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStoreTest.java
+++
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CachingSchemaStoreTest.java
@@ -40,6 +40,7 @@ import org.apache.cassandra.cdc.api.TableIdLookup;
import org.apache.cassandra.cdc.avro.AvroSchemas;
import org.apache.cassandra.cdc.avro.CqlToAvroSchemaConverter;
import org.apache.cassandra.cdc.schemastore.SchemaStorePublisherFactory;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.sidecar.bridge.CassandraBridgeFactory;
import org.apache.cassandra.sidecar.db.TableHistoryDatabaseAccessor;
import org.apache.cassandra.sidecar.db.schema.SidecarSchema;
diff --git
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntryTest.java
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntryTest.java
index 9891905f..0d9735cf 100644
---
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntryTest.java
+++
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcConsumerEntryTest.java
@@ -21,6 +21,7 @@ package org.apache.cassandra.sidecar.cdc;
import org.junit.jupiter.api.Test;
import org.apache.cassandra.cdc.sidecar.SidecarCdc;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.sidecar.SidecarStatePersister;
import org.mockito.InOrder;
@@ -32,32 +33,36 @@ import static org.mockito.Mockito.mock;
public class CdcConsumerEntryTest
{
@Test
- void startCallsPersisterBeforeConsumer()
+ void startCallsPersisterBeforeConsumerThenCapturesConsumerStarted()
{
SidecarCdc consumer = mock(SidecarCdc.class);
SidecarStatePersister persister = mock(SidecarStatePersister.class);
- CdcConsumerEntry entry = new CdcConsumerEntry(consumer, persister);
+ SidecarCdcStats sidecarCdcStats = mock(SidecarCdcStats.class);
+ CdcConsumerEntry entry = new CdcConsumerEntry(consumer, persister,
sidecarCdcStats);
entry.start();
- InOrder order = inOrder(persister, consumer);
+ InOrder order = inOrder(persister, consumer, sidecarCdcStats);
order.verify(persister).start();
order.verify(consumer).initSchema();
order.verify(consumer).start();
+ order.verify(sidecarCdcStats).captureCdcConsumerStarted();
}
@Test
- void stopCallsConsumerBeforePersister()
+ void stopCallsConsumerBeforePersisterThenCapturesConsumerStopped()
{
SidecarCdc consumer = mock(SidecarCdc.class);
SidecarStatePersister persister = mock(SidecarStatePersister.class);
- CdcConsumerEntry entry = new CdcConsumerEntry(consumer, persister);
+ SidecarCdcStats sidecarCdcStats = mock(SidecarCdcStats.class);
+ CdcConsumerEntry entry = new CdcConsumerEntry(consumer, persister,
sidecarCdcStats);
entry.stop();
- InOrder order = inOrder(consumer, persister);
+ InOrder order = inOrder(consumer, persister, sidecarCdcStats);
order.verify(consumer).stop();
order.verify(persister).stop(true);
+ order.verify(sidecarCdcStats).captureCdcConsumerStopped();
}
@Test
@@ -65,7 +70,8 @@ public class CdcConsumerEntryTest
{
SidecarCdc consumer = mock(SidecarCdc.class);
SidecarStatePersister persister = mock(SidecarStatePersister.class);
- CdcConsumerEntry entry = new CdcConsumerEntry(consumer, persister);
+ SidecarCdcStats sidecarCdcStats = mock(SidecarCdcStats.class);
+ CdcConsumerEntry entry = new CdcConsumerEntry(consumer, persister,
sidecarCdcStats);
assertThat(entry.consumer()).isSameAs(consumer);
assertThat(entry.persister()).isSameAs(persister);
diff --git
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java
index 0701e855..0142e4ee 100644
--- a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java
+++ b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcManagerTest.java
@@ -40,6 +40,7 @@ import org.apache.cassandra.cdc.api.SchemaSupplier;
import org.apache.cassandra.cdc.sidecar.ClusterConfigProvider;
import org.apache.cassandra.cdc.sidecar.SidecarCdc;
import org.apache.cassandra.cdc.sidecar.SidecarCdcClient;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.sidecar.SidecarStatePersister;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.sidecar.cluster.instance.InstanceMetadata;
@@ -87,6 +88,8 @@ public class CdcManagerTest
@Mock
private ICdcStats cdcStats;
@Mock
+ private SidecarCdcStats sidecarCdcStats;
+ @Mock
private TaskExecutorPool taskExecutorPool;
@Mock
private CdcDatabaseAccessor cdcDatabaseAccessor;
@@ -109,6 +112,7 @@ public class CdcManagerTest
clusterConfigProvider,
sidecarCdcClient,
cdcStats,
+ sidecarCdcStats,
taskExecutorPool,
cdcDatabaseAccessor,
cdcOptions
@@ -152,7 +156,7 @@ public class CdcManagerTest
when(cdcConfig.jobId()).thenReturn("test-job");
CdcManager spyManager = spy(cdcManager);
- CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
+ CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
doReturn(mockEntry).when(spyManager).buildConsumer(
any(), anyInt(), any(), any(), any(), any(), any(), any()
);
@@ -183,8 +187,8 @@ public class CdcManagerTest
when(cdcConfig.jobId()).thenReturn("test-job");
CdcManager spyManager = spy(cdcManager);
- CdcConsumerEntry mockEntry1 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
- CdcConsumerEntry mockEntry2 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
+ CdcConsumerEntry mockEntry1 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
+ CdcConsumerEntry mockEntry2 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
doReturn(mockEntry1, mockEntry2).when(spyManager).buildConsumer(
any(), anyInt(), any(), any(), any(), any(), any(), any()
);
@@ -218,8 +222,8 @@ public class CdcManagerTest
when(cdcConfig.jobId()).thenReturn("test-job");
CdcManager spyManager = spy(cdcManager);
- CdcConsumerEntry mockEntry1 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
- CdcConsumerEntry mockEntry2 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
+ CdcConsumerEntry mockEntry1 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
+ CdcConsumerEntry mockEntry2 = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
doReturn(mockEntry1, mockEntry2).when(spyManager).buildConsumer(
any(), anyInt(), any(), any(), any(), any(), any(), any()
);
@@ -251,7 +255,7 @@ public class CdcManagerTest
when(cdcConfig.jobId()).thenReturn("test-job");
CdcManager spyManager = spy(cdcManager);
- CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
+ CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
doReturn(mockEntry).when(spyManager).buildConsumer(
any(), anyInt(), any(), any(), any(), any(), any(), any()
);
@@ -275,7 +279,7 @@ public class CdcManagerTest
// Spy to mock buildConsumer - will be called with instanceId = -1
CdcManager spyManager = spy(cdcManager);
- CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
+ CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
doReturn(mockEntry).when(spyManager).buildConsumer(
any(), anyInt(), any(), any(), any(), any(), any(), any()
);
@@ -321,7 +325,7 @@ public class CdcManagerTest
when(cdcConfig.jobId()).thenReturn("test-job");
CdcManager spyManager = spy(cdcManager);
- CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class));
+ CdcConsumerEntry mockEntry = new
CdcConsumerEntry(mock(SidecarCdc.class), mock(SidecarStatePersister.class),
mock(SidecarCdcStats.class));
doReturn(mockEntry).when(spyManager).buildConsumer(
any(), anyInt(), any(), any(), any(), any(), any(), any()
);
diff --git
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcPublisherTests.java
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcPublisherTests.java
index 94d20037..8bf52619 100644
---
a/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcPublisherTests.java
+++
b/server/src/test/java/org/apache/cassandra/sidecar/cdc/CdcPublisherTests.java
@@ -41,6 +41,7 @@ import org.apache.cassandra.cdc.kafka.KafkaProducerFactory;
import org.apache.cassandra.cdc.kafka.TopicSupplier;
import org.apache.cassandra.cdc.sidecar.ClusterConfigProvider;
import org.apache.cassandra.cdc.sidecar.SidecarCdcClient;
+import org.apache.cassandra.cdc.sidecar.SidecarCdcStats;
import org.apache.cassandra.cdc.stats.ICdcStats;
import org.apache.cassandra.sidecar.bridge.CassandraBridgeFactory;
import org.apache.cassandra.sidecar.cluster.instance.InstanceMetadata;
@@ -158,6 +159,49 @@ public class CdcPublisherTests
assertThat(result).isInstanceOf(CdcEventConsumer.class);
}
+ /**
+ * Verifies that construction captures the Cdc enabled/disabled lifecycle
stat depending on
+ * {@link CdcConfig#cdcEnabled()} at the time the {@link CdcPublisher} is
built, mirroring the
+ * javadoc contract on {@link SidecarCdcStats#captureCdcDisabled()}: "Cdc
disabled and CdcPublisher
+ * will not be started."
+ */
+ @Test
+ void testConstructorCapturesCdcEnabledWhenConfigEnablesCdc()
+ {
+ CdcConfig enabledConfig = mock(CdcConfig.class, RETURNS_DEEP_STUBS);
+ when(enabledConfig.cdcEnabled()).thenReturn(true);
+ // setUp()'s own CdcPublisher construction already recorded a call
against the shared
+ // sidecarCdcStats mock (cdcConfig defaults to cdcEnabled() == false)
— use a fresh mock
+ // here so this test only observes the constructor call under test.
+ SidecarCdcStats freshSidecarCdcStats = mock(SidecarCdcStats.class);
+
+ new CdcPublisher(
+ vertx, executorPools, clusterConfigProvider, schemaSupplier,
instanceMetadataFetcher,
+ enabledConfig, databaseAccessor, cdcStats, systemViews,
freshSidecarCdcStats, rangeManager,
+ cassandraBridgeFactory, () -> sidecarCdcClient, schemaStore,
kafkaProducerFactory, cdcOptions
+ );
+
+ verify(freshSidecarCdcStats).captureCdcEnabled();
+ verify(freshSidecarCdcStats, never()).captureCdcDisabled();
+ }
+
+ @Test
+ void testConstructorCapturesCdcDisabledWhenConfigDisablesCdc()
+ {
+ CdcConfig disabledConfig = mock(CdcConfig.class, RETURNS_DEEP_STUBS);
+ when(disabledConfig.cdcEnabled()).thenReturn(false);
+ SidecarCdcStats freshSidecarCdcStats = mock(SidecarCdcStats.class);
+
+ new CdcPublisher(
+ vertx, executorPools, clusterConfigProvider, schemaSupplier,
instanceMetadataFetcher,
+ disabledConfig, databaseAccessor, cdcStats, systemViews,
freshSidecarCdcStats, rangeManager,
+ cassandraBridgeFactory, () -> sidecarCdcClient, schemaStore,
kafkaProducerFactory, cdcOptions
+ );
+
+ verify(freshSidecarCdcStats).captureCdcDisabled();
+ verify(freshSidecarCdcStats, never()).captureCdcEnabled();
+ }
+
/**
* Verifies that {@link CdcPublisher#eventConsumer(CdcConfig)} routes to
the correct
* {@link TopicSupplier} factory for every {@link
CdcConfig.TopicFormatType} value.
diff --git
a/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java
b/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java
index 6c0dc734..34e48435 100644
---
a/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java
+++
b/server/src/test/java/org/apache/cassandra/sidecar/tasks/CdcRawDirectorySpaceCleanerTest.java
@@ -131,6 +131,8 @@ class CdcRawDirectorySpaceCleanerTest
assertThat(cdcMetrics.criticalCdcRawSpace.metric.getValue()).isOne();
assertThat(cdcMetrics.orphanedIdx.metric.getValue()).isOne();
assertThat(cdcMetrics.totalCdcSpaceUsed.metric.getValue()).isGreaterThan(2097152L);
+ // matches the "cdc_total_space": "1MiB" setting stubbed above
+
assertThat(cdcMetrics.maxCdcSpaceConfigured.metric.getValue()).isEqualTo(1024L
* 1024L);
assertThat(cdcMetrics.deletedSegment.metric.getValue()).isGreaterThan(2097152L);
assertThat(cdcMetrics.oldestSegmentAge.metric.getValue()).isZero();
@@ -138,6 +140,9 @@ class CdcRawDirectorySpaceCleanerTest
// We do not expect all CDC file to be cleaned up in a running system.
But test it for robustness.
Files.deleteIfExists(Paths.get(tempDir.toString(),
CdcRawDirectorySpaceCleaner.CDC_DIR_NAME, TEST_INTACT_SEGMENT_FILE_NAME));
cleaner.routineCleanUp(); // it should run fine.
+
+ // configured limit is unaffected by usage dropping after the final
cleanup
+
assertThat(cdcMetrics.maxCdcSpaceConfigured.metric.getValue()).isEqualTo(1024L
* 1024L);
}
@Test
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]