This is an automated email from the ASF dual-hosted git repository.
kfaraz pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/druid.git
The following commit(s) were added to refs/heads/master by this push:
new 8655bd6d0d2 minor: Emit taskType on task count metrics from the statsd
and prometheus emitters (#20123)
8655bd6d0d2 is described below
commit 8655bd6d0d2908671591487c87f60b8552d57ebd
Author: Andreas Maechler <[email protected]>
AuthorDate: Sun Aug 23 23:12:20 2026 -0600
minor: Emit taskType on task count metrics from the statsd and prometheus
emitters (#20123)
Emit the `taskType` dimension for the following metrics with statsd and
prometheus emitters:
- `task/running/count`
- `task/waiting/count`
- `task/pending/count`
- `task/success/count`
- `task/failed/count`
---
.../src/main/resources/defaultMetrics.json | 10 ++---
.../druid/emitter/prometheus/MetricsTest.java | 22 +++++++++++
.../main/resources/defaultMetricDimensions.json | 10 ++---
.../emitter/statsd/DimensionConverterTest.java | 44 ++++++++++++++++++++++
4 files changed, 76 insertions(+), 10 deletions(-)
diff --git
a/extensions-contrib/prometheus-emitter/src/main/resources/defaultMetrics.json
b/extensions-contrib/prometheus-emitter/src/main/resources/defaultMetrics.json
index cd7b609bb58..a860680d7f4 100644
---
a/extensions-contrib/prometheus-emitter/src/main/resources/defaultMetrics.json
+++
b/extensions-contrib/prometheus-emitter/src/main/resources/defaultMetrics.json
@@ -155,11 +155,11 @@
"segment/added/bytes" : { "dimensions" : ["dataSource", "taskType"], "type"
: "count", "help": "Size in bytes of new segments created." },
"segment/moved/bytes" : { "dimensions" : ["dataSource", "taskType"], "type"
: "count", "help": "Size in bytes of segments moved/archived via the Move
Task." },
"segment/nuked/bytes" : { "dimensions" : ["dataSource", "taskType"], "type"
: "count", "help": "Size in bytes of segments deleted via the Kill Task." },
- "task/success/count" : { "dimensions" : ["dataSource"], "type" : "count",
"help": "Number of successful tasks per emission period."},
- "task/failed/count" : { "dimensions" : ["dataSource"], "type" : "count",
"help": "Number of failed tasks per emission period."},
- "task/running/count" : { "dimensions" : ["dataSource"], "type" : "count",
"help": "Number of current running tasks."},
- "task/pending/count" : { "dimensions" : ["dataSource"], "type" : "count",
"help": "Number of current pending tasks."},
- "task/waiting/count" : { "dimensions" : ["dataSource"], "type" : "count",
"help": "Number of current waiting tasks."},
+ "task/success/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"count", "help": "Number of successful tasks per emission period."},
+ "task/failed/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"count", "help": "Number of failed tasks per emission period."},
+ "task/running/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"count", "help": "Number of current running tasks."},
+ "task/pending/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"count", "help": "Number of current pending tasks."},
+ "task/waiting/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"count", "help": "Number of current waiting tasks."},
"supervisor/count" : { "dimensions" : ["supervisorId", "type", "state",
"detailedState"], "type" : "gauge", "help": "Count of active supervisors. Each
supervisor emits 1, tagged with its state. Available only if the
SupervisorStatsMonitor module is included."},
"segment/assigned/count" : { "dimensions" : ["tier"], "type" : "count",
"help": "Number of segments assigned to be loaded in the cluster."},
diff --git
a/extensions-contrib/prometheus-emitter/src/test/java/org/apache/druid/emitter/prometheus/MetricsTest.java
b/extensions-contrib/prometheus-emitter/src/test/java/org/apache/druid/emitter/prometheus/MetricsTest.java
index b354d83f3ad..e4567c9aaf6 100644
---
a/extensions-contrib/prometheus-emitter/src/test/java/org/apache/druid/emitter/prometheus/MetricsTest.java
+++
b/extensions-contrib/prometheus-emitter/src/test/java/org/apache/druid/emitter/prometheus/MetricsTest.java
@@ -138,4 +138,26 @@ public class MetricsTest
Assertions.assertArrayEquals(expectedHistogramBuckets,
dimensionsAndCollector.getHistogramBuckets(), 0.0);
}
+ @Test
+ public void testTaskCountMetricsHaveTaskTypeLabel()
+ {
+ PrometheusEmitterConfig config = new PrometheusEmitterConfig(null,
"test_7", null, null, null, true, true, null, null, null, null);
+ Metrics metrics = new Metrics(config);
+ for (String metric : new String[]{
+ "task/success/count",
+ "task/failed/count",
+ "task/running/count",
+ "task/pending/count",
+ "task/waiting/count"
+ }) {
+ DimensionsAndCollector dimensionsAndCollector =
metrics.getByName(metric, "overlord");
+ Assertions.assertNotNull(dimensionsAndCollector, metric);
+ Assertions.assertArrayEquals(
+ new String[]{"dataSource", "druid_service", "host_name", "taskType"},
+ dimensionsAndCollector.getDimensions(),
+ metric
+ );
+ }
+ }
+
}
diff --git
a/extensions-contrib/statsd-emitter/src/main/resources/defaultMetricDimensions.json
b/extensions-contrib/statsd-emitter/src/main/resources/defaultMetricDimensions.json
index 226ec036ec7..f1eb96a40a5 100644
---
a/extensions-contrib/statsd-emitter/src/main/resources/defaultMetricDimensions.json
+++
b/extensions-contrib/statsd-emitter/src/main/resources/defaultMetricDimensions.json
@@ -79,11 +79,11 @@
"ingest/pause/time" : { "dimensions" : ["dataSource", "taskId"], "type" :
"timer" },
- "task/success/count" : { "dimensions" : ["dataSource"], "type" : "count" },
- "task/failed/count" : { "dimensions" : ["dataSource"], "type" : "count" },
- "task/running/count" : { "dimensions" : ["dataSource"], "type" : "gauge" },
- "task/pending/count" : { "dimensions" : ["dataSource"], "type" : "gauge" },
- "task/waiting/count" : { "dimensions" : ["dataSource"], "type" : "gauge" },
+ "task/success/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"count" },
+ "task/failed/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"count" },
+ "task/running/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"gauge" },
+ "task/pending/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"gauge" },
+ "task/waiting/count" : { "dimensions" : ["dataSource", "taskType"], "type" :
"gauge" },
"task/action/run/time": { "dimensions" : ["dataSource", "taskActionType"],
"type" : "timer" },
"task/status/queue/count": { "dimensions" : [], "type" : "gauge" },
diff --git
a/extensions-contrib/statsd-emitter/src/test/java/org/apache/druid/emitter/statsd/DimensionConverterTest.java
b/extensions-contrib/statsd-emitter/src/test/java/org/apache/druid/emitter/statsd/DimensionConverterTest.java
index 0ae7b95ade5..3f43c011f25 100644
---
a/extensions-contrib/statsd-emitter/src/test/java/org/apache/druid/emitter/statsd/DimensionConverterTest.java
+++
b/extensions-contrib/statsd-emitter/src/test/java/org/apache/druid/emitter/statsd/DimensionConverterTest.java
@@ -25,6 +25,8 @@ import
org.apache.druid.java.util.emitter.service.ServiceMetricEvent;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import java.util.List;
+
public class DimensionConverterTest
{
@Test
@@ -58,4 +60,46 @@ public class DimensionConverterTest
expected.put("type", "groupBy");
Assertions.assertEquals(expected.build(), actual.build(), "correct
Dimensions");
}
+
+ @Test
+ public void testConvertTaskCountMetrics()
+ {
+ DimensionConverter dimensionConverter = new DimensionConverter(new
ObjectMapper(), null);
+ for (String metric : new String[]{
+ "task/success/count",
+ "task/failed/count",
+ "task/running/count",
+ "task/pending/count",
+ "task/waiting/count"
+ }) {
+ ServiceMetricEvent event = new ServiceMetricEvent.Builder()
+ .setDimension("dataSource", "data-source")
+ .setDimension("taskType", "index_kafka")
+ .setDimension("supervisorId", "supervisor-1")
+ .setMetric(metric, 1)
+ .build("overlord", "overlordHost1");
+
+ ImmutableMap.Builder<String, String> actual = new
ImmutableMap.Builder<>();
+ StatsDMetric statsDMetric = dimensionConverter.addFilteredUserDims(
+ event.getService(),
+ event.getMetric(),
+ event.getUserDims(),
+ actual
+ );
+ Assertions.assertNotNull(statsDMetric, metric + " is mapped");
+ final ImmutableMap<String, String> dims = actual.build();
+ Assertions.assertEquals(
+ ImmutableMap.of("dataSource", "data-source", "taskType",
"index_kafka"),
+ dims,
+ "correct Dimensions for " + metric
+ );
+ // Dimensions are iterated in sorted order, and for non-dogstatsd output
their values are
+ // appended to the dotted metric name in that order, so the emitted
order is user-visible.
+ Assertions.assertEquals(
+ List.of("dataSource", "taskType"),
+ List.copyOf(dims.keySet()),
+ "correct Dimension order for " + metric
+ );
+ }
+ }
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]