diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java index d4e6358e2ea99..dcfc28c5496b9 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java @@ -44,6 +44,7 @@ import org.apache.kafka.connect.util.ConnectUtils; import org.apache.kafka.connect.util.ConnectorTaskId; import org.apache.kafka.connect.util.KafkaBasedLog; +import org.apache.kafka.connect.util.SharedTopicAdmin; import org.apache.kafka.connect.util.TopicAdmin; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -226,6 +227,7 @@ public static String COMMIT_TASKS_KEY(String connectorName) { private final Map> connectorConfigs = new HashMap<>(); private final Map> taskConfigs = new HashMap<>(); private final Supplier topicAdminSupplier; + private SharedTopicAdmin ownTopicAdmin; // Set of connectors where we saw a task commit with an incomplete set of task config updates, indicating the data // is in an inconsistent state and we cannot safely use them until they have been refreshed. @@ -291,7 +293,13 @@ public void start() { @Override public void stop() { log.info("Closing KafkaConfigBackingStore"); - configLog.stop(); + try { + configLog.stop(); + } finally { + if (ownTopicAdmin != null) { + ownTopicAdmin.close(); + } + } log.info("Closed KafkaConfigBackingStore"); } @@ -479,7 +487,14 @@ KafkaBasedLog setupAndCreateKafkaBasedLog(String topic, final Wo Map adminProps = new HashMap<>(originals); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); - Supplier adminSupplier = topicAdminSupplier != null ? topicAdminSupplier : () -> new TopicAdmin(adminProps); + Supplier adminSupplier; + if (topicAdminSupplier != null) { + adminSupplier = topicAdminSupplier; + } else { + // Create our own topic admin supplier that we'll close when we're stopped + ownTopicAdmin = new SharedTopicAdmin(adminProps); + adminSupplier = ownTopicAdmin; + } Map topicSettings = config instanceof DistributedConfig ? ((DistributedConfig) config).configStorageTopicSettings() : Collections.emptyMap(); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java index 26b47f996b18a..313baf72c58c0 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java @@ -32,6 +32,7 @@ import org.apache.kafka.connect.util.ConnectUtils; import org.apache.kafka.connect.util.ConvertingFutureCallback; import org.apache.kafka.connect.util.KafkaBasedLog; +import org.apache.kafka.connect.util.SharedTopicAdmin; import org.apache.kafka.connect.util.TopicAdmin; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -65,6 +66,7 @@ public class KafkaOffsetBackingStore implements OffsetBackingStore { private KafkaBasedLog offsetLog; private HashMap data; private final Supplier topicAdminSupplier; + private SharedTopicAdmin ownTopicAdmin; @Deprecated public KafkaOffsetBackingStore() { @@ -98,7 +100,14 @@ public void configure(final WorkerConfig config) { Map adminProps = new HashMap<>(originals); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); - Supplier adminSupplier = topicAdminSupplier != null ? topicAdminSupplier : () -> new TopicAdmin(adminProps); + Supplier adminSupplier; + if (topicAdminSupplier != null) { + adminSupplier = topicAdminSupplier; + } else { + // Create our own topic admin supplier that we'll close when we're stopped + ownTopicAdmin = new SharedTopicAdmin(adminProps); + adminSupplier = ownTopicAdmin; + } Map topicSettings = config instanceof DistributedConfig ? ((DistributedConfig) config).offsetStorageTopicSettings() : Collections.emptyMap(); @@ -140,7 +149,13 @@ public void start() { @Override public void stop() { log.info("Stopping KafkaOffsetBackingStore"); - offsetLog.stop(); + try { + offsetLog.stop(); + } finally { + if (ownTopicAdmin != null) { + ownTopicAdmin.close(); + } + } log.info("Stopped KafkaOffsetBackingStore"); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index efa405f3a4b9b..eadbe18786717 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -44,6 +44,7 @@ import org.apache.kafka.connect.util.ConnectUtils; import org.apache.kafka.connect.util.ConnectorTaskId; import org.apache.kafka.connect.util.KafkaBasedLog; +import org.apache.kafka.connect.util.SharedTopicAdmin; import org.apache.kafka.connect.util.Table; import org.apache.kafka.connect.util.TopicAdmin; import org.slf4j.Logger; @@ -134,6 +135,7 @@ public class KafkaStatusBackingStore implements StatusBackingStore { private String statusTopic; private KafkaBasedLog kafkaLog; private int generation; + private SharedTopicAdmin ownTopicAdmin; @Deprecated public KafkaStatusBackingStore(Time time, Converter converter) { @@ -177,7 +179,14 @@ public void configure(final WorkerConfig config) { Map adminProps = new HashMap<>(originals); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); - Supplier adminSupplier = topicAdminSupplier != null ? topicAdminSupplier : () -> new TopicAdmin(adminProps); + Supplier adminSupplier; + if (topicAdminSupplier != null) { + adminSupplier = topicAdminSupplier; + } else { + // Create our own topic admin supplier that we'll close when we're stopped + ownTopicAdmin = new SharedTopicAdmin(adminProps); + adminSupplier = ownTopicAdmin; + } Map topicSettings = config instanceof DistributedConfig ? ((DistributedConfig) config).statusStorageTopicSettings() @@ -221,7 +230,13 @@ public void start() { @Override public void stop() { - kafkaLog.stop(); + try { + kafkaLog.stop(); + } finally { + if (ownTopicAdmin != null) { + ownTopicAdmin.close(); + } + } } @Override