From 5f31ce218f5a95fa413f59b79e7bb80eb37442a5 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 15 Oct 2019 16:34:43 -0700 Subject: [PATCH 1/4] KAFKA-9046: Use top-level worker configs for connector admin clients --- .../main/java/org/apache/kafka/connect/runtime/Worker.java | 4 +++- .../java/org/apache/kafka/connect/runtime/WorkerTest.java | 1 + 2 files changed, 4 insertions(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index 27955477fc573..54deaf14fa1be 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -606,7 +606,9 @@ static Map adminConfigs(ConnectorTaskId id, ConnectorConfig connConfig, Class connectorClass, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) { - Map adminProps = new HashMap<>(); + // Use the top-level worker configs to retain backwards compatibility with older releases which + // did not require a prefix for connector admin client configs in the worker configuration file + Map adminProps = new HashMap<>(config.originals()); adminProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); // User-specified overrides adminProps.putAll(config.originalsWithPrefix("admin.")); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index 36a9b66d7cf04..f5ff0fcfed07f 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -1144,6 +1144,7 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { connConfig.put("metadata.max.age.ms", "10000"); Map expectedConfigs = new HashMap<>(); + expectedConfigs.putAll(props); expectedConfigs.put("bootstrap.servers", "localhost:9092"); expectedConfigs.put("client.id", "testid"); expectedConfigs.put("metadata.max.age.ms", "10000"); From 575f9e0fbfbaaa34c32af6f54271f832bf836e92 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Wed, 23 Oct 2019 13:14:49 -0700 Subject: [PATCH 2/4] KAFKA-9046: Address review comments --- .../org/apache/kafka/connect/runtime/Worker.java | 15 ++++++++++++--- .../apache/kafka/connect/runtime/WorkerTest.java | 5 +++-- 2 files changed, 15 insertions(+), 5 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index 54deaf14fa1be..d80b8b4d63486 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.connect.runtime; +import org.apache.kafka.clients.admin.AdminClientConfig; import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.KafkaProducer; @@ -606,11 +607,19 @@ static Map adminConfigs(ConnectorTaskId id, ConnectorConfig connConfig, Class connectorClass, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) { + Map adminProps = new HashMap<>(); // Use the top-level worker configs to retain backwards compatibility with older releases which // did not require a prefix for connector admin client configs in the worker configuration file - Map adminProps = new HashMap<>(config.originals()); - adminProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); - // User-specified overrides + // Only include configs that do not begin with "admin." since those will be added next, but with + // the prefix stripped + Map nonPrefixedWorkerConfigs = config.originals().entrySet().stream() + .filter(e -> !e.getKey().startsWith("admin.")) + .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); + adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, + Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); + adminProps.putAll(nonPrefixedWorkerConfigs); + + // Admin client-specific overrides in the worker config adminProps.putAll(config.originalsWithPrefix("admin.")); // Connector-specified overrides diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index f5ff0fcfed07f..4c041f851380e 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -1145,6 +1145,9 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { Map expectedConfigs = new HashMap<>(); expectedConfigs.putAll(props); + expectedConfigs.remove("admin.client.id"); + expectedConfigs.remove("admin.metadata.max.age.ms"); + expectedConfigs.put("bootstrap.servers", "localhost:9092"); expectedConfigs.put("client.id", "testid"); expectedConfigs.put("metadata.max.age.ms", "10000"); @@ -1154,7 +1157,6 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { PowerMock.replayAll(); assertEquals(expectedConfigs, Worker.adminConfigs(new ConnectorTaskId("test", 1), configWithOverrides, connectorConfig, null, allConnectorClientConfigOverridePolicy)); - } @Test(expected = ConnectException.class) @@ -1172,7 +1174,6 @@ public void testAdminConfigsClientOverridesWithNonePolicy() { PowerMock.replayAll(); Worker.adminConfigs(new ConnectorTaskId("test", 1), configWithOverrides, connectorConfig, null, noneConnectorClientConfigOverridePolicy); - } private void assertStatusMetrics(long expected, String metricName) { From 6ef7a4b0da2f4402b0e2eef2a21f7760d8efbe32 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Wed, 23 Oct 2019 15:56:38 -0700 Subject: [PATCH 3/4] KAFKA-9046: Address review comments --- .../java/org/apache/kafka/connect/runtime/WorkerTest.java | 5 +---- 1 file changed, 1 insertion(+), 4 deletions(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index 4c041f851380e..873f49b454c72 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -1143,10 +1143,7 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { Map connConfig = new HashMap(); connConfig.put("metadata.max.age.ms", "10000"); - Map expectedConfigs = new HashMap<>(); - expectedConfigs.putAll(props); - expectedConfigs.remove("admin.client.id"); - expectedConfigs.remove("admin.metadata.max.age.ms"); + Map expectedConfigs = new HashMap<>(workerProps); expectedConfigs.put("bootstrap.servers", "localhost:9092"); expectedConfigs.put("client.id", "testid"); From ae232d262215c0c80aea45bbe5ffed4f9d05fb88 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Wed, 13 Nov 2019 11:52:39 -0800 Subject: [PATCH 4/4] KAFKA-9046: Address review comments --- .../java/org/apache/kafka/connect/runtime/Worker.java | 9 ++++++--- .../org/apache/kafka/connect/runtime/WorkerTest.java | 2 ++ 2 files changed, 8 insertions(+), 3 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index d80b8b4d63486..6bdc08b6d8be7 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -610,10 +610,13 @@ static Map adminConfigs(ConnectorTaskId id, Map adminProps = new HashMap<>(); // Use the top-level worker configs to retain backwards compatibility with older releases which // did not require a prefix for connector admin client configs in the worker configuration file - // Only include configs that do not begin with "admin." since those will be added next, but with - // the prefix stripped + // Ignore configs that begin with "admin." since those will be added next (with the prefix stripped) + // and those that begin with "producer." and "consumer.", since we know they aren't intended for + // the admin client Map nonPrefixedWorkerConfigs = config.originals().entrySet().stream() - .filter(e -> !e.getKey().startsWith("admin.")) + .filter(e -> !e.getKey().startsWith("admin.") + && !e.getKey().startsWith("producer.") + && !e.getKey().startsWith("consumer.")) .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index 873f49b454c72..7021503adbed6 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -1138,6 +1138,8 @@ public void testAdminConfigsClientOverridesWithAllPolicy() { Map props = new HashMap<>(workerProps); props.put("admin.client.id", "testid"); props.put("admin.metadata.max.age.ms", "5000"); + props.put("producer.bootstrap.servers", "cbeauho.com"); + props.put("consumer.bootstrap.servers", "localhost:4761"); WorkerConfig configWithOverrides = new StandaloneConfig(props); Map connConfig = new HashMap();