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

rzo1 pushed a commit to branch fix/spout-lag-consumer-property-filter
in repository https://gitbox.apache.org/repos/asf/storm.git

commit fe2b3f4be77fd9df50d30aa350525250a81c5893
Author: Richard Zowalla <[email protected]>
AuthorDate: Fri Aug 21 14:21:44 2026 +0200

    Pass only the required Kafka consumer properties to the spout lag monitor
---
 .../org/apache/storm/utils/TopologySpoutLag.java   | 39 +++++++++-
 .../apache/storm/utils/TopologySpoutLagTest.java   | 89 ++++++++++++++++++++++
 2 files changed, 126 insertions(+), 2 deletions(-)

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 d6bc36052..d9bbc92b7 100644
--- a/storm-core/src/jvm/org/apache/storm/utils/TopologySpoutLag.java
+++ b/storm-core/src/jvm/org/apache/storm/utils/TopologySpoutLag.java
@@ -16,6 +16,7 @@
 
 package org.apache.storm.utils;
 
+import com.google.common.annotations.VisibleForTesting;
 import com.google.common.base.Strings;
 
 import java.io.File;
@@ -52,6 +53,28 @@ public class TopologySpoutLag {
     private static final String SECURITY_PROTOCOL_CONFIG = CONFIG_KEY_PREFIX + 
"security.protocol";
     private static final Set<String> ALL_CONFIGS = new 
HashSet<>(Arrays.asList(TOPICS_CONFIG, GROUPID_CONFIG,
             BOOTSTRAP_CONFIG, SECURITY_PROTOCOL_CONFIG));
+    // The spout json_conf comes from the submitted topology, while 
storm-kafka-monitor runs on the UI host.
+    // Only the connection settings the monitor actually needs to reach the 
brokers are forwarded to it; anything
+    // else (deserializers, interceptor.classes, metric.reporters, 
sasl.jaas.config, callback handler classes, ...)
+    // is dropped, so the monitor keeps using its own defaults rather than 
classes named by the topology.
+    private static final Set<String> ALLOWED_EXTRA_PROPERTIES = new 
HashSet<>(Arrays.asList(
+            "client.id",
+            "request.timeout.ms",
+            "session.timeout.ms",
+            "sasl.mechanism",
+            "sasl.kerberos.service.name",
+            "ssl.protocol",
+            "ssl.provider",
+            "ssl.enabled.protocols",
+            "ssl.cipher.suites",
+            "ssl.endpoint.identification.algorithm",
+            "ssl.truststore.type",
+            "ssl.truststore.location",
+            "ssl.truststore.password",
+            "ssl.keystore.type",
+            "ssl.keystore.location",
+            "ssl.keystore.password",
+            "ssl.key.password"));
     private static final Logger LOGGER = 
LoggerFactory.getLogger(TopologySpoutLag.class);
 
     // The storm-kafka-monitor jars ship only in the full binary distribution; 
users of the lite
@@ -110,14 +133,26 @@ public class TopologySpoutLag {
         return commands;
     }
 
-    private static File createExtraPropertiesFile(Map<String, Object> 
jsonConf) {
+    @VisibleForTesting
+    static File createExtraPropertiesFile(Map<String, Object> jsonConf) {
         File file = null;
         Map<String, String> extraProperties = new HashMap<>();
+        List<String> droppedProperties = new ArrayList<>();
         for (Map.Entry<String, Object> conf : jsonConf.entrySet()) {
             if (conf.getKey().startsWith(CONFIG_KEY_PREFIX) && 
!ALL_CONFIGS.contains(conf.getKey())) {
-                
extraProperties.put(conf.getKey().substring(CONFIG_KEY_PREFIX.length()), 
conf.getValue().toString());
+                String consumerKey = 
conf.getKey().substring(CONFIG_KEY_PREFIX.length());
+                if (ALLOWED_EXTRA_PROPERTIES.contains(consumerKey)) {
+                    extraProperties.put(consumerKey, 
conf.getValue().toString());
+                } else {
+                    droppedProperties.add(consumerKey);
+                }
             }
         }
+        // The UI polls the lag endpoint, so log the dropped keys once per 
call rather than one line each.
+        if (!droppedProperties.isEmpty()) {
+            LOGGER.info("Not passing consumer properties {} from the topology 
to the Kafka spout lag monitor, "
+                + "only these properties are passed on: {}", 
droppedProperties, ALLOWED_EXTRA_PROPERTIES);
+        }
         if (!extraProperties.isEmpty()) {
             try {
                 file = Files.createTempFile("kafka-consumer-extra", 
"props").toFile();
diff --git 
a/storm-core/test/jvm/org/apache/storm/utils/TopologySpoutLagTest.java 
b/storm-core/test/jvm/org/apache/storm/utils/TopologySpoutLagTest.java
new file mode 100644
index 000000000..864596877
--- /dev/null
+++ b/storm-core/test/jvm/org/apache/storm/utils/TopologySpoutLagTest.java
@@ -0,0 +1,89 @@
+/*
+ * 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.utils;
+
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNull;
+
+import java.io.File;
+import java.io.FileInputStream;
+import java.io.IOException;
+import java.io.InputStream;
+import java.util.HashMap;
+import java.util.Map;
+import java.util.Properties;
+
+import org.junit.jupiter.api.Test;
+
+public class TopologySpoutLagTest {
+
+    private Properties loadProperties(File file) throws IOException {
+        Properties properties = new Properties();
+        try (InputStream in = new FileInputStream(file)) {
+            properties.load(in);
+        }
+        return properties;
+    }
+
+    @Test
+    public void testOnlyKnownConsumerPropertiesArePassedToTheMonitor() throws 
Exception {
+        Map<String, Object> jsonConf = new HashMap<>();
+        // Passed to the monitor as command line options instead of via the 
properties file.
+        jsonConf.put("config.topics", "topic");
+        jsonConf.put("config.groupid", "group");
+        jsonConf.put("config.bootstrap.servers", "broker:9092");
+        jsonConf.put("config.security.protocol", "SASL_SSL");
+        // Connection settings the monitor needs.
+        jsonConf.put("config.sasl.mechanism", "SCRAM-SHA-512");
+        jsonConf.put("config.ssl.truststore.location", 
"/etc/storm/truststore.jks");
+        jsonConf.put("config.client.id", "lag-monitor");
+        // Properties that would make the monitor load classes named by the 
topology.
+        jsonConf.put("config.key.deserializer", "org.example.MyDeserializer");
+        jsonConf.put("config.value.deserializer", 
"org.example.MyDeserializer");
+        jsonConf.put("config.interceptor.classes", 
"org.example.MyInterceptor");
+        jsonConf.put("config.metric.reporters", "org.example.MyReporter");
+        jsonConf.put("config.sasl.jaas.config", "org.example.MyLoginModule 
required;");
+        jsonConf.put("config.sasl.login.callback.handler.class", 
"org.example.MyCallbackHandler");
+        // Keys without the config. prefix are not consumer properties at all.
+        jsonConf.put("topology.name", "test");
+
+        File file = TopologySpoutLag.createExtraPropertiesFile(jsonConf);
+        try {
+            Properties properties = loadProperties(file);
+            Properties expected = new Properties();
+            expected.put("sasl.mechanism", "SCRAM-SHA-512");
+            expected.put("ssl.truststore.location", 
"/etc/storm/truststore.jks");
+            expected.put("client.id", "lag-monitor");
+            assertEquals(expected, properties);
+        } finally {
+            file.delete();
+        }
+    }
+
+    @Test
+    public void testNoPropertiesFileWhenNoKnownConsumerPropertiesAreSet() {
+        Map<String, Object> jsonConf = new HashMap<>();
+        jsonConf.put("config.topics", "topic");
+        jsonConf.put("config.groupid", "group");
+        jsonConf.put("config.bootstrap.servers", "broker:9092");
+        jsonConf.put("config.interceptor.classes", 
"org.example.MyInterceptor");
+
+        assertNull(TopologySpoutLag.createExtraPropertiesFile(jsonConf));
+    }
+}

Reply via email to