From f56595c7852a8ac56c85902549aaf25903ca710a Mon Sep 17 00:00:00 2001 From: Yash Mayya Date: Sun, 5 May 2024 19:02:58 +0530 Subject: [PATCH 1/2] MINOR: Remove deprecated constructors from Connect's Kafka*BackingStore classes --- .../kafka/connect/storage/KafkaConfigBackingStore.java | 5 ----- .../kafka/connect/storage/KafkaStatusBackingStore.java | 7 +------ 2 files changed, 1 insertion(+), 11 deletions(-) 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 e26f7d88f198d..da7357bf04f2f 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 @@ -334,11 +334,6 @@ void setConfigLog(KafkaBasedLog configLog) { this.configLog = configLog; } - @Deprecated - public KafkaConfigBackingStore(Converter converter, DistributedConfig config, WorkerConfigTransformer configTransformer) { - this(converter, config, configTransformer, null, "connect-distributed-"); - } - public KafkaConfigBackingStore(Converter converter, DistributedConfig config, WorkerConfigTransformer configTransformer, Supplier adminSupplier, String clientIdBase) { this(converter, config, configTransformer, adminSupplier, clientIdBase, Time.SYSTEM); } 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 0ffc3eed4f346..bd129efcbaaa0 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 @@ -142,11 +142,6 @@ public class KafkaStatusBackingStore extends KafkaTopicBasedBackingStore impleme private SharedTopicAdmin ownTopicAdmin; private ExecutorService sendRetryExecutor; - @Deprecated - public KafkaStatusBackingStore(Time time, Converter converter) { - this(time, converter, null, "connect-distributed-"); - } - public KafkaStatusBackingStore(Time time, Converter converter, Supplier topicAdminSupplier, String clientIdBase) { this.time = time; this.converter = converter; @@ -159,7 +154,7 @@ public KafkaStatusBackingStore(Time time, Converter converter, Supplier kafkaLog) { - this(time, converter); + this(time, converter, null, "connect-distributed-"); this.kafkaLog = kafkaLog; this.statusTopic = statusTopic; sendRetryExecutor = Executors.newSingleThreadExecutor( From 73d82d0c507aec9be51ff253319f1b8f41e6a3f8 Mon Sep 17 00:00:00 2001 From: Yash Mayya Date: Thu, 16 May 2024 16:29:46 +0530 Subject: [PATCH 2/2] Remove testing-only constructor in KafkaStatusBackingStore; remove 'ownTopicAdmin' cruft from Kafka*BackingStore classes --- .../storage/KafkaConfigBackingStore.java | 14 ++----------- .../storage/KafkaOffsetBackingStore.java | 20 ++----------------- .../storage/KafkaStatusBackingStore.java | 18 +++-------------- .../KafkaConfigBackingStoreMockitoTest.java | 2 +- .../storage/KafkaConfigBackingStoreTest.java | 3 ++- .../KafkaStatusBackingStoreFormatTest.java | 7 +++---- .../storage/KafkaStatusBackingStoreTest.java | 2 +- 7 files changed, 14 insertions(+), 52 deletions(-) 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 da7357bf04f2f..c24dd6c7907a0 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 @@ -53,7 +53,6 @@ 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; @@ -300,7 +299,6 @@ public static String LOGGER_CLUSTER_KEY(String namespace) { final Map> connectorConfigs = new HashMap<>(); final Map> taskConfigs = new HashMap<>(); private final Supplier topicAdminSupplier; - private SharedTopicAdmin ownTopicAdmin; private final String clientId; // Set of connectors where we saw a task commit with an incomplete set of task config updates, indicating the data @@ -406,7 +404,6 @@ public void stop() { log.info("Closing KafkaConfigBackingStore"); relinquishWritePrivileges(); - Utils.closeQuietly(ownTopicAdmin, "admin for config topic"); Utils.closeQuietly(configLog::stop, "KafkaBasedLog for config topic"); log.info("Closed KafkaConfigBackingStore"); @@ -789,14 +786,7 @@ KafkaBasedLog setupAndCreateKafkaBasedLog(String topic, final Wo Map adminProps = new HashMap<>(originals); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); adminProps.put(CommonClientConfigs.CLIENT_ID_CONFIG, clientId); - 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(); @@ -807,7 +797,7 @@ KafkaBasedLog setupAndCreateKafkaBasedLog(String topic, final Wo .replicationFactor(config.getShort(DistributedConfig.CONFIG_STORAGE_REPLICATION_FACTOR_CONFIG)) .build(); - return createKafkaBasedLog(topic, producerProps, consumerProps, new ConsumeCallback(), topicDescription, adminSupplier, config, time); + return createKafkaBasedLog(topic, producerProps, consumerProps, new ConsumeCallback(), topicDescription, topicAdminSupplier, config, time); } /** 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 4d447759eddfa..e3960748b9c4e 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 @@ -37,7 +37,6 @@ 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; @@ -149,7 +148,6 @@ private static String noClientId() { private Converter keyConverter; private final Supplier topicAdminSupplier; private final Supplier clientIdBase; - private SharedTopicAdmin ownTopicAdmin; protected boolean exactlyOnce; /** @@ -211,17 +209,9 @@ public void configure(final WorkerConfig config) { Map adminProps = new HashMap<>(originals); adminProps.put(CommonClientConfigs.CLIENT_ID_CONFIG, clientId); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); - Supplier adminSupplier; - if (topicAdminSupplier != null) { - adminSupplier = topicAdminSupplier; - } else { - // Create our own topic admin supplier that we'll close when we're stopped - this.ownTopicAdmin = new SharedTopicAdmin(adminProps); - adminSupplier = ownTopicAdmin; - } NewTopic topicDescription = newTopicDescription(topic, config); - this.offsetLog = createKafkaBasedLog(topic, producerProps, consumerProps, consumedCallback, topicDescription, adminSupplier, config, Time.SYSTEM); + this.offsetLog = createKafkaBasedLog(topic, producerProps, consumerProps, consumedCallback, topicDescription, topicAdminSupplier, config, Time.SYSTEM); } protected NewTopic newTopicDescription(final String topic, final WorkerConfig config) { @@ -268,13 +258,7 @@ public void start() { @Override public void stop() { log.info("Stopping KafkaOffsetBackingStore"); - try { - offsetLog.stop(); - } finally { - if (ownTopicAdmin != null) { - ownTopicAdmin.close(); - } - } + offsetLog.stop(); 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 bd129efcbaaa0..4f2d832fc8379 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,7 +44,6 @@ 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; @@ -139,7 +138,6 @@ public class KafkaStatusBackingStore extends KafkaTopicBasedBackingStore impleme private String statusTopic; private KafkaBasedLog kafkaLog; private int generation; - private SharedTopicAdmin ownTopicAdmin; private ExecutorService sendRetryExecutor; public KafkaStatusBackingStore(Time time, Converter converter, Supplier topicAdminSupplier, String clientIdBase) { @@ -153,7 +151,8 @@ public KafkaStatusBackingStore(Time time, Converter converter, Supplier kafkaLog) { + KafkaStatusBackingStore(Time time, Converter converter, String statusTopic, Supplier topicAdminSupplier, + KafkaBasedLog kafkaLog) { this(time, converter, null, "connect-distributed-"); this.kafkaLog = kafkaLog; this.statusTopic = statusTopic; @@ -194,14 +193,6 @@ public void configure(final WorkerConfig config) { Map adminProps = new HashMap<>(originals); adminProps.put(CommonClientConfigs.CLIENT_ID_CONFIG, clientId); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); - 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() @@ -214,7 +205,7 @@ public void configure(final WorkerConfig config) { .build(); Callback> readCallback = (error, record) -> read(record); - this.kafkaLog = createKafkaBasedLog(statusTopic, producerProps, consumerProps, readCallback, topicDescription, adminSupplier, config, time); + this.kafkaLog = createKafkaBasedLog(statusTopic, producerProps, consumerProps, readCallback, topicDescription, topicAdminSupplier, config, time); } @Override @@ -231,9 +222,6 @@ public void stop() { kafkaLog.stop(); } finally { ThreadUtils.shutdownExecutorServiceQuietly(sendRetryExecutor, 10, TimeUnit.SECONDS); - if (ownTopicAdmin != null) { - ownTopicAdmin.close(); - } } } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreMockitoTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreMockitoTest.java index c8bdea1fd8c54..76fcb468ddc60 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreMockitoTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreMockitoTest.java @@ -218,7 +218,7 @@ private void createStore() { doReturn("test-cluster").when(config).kafkaClusterId(); configStorage = Mockito.spy( new KafkaConfigBackingStore( - converter, config, null, null, CLIENT_ID_BASE, time) + converter, config, null, () -> null, CLIENT_ID_BASE, time) ); configStorage.setConfigLog(configLog); configStorage.setUpdateListener(configUpdateListener); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreTest.java index 637dac4d679df..e57c31427a44d 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaConfigBackingStoreTest.java @@ -179,10 +179,11 @@ private void createStore() { // The kafkaClusterId is used in the constructor for KafkaConfigBackingStore // So temporarily enter replay mode in order to mock that call EasyMock.replay(config); + Supplier topicAdminSupplier = () -> null; configStorage = PowerMock.createPartialMock( KafkaConfigBackingStore.class, new String[]{"createKafkaBasedLog", "createFencableProducer"}, - converter, config, null, null, CLIENT_ID_BASE, time); + converter, config, null, topicAdminSupplier, CLIENT_ID_BASE, time); Whitebox.setInternalState(configStorage, "configLog", storeLog); configStorage.setUpdateListener(configUpdateListener); // The mock must be reset and re-mocked for the remainder of the test. diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreFormatTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreFormatTest.java index a13b76e400fe3..ee53a4942041c 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreFormatTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreFormatTest.java @@ -64,16 +64,15 @@ public class KafkaStatusBackingStoreFormatTest { private Time time; private KafkaStatusBackingStore store; - private JsonConverter converter; - private KafkaBasedLog kafkaBasedLog = mock(KafkaBasedLog.class); + private final KafkaBasedLog kafkaBasedLog = mock(KafkaBasedLog.class); @Before public void setup() { time = new MockTime(); - converter = new JsonConverter(); + JsonConverter converter = new JsonConverter(); converter.configure(Collections.singletonMap(SCHEMAS_ENABLE_CONFIG, false), false); - store = new KafkaStatusBackingStore(new MockTime(), converter, STATUS_TOPIC, kafkaBasedLog); + store = new KafkaStatusBackingStore(new MockTime(), converter, STATUS_TOPIC, () -> null, kafkaBasedLog); } @Test diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java index 11a4b378c962d..76002b763711a 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/storage/KafkaStatusBackingStoreTest.java @@ -80,7 +80,7 @@ public class KafkaStatusBackingStoreTest { @Before public void setup() { - store = new KafkaStatusBackingStore(new MockTime(), converter, STATUS_TOPIC, kafkaBasedLog); + store = new KafkaStatusBackingStore(new MockTime(), converter, STATUS_TOPIC, () -> null, kafkaBasedLog); } @Test