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

SvenO3 pushed a commit to branch fix-adapter-producer-lifecycle
in repository https://gitbox.apache.org/repos/asf/streampipes.git

commit 85ef99bb5fcca194e5ca04e47e37eaeff5f000c3
Author: Sven Oehler <[email protected]>
AuthorDate: Mon Jul 13 14:56:50 2026 +0200

    Close kafka admin client
---
 .../messaging/kafka/SpKafkaProducer.java           | 28 +++++++++++-----------
 1 file changed, 14 insertions(+), 14 deletions(-)

diff --git 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
index 68637be02c..1886a9a8c5 100644
--- 
a/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
+++ 
b/streampipes-messaging-kafka/src/main/java/org/apache/streampipes/messaging/kafka/SpKafkaProducer.java
@@ -119,24 +119,24 @@ public class SpKafkaProducer implements EventProducer, 
Serializable {
     Properties props = new Properties();
     props.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerUrl);
 
-    AdminClient adminClient = KafkaAdminClient.create(props);
+    try (AdminClient adminClient = KafkaAdminClient.create(props)) {
+      ListTopicsResult topics = adminClient.listTopics();
 
-    ListTopicsResult topics = adminClient.listTopics();
+      if (!topicExists(topics)) {
+        Map<String, String> topicConfig = new HashMap<>();
+        String retentionTime = 
Environments.getEnvironment().getKafkaRetentionTimeMs().getValueOrDefault();
+        topicConfig.put(TopicConfig.RETENTION_MS_CONFIG, retentionTime);
 
-    if (!topicExists(topics)) {
-      Map<String, String> topicConfig = new HashMap<>();
-      String retentionTime = 
Environments.getEnvironment().getKafkaRetentionTimeMs().getValueOrDefault();
-      topicConfig.put(TopicConfig.RETENTION_MS_CONFIG, retentionTime);
+        final NewTopic newTopic = new NewTopic(topic, 1, (short) 1);
+        newTopic.configs(topicConfig);
 
-      final NewTopic newTopic = new NewTopic(topic, 1, (short) 1);
-      newTopic.configs(topicConfig);
+        final CreateTopicsResult createTopicsResult = 
adminClient.createTopics(Collections.singleton(newTopic));
+        createTopicsResult.values().get(topic).get();
+        LOG.info("Successfully created Kafka topic " + topic);
 
-      final CreateTopicsResult createTopicsResult = 
adminClient.createTopics(Collections.singleton(newTopic));
-      createTopicsResult.values().get(topic).get();
-      LOG.info("Successfully created Kafka topic " + topic);
-
-    } else {
-      LOG.debug("Topic " + topic + "already exists in the broker, skipping 
topic creation");
+      } else {
+        LOG.debug("Topic " + topic + "already exists in the broker, skipping 
topic creation");
+      }
     }
   }
 

Reply via email to