From 96caaa9a4934bcef78d7b145d18aa1718cb10009 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Thu, 9 Apr 2020 08:32:37 -0700 Subject: [PATCH 1/7] KAFKA-9845: Fix plugin.path when config provider is used --- .../org/apache/kafka/connect/cli/ConnectDistributed.java | 5 +++-- .../java/org/apache/kafka/connect/cli/ConnectStandalone.java | 4 ++-- 2 files changed, 5 insertions(+), 4 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java index 22c1ad82d6138..a4a04ec3adcdd 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java @@ -87,10 +87,11 @@ public static void main(String[] args) { } public Connect startConnect(Map workerProps) { + DistributedConfig config = new DistributedConfig(workerProps); + log.info("Scanning for plugin classes. This might take a moment ..."); - Plugins plugins = new Plugins(workerProps); + Plugins plugins = new Plugins(config.originalsStrings()); plugins.compareAndSwapWithDelegatingLoader(); - DistributedConfig config = new DistributedConfig(workerProps); String kafkaClusterId = ConnectUtils.lookupKafkaClusterId(config); log.debug("Kafka cluster ID: {}", kafkaClusterId); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java index cf7b93bd7c838..65d40bdfb521c 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java @@ -74,11 +74,11 @@ public static void main(String[] args) { String workerPropsFile = args[0]; Map workerProps = !workerPropsFile.isEmpty() ? Utils.propsToStringMap(Utils.loadProps(workerPropsFile)) : Collections.emptyMap(); + StandaloneConfig config = new StandaloneConfig(workerProps); log.info("Scanning for plugin classes. This might take a moment ..."); - Plugins plugins = new Plugins(workerProps); + Plugins plugins = new Plugins(config.originalsStrings()); plugins.compareAndSwapWithDelegatingLoader(); - StandaloneConfig config = new StandaloneConfig(workerProps); String kafkaClusterId = ConnectUtils.lookupKafkaClusterId(config); log.debug("Kafka cluster ID: {}", kafkaClusterId); From ab16154ff778e36eb509a43639438e8692cb048e Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Thu, 9 Apr 2020 10:37:48 -0700 Subject: [PATCH 2/7] Revert "KAFKA-9845: Fix plugin.path when config provider is used" This reverts commit 96caaa9a4934bcef78d7b145d18aa1718cb10009. --- .../org/apache/kafka/connect/cli/ConnectDistributed.java | 5 ++--- .../java/org/apache/kafka/connect/cli/ConnectStandalone.java | 4 ++-- 2 files changed, 4 insertions(+), 5 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java index a4a04ec3adcdd..22c1ad82d6138 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java @@ -87,11 +87,10 @@ public static void main(String[] args) { } public Connect startConnect(Map workerProps) { - DistributedConfig config = new DistributedConfig(workerProps); - log.info("Scanning for plugin classes. This might take a moment ..."); - Plugins plugins = new Plugins(config.originalsStrings()); + Plugins plugins = new Plugins(workerProps); plugins.compareAndSwapWithDelegatingLoader(); + DistributedConfig config = new DistributedConfig(workerProps); String kafkaClusterId = ConnectUtils.lookupKafkaClusterId(config); log.debug("Kafka cluster ID: {}", kafkaClusterId); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java index 65d40bdfb521c..cf7b93bd7c838 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java @@ -74,11 +74,11 @@ public static void main(String[] args) { String workerPropsFile = args[0]; Map workerProps = !workerPropsFile.isEmpty() ? Utils.propsToStringMap(Utils.loadProps(workerPropsFile)) : Collections.emptyMap(); - StandaloneConfig config = new StandaloneConfig(workerProps); log.info("Scanning for plugin classes. This might take a moment ..."); - Plugins plugins = new Plugins(config.originalsStrings()); + Plugins plugins = new Plugins(workerProps); plugins.compareAndSwapWithDelegatingLoader(); + StandaloneConfig config = new StandaloneConfig(workerProps); String kafkaClusterId = ConnectUtils.lookupKafkaClusterId(config); log.debug("Kafka cluster ID: {}", kafkaClusterId); From f9387144bffd9257df08c49aa8a177ee7fc5c451 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Thu, 9 Apr 2020 14:46:19 -0700 Subject: [PATCH 3/7] KAFKA-9845: Emit ERROR-level log message when config provider is used for plugin.path property --- .../kafka/connect/runtime/WorkerConfig.java | 16 ++++++++++++++++ 1 file changed, 16 insertions(+) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index 347e250cefbf8..7b1d248164060 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -37,6 +37,7 @@ import java.util.Collections; import java.util.List; import java.util.Map; +import java.util.Objects; import java.util.regex.Pattern; import static org.apache.kafka.common.config.ConfigDef.Range.atLeast; @@ -379,6 +380,20 @@ private void logDeprecatedProperty(String propName, String propValue, String def } } + private void logPluginPathConfigProviderWarning(Map rawOriginals) { + String rawPluginPath = rawOriginals.get(PLUGIN_PATH_CONFIG); + String transformedPluginPath = originalsStrings().get(PLUGIN_PATH_CONFIG); + if (!Objects.equals(rawPluginPath, transformedPluginPath)) { + log.error( + "Config providers do not work with the plugin.path property. The raw value '{}' " + + "will be used for plugin scanning, as opposed to the transformed value '{}'. " + + "See https://issues.apache.org/jira/browse/KAFKA-9845 for more information.", + rawPluginPath, + transformedPluginPath + ); + } + } + public Integer getRebalanceTimeout() { return null; } @@ -398,6 +413,7 @@ public static List pluginLocations(Map props) { public WorkerConfig(ConfigDef definition, Map props) { super(definition, props); logInternalConverterDeprecationWarnings(props); + logPluginPathConfigProviderWarning(props); } private static class AdminListenersValidator implements ConfigDef.Validator { From 1e9dae6efad90e6588aeddf030253c569c458731 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Thu, 9 Apr 2020 22:26:17 -0700 Subject: [PATCH 4/7] KAFKA-9845: Demote log message level from ERROR to WARN Co-Authored-By: Nigel Liang --- .../java/org/apache/kafka/connect/runtime/WorkerConfig.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index 7b1d248164060..62f9383095579 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -384,7 +384,7 @@ private void logPluginPathConfigProviderWarning(Map rawOriginals String rawPluginPath = rawOriginals.get(PLUGIN_PATH_CONFIG); String transformedPluginPath = originalsStrings().get(PLUGIN_PATH_CONFIG); if (!Objects.equals(rawPluginPath, transformedPluginPath)) { - log.error( + log.warn( "Config providers do not work with the plugin.path property. The raw value '{}' " + "will be used for plugin scanning, as opposed to the transformed value '{}'. " + "See https://issues.apache.org/jira/browse/KAFKA-9845 for more information.", From a2f339c438e1c91ef7b646941fc1ae584d0b2cda Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Fri, 10 Apr 2020 10:54:21 -0700 Subject: [PATCH 5/7] KAFKA-94845: Fix failing unit tests --- .../java/org/apache/kafka/connect/runtime/WorkerConfig.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index 62f9383095579..0605a77ec4fbd 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -382,7 +382,9 @@ private void logDeprecatedProperty(String propName, String propValue, String def private void logPluginPathConfigProviderWarning(Map rawOriginals) { String rawPluginPath = rawOriginals.get(PLUGIN_PATH_CONFIG); - String transformedPluginPath = originalsStrings().get(PLUGIN_PATH_CONFIG); + // Can't use AbstractConfig::originalsStrings here since some values may be null, which + // causes that method to fail + String transformedPluginPath = Objects.toString(originals().get(PLUGIN_PATH_CONFIG)); if (!Objects.equals(rawPluginPath, transformedPluginPath)) { log.warn( "Config providers do not work with the plugin.path property. The raw value '{}' " From 5d8af18e03e17361d430727a8b44e3f7a4077af4 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Fri, 10 Apr 2020 11:35:34 -0700 Subject: [PATCH 6/7] KAFKA-9845: Add warning message to docstring for plugin.path config --- .../java/org/apache/kafka/connect/runtime/WorkerConfig.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index 0605a77ec4fbd..b7fe6941a95d6 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -206,7 +206,9 @@ public class WorkerConfig extends AbstractConfig { + "plugins and their dependencies\n" + "Note: symlinks will be followed to discover dependencies or plugins.\n" + "Examples: plugin.path=/usr/local/share/java,/usr/local/share/kafka/plugins," - + "/opt/connectors"; + + "/opt/connectors\n" + + "Warning: Config providers will not take effect if used for the value of this " + + "property, and instead the raw, non-transformed value will be used."; public static final String CONFIG_PROVIDERS_CONFIG = "config.providers"; protected static final String CONFIG_PROVIDERS_DOC = From ccec713d0bbbd5c6ab85f76dcbfd2592c57c2158 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Wed, 10 Jun 2020 19:54:49 -0700 Subject: [PATCH 7/7] KAFKA-9845: Apply suggestions from code review Co-authored-by: Randall Hauch --- .../apache/kafka/connect/runtime/WorkerConfig.java | 12 +++++++----- 1 file changed, 7 insertions(+), 5 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index b7fe6941a95d6..c6e3bdc375c3b 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -207,8 +207,9 @@ public class WorkerConfig extends AbstractConfig { + "Note: symlinks will be followed to discover dependencies or plugins.\n" + "Examples: plugin.path=/usr/local/share/java,/usr/local/share/kafka/plugins," + "/opt/connectors\n" - + "Warning: Config providers will not take effect if used for the value of this " - + "property, and instead the raw, non-transformed value will be used."; + + "Do not use config provider variables in this property, since the raw path is used " + + "by the worker's scanner before config providers are initialized and used to " + + "replace variables."; public static final String CONFIG_PROVIDERS_CONFIG = "config.providers"; protected static final String CONFIG_PROVIDERS_DOC = @@ -389,9 +390,10 @@ private void logPluginPathConfigProviderWarning(Map rawOriginals String transformedPluginPath = Objects.toString(originals().get(PLUGIN_PATH_CONFIG)); if (!Objects.equals(rawPluginPath, transformedPluginPath)) { log.warn( - "Config providers do not work with the plugin.path property. The raw value '{}' " - + "will be used for plugin scanning, as opposed to the transformed value '{}'. " - + "See https://issues.apache.org/jira/browse/KAFKA-9845 for more information.", + "Variables cannot be used in the 'plugin.path' property, since the property is " + + "used by plugin scanning before the config providers that replace the " + + "variables are initialized. The raw value '{}' was used for plugin scanning, as " + + "opposed to the transformed value '{}', and this may cause unexpected results.", rawPluginPath, transformedPluginPath );