From 7318e4840f18269b662947c8d7f303ae10fcefa4 Mon Sep 17 00:00:00 2001 From: Thunder Stumpges Date: Wed, 7 Aug 2019 13:17:32 -0700 Subject: [PATCH] Fix for topic lag metrics where topic contains period (.) (messages-behind-high-watermark and high-watermark) are calculated from kafka's "records-lag" consumer metric. In version 1.1 or so, when metrics moved from including the topic name in the metric name to using tags, they added a replacement of period to underscore. See commit : https://github.com/apache/kafka/commit/5d81639907869ce7355c40d2bac176a655e52074#diff-b45245913eaae46aa847d2615d62cde0R1331 --- .../org/apache/samza/system/kafka/KafkaConsumerProxy.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java b/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java index f85f5bf08f..88f8510a7f 100644 --- a/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java +++ b/samza-kafka/src/main/scala/org/apache/samza/system/kafka/KafkaConsumerProxy.java @@ -397,7 +397,9 @@ private void populateCurrentLags(Set ssps) { // These are required by the KafkaConsumer to get the metrics HashMap tags = new HashMap<>(); tags.put("client-id", clientId); - tags.put("topic", tp.topic()); + // kafka replaces '.' with underscore '_' in many/all of their metrics tags for topic names. + // see https://github.com/apache/kafka/commit/5d81639907869ce7355c40d2bac176a655e52074#diff-b45245913eaae46aa847d2615d62cde0R1331 + tags.put("topic", tp.topic().replace('.', '_')); tags.put("partition", Integer.toString(tp.partition())); perPartitionMetrics.put(ssp, new MetricName("records-lag", "consumer-fetch-manager-metrics", "", tags));