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 {