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"); + } } }
