This is an automated email from the ASF dual-hosted git repository.

reiabreu pushed a commit to branch 2.x
in repository https://gitbox.apache.org/repos/asf/storm.git


The following commit(s) were added to refs/heads/2.x by this push:
     new 5126b4976 STORM-4041: fix(topology_lag): Kafka Topology Lag breaking 
when no offsets are committed (Backport #8589) (#8591)
5126b4976 is described below

commit 5126b4976e39494af9053f1a42230aa7d0acb314
Author: reiabreu <[email protected]>
AuthorDate: Sat May 9 18:39:04 2026 +0100

    STORM-4041: fix(topology_lag): Kafka Topology Lag breaking when no offsets 
are committed (Backport #8589) (#8591)
---
 external/storm-kafka-monitor/pom.xml               |  12 +++
 .../storm/kafka/monitor/KafkaOffsetLagUtil.java    |  20 ++--
 .../kafka/monitor/KafkaOffsetLagUtilTest.java      | 115 +++++++++++++++++++++
 .../org/apache/storm/utils/TopologySpoutLag.java   |  11 +-
 4 files changed, 142 insertions(+), 16 deletions(-)

diff --git a/external/storm-kafka-monitor/pom.xml 
b/external/storm-kafka-monitor/pom.xml
index b4dc67660..fdd526941 100644
--- a/external/storm-kafka-monitor/pom.xml
+++ b/external/storm-kafka-monitor/pom.xml
@@ -67,6 +67,18 @@
             <artifactId>jakarta.xml.bind-api</artifactId>
             <!-- version will be inherited -->
         </dependency>
+        <dependency>
+            <groupId>org.testcontainers</groupId>
+            <artifactId>kafka</artifactId>
+            <version>${testcontainers.version}</version>
+            <scope>test</scope>
+        </dependency>
+        <dependency>
+            <groupId>org.testcontainers</groupId>
+            <artifactId>junit-jupiter</artifactId>
+            <version>${testcontainers.version}</version>
+            <scope>test</scope>
+        </dependency>
     </dependencies>
 
     <build>
diff --git 
a/external/storm-kafka-monitor/src/main/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtil.java
 
b/external/storm-kafka-monitor/src/main/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtil.java
index fa06ffa3e..d5918f784 100644
--- 
a/external/storm-kafka-monitor/src/main/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtil.java
+++ 
b/external/storm-kafka-monitor/src/main/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtil.java
@@ -19,9 +19,8 @@
 package org.apache.storm.kafka.monitor;
 
 import java.util.ArrayList;
-import java.util.Collection;
-import java.util.Collections;
 import java.util.HashMap;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
 import java.util.Properties;
@@ -172,12 +171,13 @@ public class KafkaOffsetLagUtil {
                 }
             }
             consumer.assign(topicPartitionList);
+            Map<TopicPartition, OffsetAndMetadata> committedOffsets = 
consumer.committed(new HashSet<>(topicPartitionList));
+            consumer.seekToEnd(topicPartitionList);
             for (TopicPartition topicPartition : topicPartitionList) {
-                Map<TopicPartition, OffsetAndMetadata> offsetAndMetadata = 
consumer.committed(Collections.singleton(topicPartition));
-                long committedOffset = offsetAndMetadata != null ? 
offsetAndMetadata.get(topicPartition).offset() : -1;
-                consumer.seekToEnd(toArrayList(topicPartition));
+                OffsetAndMetadata partitionOffset = 
committedOffsets.get(topicPartition);
+                long committedOffset = partitionOffset != null ? 
partitionOffset.offset() : -1;
                 result.add(new KafkaOffsetLagResult(topicPartition.topic(), 
topicPartition.partition(), committedOffset,
-                                                    
consumer.position(topicPartition)));
+                        consumer.position(topicPartition)));
             }
         } finally {
             if (consumer != null) {
@@ -187,12 +187,4 @@ public class KafkaOffsetLagUtil {
         return result;
     }
 
-    private static Collection<TopicPartition> toArrayList(final TopicPartition 
tp) {
-        return new ArrayList<TopicPartition>(1) {
-            {
-                add(tp);
-            }
-        };
-    }
-
 }
diff --git 
a/external/storm-kafka-monitor/src/test/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtilTest.java
 
b/external/storm-kafka-monitor/src/test/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtilTest.java
new file mode 100644
index 000000000..c70a2db4f
--- /dev/null
+++ 
b/external/storm-kafka-monitor/src/test/java/org/apache/storm/kafka/monitor/KafkaOffsetLagUtilTest.java
@@ -0,0 +1,115 @@
+/*
+ * 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.storm.kafka.monitor;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotNull;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+import java.util.Properties;
+import org.apache.kafka.clients.admin.Admin;
+import org.apache.kafka.clients.admin.NewTopic;
+import org.apache.kafka.clients.consumer.ConsumerConfig;
+import org.apache.kafka.clients.consumer.KafkaConsumer;
+import org.apache.kafka.clients.consumer.OffsetAndMetadata;
+import org.apache.kafka.clients.producer.KafkaProducer;
+import org.apache.kafka.clients.producer.ProducerConfig;
+import org.apache.kafka.clients.producer.ProducerRecord;
+import org.apache.kafka.common.TopicPartition;
+import org.apache.kafka.common.serialization.StringDeserializer;
+import org.apache.kafka.common.serialization.StringSerializer;
+import org.junit.jupiter.api.BeforeAll;
+import org.junit.jupiter.api.Test;
+import org.testcontainers.junit.jupiter.Container;
+import org.testcontainers.junit.jupiter.Testcontainers;
+import org.testcontainers.kafka.KafkaContainer;
+import org.testcontainers.utility.DockerImageName;
+
+/**
+ * Integration test — requires a Docker daemon. Skipped automatically when 
Docker is unavailable.
+ */
+@Testcontainers(disabledWithoutDocker = true)
+class KafkaOffsetLagUtilTest {
+
+    private static final String TOPIC = "lag-test-topic";
+    private static final String GROUP_ID = "lag-test-group";
+    private static final int PARTITIONS = 2;
+    private static final long COMMITTED_OFFSET_PARTITION_0 = 7L;
+
+    @Container
+    private static final KafkaContainer KAFKA = new 
KafkaContainer(DockerImageName.parse("apache/kafka:4.0.0"));
+
+    @BeforeAll
+    static void seedKafka() throws Exception {
+        Properties adminProps = new Properties();
+        adminProps.put("bootstrap.servers", KAFKA.getBootstrapServers());
+        try (Admin admin = Admin.create(adminProps)) {
+            admin.createTopics(Collections.singletonList(new NewTopic(TOPIC, 
PARTITIONS, (short) 1))).all().get();
+        }
+        produceOneRecordToEachPartition();
+        commitOffsetForPartitionZeroOnly();
+    }
+
+    @Test
+    void 
getOffsetLagsReportsCommittedOffsetForCommittedPartitionsAndMinusOneForUncommittedPartitions()
 throws Exception {
+        NewKafkaSpoutOffsetQuery query = new NewKafkaSpoutOffsetQuery(
+                TOPIC, KAFKA.getBootstrapServers(), GROUP_ID, null, null, 
null);
+
+        List<KafkaOffsetLagResult> results = 
KafkaOffsetLagUtil.getOffsetLags(query);
+
+        assertEquals(PARTITIONS, results.size(), "Expected one result per 
partition");
+        Map<Integer, Long> committedByPartition = new HashMap<>();
+        for (KafkaOffsetLagResult r : results) {
+            committedByPartition.put(r.getPartition(), 
r.getConsumerCommittedOffset());
+        }
+        assertNotNull(committedByPartition.get(0));
+        assertEquals(COMMITTED_OFFSET_PARTITION_0, committedByPartition.get(0),
+                "Partition with a committed offset should report it");
+        // Regression assertion: pre-fix this NPE'd inside the monitor and 
surfaced as ClassCastException upstream.
+        assertEquals(-1L, committedByPartition.get(1),
+                "Partition with no committed offset should report -1, not 
throw");
+    }
+
+    private static void produceOneRecordToEachPartition() {
+        Properties producerProps = new Properties();
+        producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, 
KAFKA.getBootstrapServers());
+        producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, 
StringSerializer.class.getName());
+        producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, 
StringSerializer.class.getName());
+        try (KafkaProducer<String, String> producer = new 
KafkaProducer<>(producerProps)) {
+            for (int p = 0; p < PARTITIONS; p++) {
+                producer.send(new ProducerRecord<>(TOPIC, p, "k", "v"));
+            }
+            producer.flush();
+        }
+    }
+
+    private static void commitOffsetForPartitionZeroOnly() {
+        Properties consumerProps = new Properties();
+        consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, 
KAFKA.getBootstrapServers());
+        consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, GROUP_ID);
+        consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false");
+        consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, 
StringDeserializer.class.getName());
+        consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, 
StringDeserializer.class.getName());
+        try (KafkaConsumer<String, String> consumer = new 
KafkaConsumer<>(consumerProps)) {
+            TopicPartition p0 = new TopicPartition(TOPIC, 0);
+            consumer.commitSync(Collections.singletonMap(p0, new 
OffsetAndMetadata(COMMITTED_OFFSET_PARTITION_0)));
+        }
+    }
+}
diff --git a/storm-core/src/jvm/org/apache/storm/utils/TopologySpoutLag.java 
b/storm-core/src/jvm/org/apache/storm/utils/TopologySpoutLag.java
index 4859552eb..2bf1f2bd8 100644
--- a/storm-core/src/jvm/org/apache/storm/utils/TopologySpoutLag.java
+++ b/storm-core/src/jvm/org/apache/storm/utils/TopologySpoutLag.java
@@ -171,10 +171,17 @@ public class TopologySpoutLag {
                     String resultFromMonitor = new 
ShellCommandRunnerImpl().execCommand(commands.toArray(new String[0]));
 
                     try {
-                        result = (Map<String, Object>) 
JSONValue.parseWithException(resultFromMonitor);
+                        Object parsed = 
JSONValue.parseWithException(resultFromMonitor);
+                        if (parsed instanceof Map) {
+                            result = (Map<String, Object>) parsed;
+                        } else {
+                            // json-smart parses unquoted plain text leniently 
as a String, so we can land here
+                            // when the monitor printed an error message 
instead of JSON; surface it as the error.
+                            LOGGER.debug("Monitor returned non-JSON output, 
treating as error: {}", resultFromMonitor);
+                            errorMsg = resultFromMonitor;
+                        }
                     } catch (ParseException e) {
                         LOGGER.debug("JSON parsing failed, assuming message as 
error message: {}", resultFromMonitor);
-                        // json parsing fail -> error received
                         errorMsg = resultFromMonitor;
                     }
                 } finally {

Reply via email to