Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
1 change: 1 addition & 0 deletions CHANGES.txt
Original file line number Diff line number Diff line change
@@ -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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,12 +25,12 @@
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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -32,13 +32,13 @@
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;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -45,6 +45,7 @@
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;
Expand Down Expand Up @@ -179,7 +180,11 @@ private void publishSchemas()
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);
});
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;

/**
Expand All @@ -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()
Expand All @@ -51,11 +54,13 @@ void start()
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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -23,6 +23,7 @@
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;

/**
Expand All @@ -31,10 +32,12 @@
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)
Expand All @@ -45,7 +48,9 @@ public void accept(CdcEvent cdcEvent)
@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)
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -99,6 +100,7 @@ public CdcManager(EventConsumer eventConsumer,
ClusterConfigProvider clusterConfigProvider,
SidecarCdcClient sidecarCdcClient,
ICdcStats cdcStats,
SidecarCdcStats sidecarCdcStats,
TaskExecutorPool taskExecutorPool,
CdcDatabaseAccessor cdcDatabaseAccessor,
CdcOptions cdcOptions)
Expand All @@ -111,6 +113,7 @@ public CdcManager(EventConsumer eventConsumer,
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);
Expand Down Expand Up @@ -216,14 +219,14 @@ CdcConsumerEntry buildConsumer(@NotNull String jobId,
.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);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -38,6 +38,7 @@
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;
Expand Down Expand Up @@ -120,13 +121,18 @@ public CdcPublisher(Vertx vertx,

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);
vertx.eventBus().localConsumer(ON_SERVER_STOP.address(), this);
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)
Expand All @@ -146,7 +152,7 @@ public EventConsumer eventConsumer(CdcConfig conf)
conf.failOnRecordTooLargeError(),
conf.failOnKafkaError(),
CdcLogMode.FULL);
return new CdcEventConsumer(kafkaPublisher);
return new CdcEventConsumer(kafkaPublisher, sidecarCdcStats);
}

/**
Expand Down Expand Up @@ -220,6 +226,7 @@ private synchronized void run() throws IllegalStateException
clusterConfigProvider,
sidecarCdcClientProvider.get(),
cdcStats,
sidecarCdcStats,
this.executorPools,
databaseAccessor,
cdcOptions);
Expand Down
Loading
Loading