From 026a3e0ecade91692b136e7bf5a1f4d8d3cb847c Mon Sep 17 00:00:00 2001 From: Jan Korona Date: Tue, 21 Jul 2026 15:00:25 +0200 Subject: [PATCH 1/4] Add Cassandra compaction progress metrics via JMX handler --- CHANGELOG.md | 1 + .../jmx-metrics/library/cassandra.md | 26 +-- .../CassandraCompactionProgressHandler.java | 164 +++++++++++++++++ ....jmx.internal.ExperimentalJmxMetricHandler | 1 + .../jmx/rules/experimental-cassandra.yaml | 6 + ...assandraCompactionProgressHandlerTest.java | 169 ++++++++++++++++++ .../jmx/rules/CassandraTest.java | 132 +++++++++++++- 7 files changed, 486 insertions(+), 13 deletions(-) create mode 100644 instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java create mode 100644 instrumentation/jmx-metrics/library/src/main/resources/META-INF/services/io.opentelemetry.instrumentation.jmx.internal.ExperimentalJmxMetricHandler create mode 100644 instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java diff --git a/CHANGELOG.md b/CHANGELOG.md index d2242c56c877..3f47f889dee8 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,6 +18,7 @@ - Add Cassandra JMX metrics target system. ([#19080](https://github.com/open-telemetry/opentelemetry-java-instrumentation/pull/19080)) +- Add Cassandra compaction byte-progress metrics via a code-based JMX handler. (#0) - Add `captureTemplate` and `captureArguments` options to the log4j, java-util-logging, and jboss-logmanager logging instrumentations, capturing the log message template and arguments as separate `log.body.template` / `log.body.parameters` attributes. This extends the same option diff --git a/instrumentation/jmx-metrics/library/cassandra.md b/instrumentation/jmx-metrics/library/cassandra.md index a4eefe9e59e3..5b33740b7fd7 100644 --- a/instrumentation/jmx-metrics/library/cassandra.md +++ b/instrumentation/jmx-metrics/library/cassandra.md @@ -2,15 +2,17 @@ Here is the list of metrics based on MBeans exposed by Cassandra. -| Metric Name | Type | Unit | Attributes | Description | -| ------------------------------------ | ------------- | --------- | ------------------------------------- | ---------------------------------------------------------------- | -| cassandra.client.request.count | Counter | {request} | cassandra.operation | Number of requests by operation. | -| cassandra.client.request.error | Counter | {error} | cassandra.operation, cassandra.status | Number of request errors by operation. | -| cassandra.client.request.latency.p50 | Gauge | s | cassandra.operation | Request latency 50th percentile by operation. | -| cassandra.client.request.latency.p99 | Gauge | s | cassandra.operation | Request latency 99th percentile by operation. | -| cassandra.client.request.latency.max | Gauge | s | cassandra.operation | Maximum request latency by operation. | -| cassandra.compaction.tasks.completed | Counter | {task} | | Number of completed compactions since server start. | -| cassandra.compaction.tasks.pending | Gauge | {task} | | Estimated number of compactions remaining to perform. | -| cassandra.storage.load | UpDownCounter | By | | Size of the on disk data size this node manages. | -| cassandra.storage.hints.count | Counter | {hint} | | Number of hint messages written to this node since server start. | -| cassandra.storage.hints.in_progress | UpDownCounter | {hint} | | Number of hints attempting to be sent currently. | +| Metric Name | Type | Unit | Attributes | Description | +|--------------------------------------|---------------|-----------|---------------------------------------------------------------------|------------------------------------------------------------------| +| cassandra.client.request.count | Counter | {request} | cassandra.operation | Number of requests by operation. | +| cassandra.client.request.error | Counter | {error} | cassandra.operation, cassandra.status | Number of request errors by operation. | +| cassandra.client.request.latency.p50 | Gauge | s | cassandra.operation | Request latency 50th percentile by operation. | +| cassandra.client.request.latency.p99 | Gauge | s | cassandra.operation | Request latency 99th percentile by operation. | +| cassandra.client.request.latency.max | Gauge | s | cassandra.operation | Maximum request latency by operation. | +| cassandra.compaction.progress.bytes | Gauge | By | cassandra.compaction.task_type, cassandra.keyspace, cassandra.table | Bytes completed for in-flight compactions. | +| cassandra.compaction.progress.total | Gauge | By | cassandra.compaction.task_type, cassandra.keyspace, cassandra.table | Total bytes for in-flight compactions. | +| cassandra.compaction.tasks.completed | Counter | {task} | | Number of completed compactions since server start. | +| cassandra.compaction.tasks.pending | Gauge | {task} | | Estimated number of compactions remaining to perform. | +| cassandra.storage.load | UpDownCounter | By | | Size of the on disk data size this node manages. | +| cassandra.storage.hints.count | Counter | {hint} | | Number of hint messages written to this node since server start. | +| cassandra.storage.hints.in_progress | UpDownCounter | {hint} | | Number of hints attempting to be sent currently. | diff --git a/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java b/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java new file mode 100644 index 000000000000..583b9b17d7be --- /dev/null +++ b/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java @@ -0,0 +1,164 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.instrumentation.jmx.internal.handler; + +import static java.util.logging.Level.WARNING; + +import io.opentelemetry.api.common.Attributes; +import io.opentelemetry.api.common.AttributesBuilder; +import io.opentelemetry.api.metrics.Meter; +import io.opentelemetry.api.metrics.ObservableLongMeasurement; +import io.opentelemetry.instrumentation.jmx.internal.ExperimentalJmxMetricHandler; +import java.math.BigInteger; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.function.Supplier; +import java.util.logging.Logger; +import javax.annotation.Nullable; +import javax.management.MBeanServerConnection; +import javax.management.ObjectName; + +/** + * JMX metric handler that reports per-compaction byte progress from Cassandra's CompactionManager. + * + *

Queries {@code org.apache.cassandra.db:type=CompactionManager}, reads the {@code Compactions} + * attribute (a list of maps representing in-flight compactions), groups entries by {@code taskType + * + keyspace + columnfamily}, and emits two gauges per group: + * + *

+ * + *

Both metrics carry {@code cassandra.compaction.task_type}, {@code cassandra.keyspace}, and + * {@code cassandra.table} attributes. Entries missing any of these dimension fields are skipped. + * Byte values are string-encoded in the MBean and parsed via {@link BigInteger}. + * + *

Note: these are distinct from {@code cassandra.compaction.tasks.pending/completed}, which are + * simple scalar task counts from {@code org.apache.cassandra.metrics:type=Compaction}. + * + *

This class is internal and is hence not for public use. Its APIs are unstable and can change + * at any time. + */ +public final class CassandraCompactionProgressHandler implements ExperimentalJmxMetricHandler { + + static final String HANDLER_NAME = "cassandra-compaction-progress"; + static final String METRIC_CURRENT = "cassandra.compaction.progress.bytes"; + static final String METRIC_TOTAL = "cassandra.compaction.progress.total"; + + private static final String ATTR_TASK_TYPE = "cassandra.compaction.task_type"; + private static final String ATTR_KEYSPACE = "cassandra.keyspace"; + private static final String ATTR_TABLE = "cassandra.table"; + + private static final Logger logger = + Logger.getLogger(CassandraCompactionProgressHandler.class.getName()); + + @Override + public String getName() { + return HANDLER_NAME; + } + + @Override + public AutoCloseable create(Meter meter, Supplier detectorSupplier) { + ObservableLongMeasurement currentGauge = + meter + .gaugeBuilder(METRIC_CURRENT) + .setDescription("Bytes completed for in-flight compactions") + .setUnit("By") + .ofLongs() + .buildObserver(); + + ObservableLongMeasurement totalGauge = + meter + .gaugeBuilder(METRIC_TOTAL) + .setDescription("Total bytes for in-flight compactions") + .setUnit("By") + .ofLongs() + .buildObserver(); + + return meter.batchCallback( + () -> { + for (Map.Entry entry : queryGroups(detectorSupplier).entrySet()) { + currentGauge.record(entry.getValue()[0], entry.getKey()); + totalGauge.record(entry.getValue()[1], entry.getKey()); + } + }, + currentGauge, + totalGauge); + } + + private static Map queryGroups(Supplier detectorSupplier) { + Map groups = new HashMap<>(); + Detector detector = detectorSupplier.get(); + if (detector == null) { + return groups; + } + MBeanServerConnection connection = detector.getConnection(); + for (ObjectName objectName : detector.getObjectNames()) { + queryCompactions(connection, objectName) + .forEach( + (attrs, values) -> + groups.merge(attrs, values, (a, b) -> new long[] {a[0] + b[0], a[1] + b[1]})); + } + return groups; + } + + static Map queryCompactions( + MBeanServerConnection connection, ObjectName objectName) { + Map groups = new HashMap<>(); + try { + // CompactionManager#getCompactions() returns List> via Standard MBean + @SuppressWarnings("unchecked") + List> compactions = + (List>) connection.getAttribute(objectName, "Compactions"); + + for (Map entry : compactions) { + String taskType = entry.get("taskType"); + String keyspace = entry.get("keyspace"); + String columnFamily = entry.get("columnfamily"); + String unit = entry.get("unit"); + + if (taskType == null || keyspace == null || columnFamily == null || !isByteUnit(unit)) { + continue; + } + + long completed = parseLong(entry.get("completed")); + long total = parseLong(entry.get("total")); + + Attributes attrs = buildAttributes(taskType, keyspace, columnFamily); + groups.merge( + attrs, new long[] {completed, total}, (a, b) -> new long[] {a[0] + b[0], a[1] + b[1]}); + } + } catch (Exception e) { + logger.log(WARNING, "cassandra.compaction.progress: failed to query CompactionManager", e); + } + return groups; + } + + private static boolean isByteUnit(@Nullable String unit) { + return unit != null && unit.equalsIgnoreCase("bytes"); + } + + private static long parseLong(@Nullable Object value) { + if (value == null) { + return 0; + } + try { + return new BigInteger(value.toString()).longValue(); + } catch (NumberFormatException e) { + return 0; + } + } + + private static Attributes buildAttributes(String taskType, String keyspace, String columnFamily) { + AttributesBuilder builder = Attributes.builder(); + builder.put(ATTR_TASK_TYPE, taskType); + builder.put(ATTR_KEYSPACE, keyspace); + builder.put(ATTR_TABLE, columnFamily); + return builder.build(); + } +} diff --git a/instrumentation/jmx-metrics/library/src/main/resources/META-INF/services/io.opentelemetry.instrumentation.jmx.internal.ExperimentalJmxMetricHandler b/instrumentation/jmx-metrics/library/src/main/resources/META-INF/services/io.opentelemetry.instrumentation.jmx.internal.ExperimentalJmxMetricHandler new file mode 100644 index 000000000000..82767d1f24ca --- /dev/null +++ b/instrumentation/jmx-metrics/library/src/main/resources/META-INF/services/io.opentelemetry.instrumentation.jmx.internal.ExperimentalJmxMetricHandler @@ -0,0 +1 @@ +io.opentelemetry.instrumentation.jmx.internal.handler.CassandraCompactionProgressHandler diff --git a/instrumentation/jmx-metrics/library/src/main/resources/jmx/rules/experimental-cassandra.yaml b/instrumentation/jmx-metrics/library/src/main/resources/jmx/rules/experimental-cassandra.yaml index d799f1e2a42d..f2968f926e89 100644 --- a/instrumentation/jmx-metrics/library/src/main/resources/jmx/rules/experimental-cassandra.yaml +++ b/instrumentation/jmx-metrics/library/src/main/resources/jmx/rules/experimental-cassandra.yaml @@ -120,3 +120,9 @@ rules: type: *errortype unit: *errorunit desc: *errordesc + + # Compaction byte progress - uses a code-based handler because the Compactions attribute + # returns a list of maps requiring iteration, grouping by composite key, and string-encoded + # BigInteger parsing, none of which can be expressed in declarative YAML. + - bean: org.apache.cassandra.db:type=CompactionManager + handler: cassandra-compaction-progress diff --git a/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java b/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java new file mode 100644 index 000000000000..a2dc39cd6eba --- /dev/null +++ b/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java @@ -0,0 +1,169 @@ +/* + * Copyright The OpenTelemetry Authors + * SPDX-License-Identifier: Apache-2.0 + */ + +package io.opentelemetry.instrumentation.jmx.internal.handler; + +import static java.util.Arrays.asList; +import static java.util.Collections.singletonList; +import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +import io.opentelemetry.api.common.Attributes; +import java.util.HashMap; +import java.util.Map; +import javax.management.MBeanServerConnection; +import javax.management.ObjectName; +import org.junit.jupiter.api.BeforeEach; +import org.junit.jupiter.api.Test; + +class CassandraCompactionProgressHandlerTest { + + private static final String ATTR_TASK_TYPE = "cassandra.compaction.task_type"; + private static final String ATTR_KEYSPACE = "cassandra.keyspace"; + private static final String ATTR_TABLE = "cassandra.table"; + + private MBeanServerConnection connection; + private ObjectName objectName; + + @BeforeEach + void setUp() throws Exception { + connection = mock(MBeanServerConnection.class); + objectName = new ObjectName("org.apache.cassandra.db:type=CompactionManager"); + } + + @Test + void groupsByCompositeKey() throws Exception { + when(connection.getAttribute(objectName, "Compactions")) + .thenReturn( + asList( + compactionEntry("COMPACTION", "ks1", "cf1", "100", "200"), + compactionEntry("COMPACTION", "ks1", "cf1", "50", "150"), + compactionEntry("COMPACTION", "ks2", "cf2", "10", "100"))); + + Map groups = + CassandraCompactionProgressHandler.queryCompactions(connection, objectName); + + assertThat(groups).hasSize(2); + assertThat(groups.get(attrs("COMPACTION", "ks1", "cf1"))).containsExactly(150L, 350L); + assertThat(groups.get(attrs("COMPACTION", "ks2", "cf2"))).containsExactly(10L, 100L); + } + + @Test + void skipsEntriesMissingDimensionFields() throws Exception { + when(connection.getAttribute(objectName, "Compactions")) + .thenReturn( + asList( + compactionEntry("COMPACTION", "ks1", "cf1", "10", "100"), + compactionEntry("COMPACTION", null, "cf1", "5", "50"), + compactionEntry("COMPACTION", "ks1", null, "5", "50"))); + + Map groups = + CassandraCompactionProgressHandler.queryCompactions(connection, objectName); + + assertThat(groups).hasSize(1); + assertThat(groups.get(attrs("COMPACTION", "ks1", "cf1"))).containsExactly(10L, 100L); + } + + @Test + void skipsEntriesWithNonByteUnits() throws Exception { + when(connection.getAttribute(objectName, "Compactions")) + .thenReturn( + asList( + compactionEntry("COMPACTION", "ks1", "cf1", "10", "100", "bytes"), + compactionEntry("VALIDATION", "ks1", "cf1", "5", "50", "keys"), + compactionEntry("ANTICOMPACTION", "ks1", "cf1", "5", "50", "ranges"), + compactionEntry("COMPACTION", "ks2", "cf2", "5", "50", null))); + + Map groups = + CassandraCompactionProgressHandler.queryCompactions(connection, objectName); + + assertThat(groups).hasSize(1); + assertThat(groups.get(attrs("COMPACTION", "ks1", "cf1"))).containsExactly(10L, 100L); + } + + @Test + void parsesBigIntegerStringValues() throws Exception { + // values larger than Long.MAX_VALUE are truncated but must not throw + String big = "99999999999999999999"; + when(connection.getAttribute(objectName, "Compactions")) + .thenReturn(singletonList(compactionEntry("COMPACTION", "ks", "cf", big, big))); + + Map groups = + CassandraCompactionProgressHandler.queryCompactions(connection, objectName); + + assertThat(groups).hasSize(1); + } + + @Test + void returnsEmptyMapOnException() throws Exception { + when(connection.getAttribute(objectName, "Compactions")) + .thenThrow(new RuntimeException("connection lost")); + + Map groups = + CassandraCompactionProgressHandler.queryCompactions(connection, objectName); + + assertThat(groups).isEmpty(); + } + + @Test + void handlerNameIsStable() { + assertThat(new CassandraCompactionProgressHandler().getName()) + .isEqualTo(CassandraCompactionProgressHandler.HANDLER_NAME); + } + + @Test + void mergesGroupsAcrossMultipleObjectNames() throws Exception { + ObjectName objectName2 = new ObjectName("org.apache.cassandra.db:type=CompactionManager,id=2"); + when(connection.getAttribute(objectName, "Compactions")) + .thenReturn(singletonList(compactionEntry("COMPACTION", "ks1", "cf1", "100", "200"))); + when(connection.getAttribute(objectName2, "Compactions")) + .thenReturn(singletonList(compactionEntry("COMPACTION", "ks1", "cf1", "50", "150"))); + + Map first = + CassandraCompactionProgressHandler.queryCompactions(connection, objectName); + Map second = + CassandraCompactionProgressHandler.queryCompactions(connection, objectName2); + + // Simulate what queryGroups does: merge across ObjectNames using the same merge function + Map merged = new HashMap<>(first); + second.forEach( + (attrs, values) -> + merged.merge(attrs, values, (a, b) -> new long[] {a[0] + b[0], a[1] + b[1]})); + + assertThat(merged).hasSize(1); + assertThat(merged.get(attrs("COMPACTION", "ks1", "cf1"))).containsExactly(150L, 350L); + } + + private static Map compactionEntry( + String taskType, String keyspace, String columnfamily, String completed, String total) { + return compactionEntry(taskType, keyspace, columnfamily, completed, total, "bytes"); + } + + private static Map compactionEntry( + String taskType, + String keyspace, + String columnfamily, + String completed, + String total, + String unit) { + Map entry = new HashMap<>(); + entry.put("taskType", taskType); + entry.put("keyspace", keyspace); + entry.put("columnfamily", columnfamily); + entry.put("completed", completed); + entry.put("total", total); + entry.put("unit", unit); + return entry; + } + + private static Attributes attrs(String taskType, String keyspace, String columnFamily) { + return Attributes.builder() + .put(ATTR_TASK_TYPE, taskType) + .put(ATTR_KEYSPACE, keyspace) + .put(ATTR_TABLE, columnFamily) + .build(); + } +} diff --git a/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/rules/CassandraTest.java b/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/rules/CassandraTest.java index e3c1c56ac1bd..26ad921deb70 100644 --- a/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/rules/CassandraTest.java +++ b/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/rules/CassandraTest.java @@ -7,13 +7,17 @@ import static io.opentelemetry.instrumentation.jmx.rules.assertions.DataPointAttributes.attribute; import static io.opentelemetry.instrumentation.jmx.rules.assertions.DataPointAttributes.attributeGroup; +import static io.opentelemetry.instrumentation.jmx.rules.assertions.DataPointAttributes.attributeWithAnyValue; import static java.util.Collections.singletonList; import io.opentelemetry.instrumentation.jmx.rules.assertions.AttributeMatcherGroup; +import java.io.IOException; import java.time.Duration; import java.util.ArrayList; import java.util.List; +import java.util.logging.Logger; import org.junit.jupiter.api.Test; +import org.testcontainers.containers.Container; import org.testcontainers.containers.GenericContainer; import org.testcontainers.containers.wait.strategy.Wait; @@ -21,6 +25,14 @@ class CassandraTest extends TargetSystemTest { private static final int CASSANDRA_PORT = 9042; + // 4 batches x 2500 rows x 2 KB = ~20 MB; at 1 MB/s throttle compaction stays visible for ~20s, + // well within the 60s await in verifyMetrics(). + private static final int COMPACTION_BATCHES = 4; + private static final int COMPACTION_ROWS_PER_BATCH = 2_500; + private static final int COMPACTION_VALUE_BYTES = 2_048; + + private static final Logger logger = Logger.getLogger(CassandraTest.class.getName()); + @Test void testCassandraMetrics() { List yamlFiles = singletonList("experimental-cassandra.yaml"); @@ -47,9 +59,103 @@ void testCassandraMetrics() { startTarget(target); + seedCompactionData(target); + triggerCompaction(target); + verifyMetrics(createMetricsVerifier()); } + private static void seedCompactionData(GenericContainer target) { + try { + // Throttle compaction to 1 MB/s so it stays active long enough to be observed. + nodetool(target, "setcompactionthroughput", "1"); + + execOrThrow( + target, + "cqlsh", + "-e", + "CREATE KEYSPACE IF NOT EXISTS test" + + " WITH replication = {'class':'SimpleStrategy','replication_factor':1};" + + "CREATE TABLE IF NOT EXISTS test.data (id uuid PRIMARY KEY, val text)" + + " WITH compression = {'enabled':'false'};"); + nodetool(target, "disableautocompaction", "test", "data"); + + for (int ignored = 0; ignored < COMPACTION_BATCHES; ignored++) { + seedCompactionBatch(target); + nodetool(target, "flush", "test", "data"); + } + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + throw new IllegalStateException("Failed to seed Cassandra data", e); + } catch (Exception e) { + throw new IllegalStateException("Failed to seed Cassandra data", e); + } + } + + private static void seedCompactionBatch(GenericContainer target) + throws IOException, InterruptedException { + execOrThrow( + target, + "bash", + "-c", + // value is generated once and reused for all rows - only byte volume matters here. + "set -euo pipefail; " + + "value=$(head -c " + + COMPACTION_VALUE_BYTES + + " /dev/zero | tr '\\0' x); " + + "rm -f /tmp/cassandra-data.csv; " + + "for i in $(seq 1 " + + COMPACTION_ROWS_PER_BATCH + + "); do " + + "printf \"%s,%s\\n\" \"$(cat /proc/sys/kernel/random/uuid)\" \"$value\"; " + + "done > /tmp/cassandra-data.csv; " + + "cqlsh -e \"COPY test.data (id, val) FROM '/tmp/cassandra-data.csv' " + + "WITH HEADER = false AND MINBATCHSIZE = 1 AND MAXBATCHSIZE = 2;\""); + } + + private static void triggerCompaction(GenericContainer target) { + Thread compactionThread = + new Thread( + () -> { + try { + nodetool(target, "compact", "--", "test", "data"); + } catch (InterruptedException e) { + Thread.currentThread().interrupt(); + logger.warning("Background compaction interrupted"); + } catch (Exception e) { + if (target.isRunning()) { + logger.warning("Background compaction failed: " + e.getMessage()); + } + } + }, + "cassandra-compaction-trigger"); + compactionThread.setDaemon(true); + compactionThread.start(); + } + + private static void nodetool(GenericContainer target, String... args) + throws IOException, InterruptedException { + String[] command = new String[args.length + 1]; + command[0] = "nodetool"; + System.arraycopy(args, 0, command, 1, args.length); + execOrThrow(target, command); + } + + private static void execOrThrow(GenericContainer target, String... command) + throws IOException, InterruptedException { + Container.ExecResult result = target.execInContainer(command); + if (result.getExitCode() != 0) { + throw new IllegalStateException( + String.join(" ", command) + + " failed with exit code " + + result.getExitCode() + + "\nstdout:\n" + + result.getStdout() + + "\nstderr:\n" + + result.getStderr()); + } + } + private static MetricsVerifier createMetricsVerifier() { return MetricsVerifier.create() .add( @@ -153,7 +259,31 @@ private static MetricsVerifier createMetricsVerifier() { errorAttributesGroup("read", "unavailable"), errorAttributesGroup("write", "timeout"), errorAttributesGroup("write", "failure"), - errorAttributesGroup("write", "unavailable"))); + errorAttributesGroup("write", "unavailable"))) + .add( + "cassandra.compaction.progress.bytes", + metric -> + metric + .hasDescription("Bytes completed for in-flight compactions") + .hasUnit("By") + .isGauge() + .hasDataPointsWithAttributes( + attributeGroup( + attributeWithAnyValue("cassandra.compaction.task_type"), + attribute("cassandra.keyspace", "test"), + attribute("cassandra.table", "data")))) + .add( + "cassandra.compaction.progress.total", + metric -> + metric + .hasDescription("Total bytes for in-flight compactions") + .hasUnit("By") + .isGauge() + .hasDataPointsWithAttributes( + attributeGroup( + attributeWithAnyValue("cassandra.compaction.task_type"), + attribute("cassandra.keyspace", "test"), + attribute("cassandra.table", "data")))); } private static AttributeMatcherGroup errorAttributesGroup(String operation, String status) { From a29b1b1003c381025973a24ed2c796f5bee0d870 Mon Sep 17 00:00:00 2001 From: Jan Korona Date: Tue, 21 Jul 2026 16:19:27 +0200 Subject: [PATCH 2/4] Apply suggestions from code review Co-authored-by: Copilot Autofix powered by AI <175728472+Copilot@users.noreply.github.com> --- CHANGELOG.md | 1 - .../internal/handler/CassandraCompactionProgressHandler.java | 2 +- 2 files changed, 1 insertion(+), 2 deletions(-) diff --git a/CHANGELOG.md b/CHANGELOG.md index 3f47f889dee8..d2242c56c877 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -18,7 +18,6 @@ - Add Cassandra JMX metrics target system. ([#19080](https://github.com/open-telemetry/opentelemetry-java-instrumentation/pull/19080)) -- Add Cassandra compaction byte-progress metrics via a code-based JMX handler. (#0) - Add `captureTemplate` and `captureArguments` options to the log4j, java-util-logging, and jboss-logmanager logging instrumentations, capturing the log message template and arguments as separate `log.body.template` / `log.body.parameters` attributes. This extends the same option diff --git a/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java b/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java index 583b9b17d7be..aeca720e5997 100644 --- a/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java +++ b/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java @@ -44,7 +44,7 @@ *

This class is internal and is hence not for public use. Its APIs are unstable and can change * at any time. */ -public final class CassandraCompactionProgressHandler implements ExperimentalJmxMetricHandler { +public class CassandraCompactionProgressHandler implements ExperimentalJmxMetricHandler { static final String HANDLER_NAME = "cassandra-compaction-progress"; static final String METRIC_CURRENT = "cassandra.compaction.progress.bytes"; From 9db27f9f3702c4ccfdbc0fb1f43025ab11e043e0 Mon Sep 17 00:00:00 2001 From: Jan Korona Date: Wed, 22 Jul 2026 13:06:35 +0200 Subject: [PATCH 3/4] throw ArithmeticException instead of silently truncating values that exceed Long.MAX_VALUE --- .../CassandraCompactionProgressHandler.java | 16 +++++++++++++--- .../CassandraCompactionProgressHandlerTest.java | 6 +++--- 2 files changed, 16 insertions(+), 6 deletions(-) diff --git a/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java b/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java index aeca720e5997..cd2b0a43fa2e 100644 --- a/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java +++ b/instrumentation/jmx-metrics/library/src/main/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandler.java @@ -126,8 +126,18 @@ static Map queryCompactions( continue; } - long completed = parseLong(entry.get("completed")); - long total = parseLong(entry.get("total")); + long completed; + long total; + try { + completed = parseLong(entry.get("completed")); + total = parseLong(entry.get("total")); + } catch (ArithmeticException e) { + logger.log( + WARNING, + "cassandra.compaction.progress: byte value overflows long range, skipping entry", + e); + continue; + } Attributes attrs = buildAttributes(taskType, keyspace, columnFamily); groups.merge( @@ -148,7 +158,7 @@ private static long parseLong(@Nullable Object value) { return 0; } try { - return new BigInteger(value.toString()).longValue(); + return new BigInteger(value.toString()).longValueExact(); } catch (NumberFormatException e) { return 0; } diff --git a/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java b/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java index a2dc39cd6eba..8f87dd5d9926 100644 --- a/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java +++ b/instrumentation/jmx-metrics/library/src/test/java/io/opentelemetry/instrumentation/jmx/internal/handler/CassandraCompactionProgressHandlerTest.java @@ -85,8 +85,8 @@ void skipsEntriesWithNonByteUnits() throws Exception { } @Test - void parsesBigIntegerStringValues() throws Exception { - // values larger than Long.MAX_VALUE are truncated but must not throw + void skipsEntriesWithValuesExceedingLongRange() throws Exception { + // values larger than Long.MAX_VALUE cannot be safely cast to long — entry must be skipped String big = "99999999999999999999"; when(connection.getAttribute(objectName, "Compactions")) .thenReturn(singletonList(compactionEntry("COMPACTION", "ks", "cf", big, big))); @@ -94,7 +94,7 @@ void parsesBigIntegerStringValues() throws Exception { Map groups = CassandraCompactionProgressHandler.queryCompactions(connection, objectName); - assertThat(groups).hasSize(1); + assertThat(groups).isEmpty(); } @Test From 17c185c89985e1ef3091411c4fe5c10c5c806a06 Mon Sep 17 00:00:00 2001 From: Jan Korona Date: Wed, 22 Jul 2026 13:06:42 +0200 Subject: [PATCH 4/4] reformat --- instrumentation/jmx-metrics/library/cassandra.md | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/instrumentation/jmx-metrics/library/cassandra.md b/instrumentation/jmx-metrics/library/cassandra.md index 5b33740b7fd7..fc56e79b35db 100644 --- a/instrumentation/jmx-metrics/library/cassandra.md +++ b/instrumentation/jmx-metrics/library/cassandra.md @@ -3,7 +3,7 @@ Here is the list of metrics based on MBeans exposed by Cassandra. | Metric Name | Type | Unit | Attributes | Description | -|--------------------------------------|---------------|-----------|---------------------------------------------------------------------|------------------------------------------------------------------| +| ------------------------------------ | ------------- | --------- | ------------------------------------------------------------------- | ---------------------------------------------------------------- | | cassandra.client.request.count | Counter | {request} | cassandra.operation | Number of requests by operation. | | cassandra.client.request.error | Counter | {error} | cassandra.operation, cassandra.status | Number of request errors by operation. | | cassandra.client.request.latency.p50 | Gauge | s | cassandra.operation | Request latency 50th percentile by operation. |