From 840f2c9a0c3188f75c0a8280be1cfbf24170c9e3 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Mon, 28 Aug 2023 11:46:59 -0400 Subject: [PATCH 01/15] KAFKA-13329: Add preflight validation for connector key and value converter classes --- .../connect/runtime/ConnectorConfig.java | 15 +++- .../util/ConcreteSubClassValidator.java | 68 ++++++++++++++ .../util/InstantiableClassValidator.java | 51 +++++++++++ .../ConnectorValidationIntegrationTest.java | 90 +++++++++++++++++++ 4 files changed, 222 insertions(+), 2 deletions(-) create mode 100644 connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java create mode 100644 connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java index 4f20a32f81b22..6902de2e0e5d4 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java @@ -28,8 +28,11 @@ import org.apache.kafka.connect.runtime.errors.ToleranceType; import org.apache.kafka.connect.runtime.isolation.PluginDesc; import org.apache.kafka.connect.runtime.isolation.Plugins; +import org.apache.kafka.connect.storage.Converter; import org.apache.kafka.connect.transforms.Transformation; import org.apache.kafka.connect.transforms.predicates.Predicate; +import org.apache.kafka.connect.util.ConcreteSubClassValidator; +import org.apache.kafka.connect.util.InstantiableClassValidator; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -82,10 +85,18 @@ public class ConnectorConfig extends AbstractConfig { public static final String KEY_CONVERTER_CLASS_CONFIG = WorkerConfig.KEY_CONVERTER_CLASS_CONFIG; public static final String KEY_CONVERTER_CLASS_DOC = WorkerConfig.KEY_CONVERTER_CLASS_DOC; public static final String KEY_CONVERTER_CLASS_DISPLAY = "Key converter class"; + private static final ConfigDef.Validator KEY_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( + ConcreteSubClassValidator.forSuperClass(Converter.class), + new InstantiableClassValidator() + ); public static final String VALUE_CONVERTER_CLASS_CONFIG = WorkerConfig.VALUE_CONVERTER_CLASS_CONFIG; public static final String VALUE_CONVERTER_CLASS_DOC = WorkerConfig.VALUE_CONVERTER_CLASS_DOC; public static final String VALUE_CONVERTER_CLASS_DISPLAY = "Value converter class"; + private static final ConfigDef.Validator VALUE_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( + ConcreteSubClassValidator.forSuperClass(Converter.class), + new InstantiableClassValidator() + ); public static final String HEADER_CONVERTER_CLASS_CONFIG = WorkerConfig.HEADER_CONVERTER_CLASS_CONFIG; public static final String HEADER_CONVERTER_CLASS_DOC = WorkerConfig.HEADER_CONVERTER_CLASS_DOC; @@ -181,8 +192,8 @@ public static ConfigDef configDef() { .define(NAME_CONFIG, Type.STRING, ConfigDef.NO_DEFAULT_VALUE, nonEmptyStringWithoutControlChars(), Importance.HIGH, NAME_DOC, COMMON_GROUP, ++orderInGroup, Width.MEDIUM, NAME_DISPLAY) .define(CONNECTOR_CLASS_CONFIG, Type.STRING, Importance.HIGH, CONNECTOR_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.LONG, CONNECTOR_CLASS_DISPLAY) .define(TASKS_MAX_CONFIG, Type.INT, TASKS_MAX_DEFAULT, atLeast(TASKS_MIN_CONFIG), Importance.HIGH, TASKS_MAX_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, TASK_MAX_DISPLAY) - .define(KEY_CONVERTER_CLASS_CONFIG, Type.CLASS, null, Importance.LOW, KEY_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, KEY_CONVERTER_CLASS_DISPLAY) - .define(VALUE_CONVERTER_CLASS_CONFIG, Type.CLASS, null, Importance.LOW, VALUE_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, VALUE_CONVERTER_CLASS_DISPLAY) + .define(KEY_CONVERTER_CLASS_CONFIG, Type.CLASS, null, KEY_CONVERTER_CLASS_VALIDATOR, Importance.LOW, KEY_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, KEY_CONVERTER_CLASS_DISPLAY) + .define(VALUE_CONVERTER_CLASS_CONFIG, Type.CLASS, null, VALUE_CONVERTER_CLASS_VALIDATOR, Importance.LOW, VALUE_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, VALUE_CONVERTER_CLASS_DISPLAY) .define(HEADER_CONVERTER_CLASS_CONFIG, Type.CLASS, HEADER_CONVERTER_CLASS_DEFAULT, Importance.LOW, HEADER_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, HEADER_CONVERTER_CLASS_DISPLAY) .define(TRANSFORMS_CONFIG, Type.LIST, Collections.emptyList(), aliasValidator("transformation"), Importance.LOW, TRANSFORMS_DOC, TRANSFORMS_GROUP, ++orderInGroup, Width.LONG, TRANSFORMS_DISPLAY) .define(PREDICATES_CONFIG, Type.LIST, Collections.emptyList(), aliasValidator("predicate"), Importance.LOW, PREDICATES_DOC, PREDICATES_GROUP, ++orderInGroup, Width.LONG, PREDICATES_DISPLAY) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java new file mode 100644 index 0000000000000..07683681e566f --- /dev/null +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java @@ -0,0 +1,68 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.connect.util; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.common.config.ConfigException; +import org.apache.kafka.common.utils.Utils; + +import java.lang.reflect.Modifier; +import java.util.stream.Collectors; +import java.util.stream.Stream; + +public class ConcreteSubClassValidator implements ConfigDef.Validator { + private final Class expectedSuperClass; + + private ConcreteSubClassValidator(Class expectedSuperClass) { + this.expectedSuperClass = expectedSuperClass; + } + + public static ConcreteSubClassValidator forSuperClass(Class expectedSuperClass) { + return new ConcreteSubClassValidator(expectedSuperClass); + } + + @Override + public void ensureValid(String name, Object value) { + if (value == null) { + // The value will be null if the class couldn't be found; no point in performing follow-up validation + return; + } + + Class cls = (Class) value; + if (!expectedSuperClass.isAssignableFrom(cls)) { + throw new ConfigException(name, String.valueOf(cls), "Not a " + expectedSuperClass.getSimpleName()); + } + + if (Modifier.isAbstract(cls.getModifiers())) { + String childClassNames = Stream.of(cls.getClasses()) + .filter(cls::isAssignableFrom) + .filter(c -> !Modifier.isAbstract(c.getModifiers())) + .filter(c -> Modifier.isPublic(c.getModifiers())) + .map(Class::getName) + .collect(Collectors.joining(", ")); + String message = Utils.isBlank(childClassNames) ? + "Class is abstract and cannot be created." : + "Class is abstract and cannot be created. Did you mean " + childClassNames + "?"; + throw new ConfigException(name, cls.getName(), message); + } + } + + @Override + public String toString() { + return "A concrete subclass of " + expectedSuperClass.getName(); + } +} diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java new file mode 100644 index 0000000000000..bc5f05c9aae82 --- /dev/null +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java @@ -0,0 +1,51 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.kafka.connect.util; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.common.config.ConfigException; +import org.apache.kafka.common.utils.Utils; + +import java.io.Closeable; + +public class InstantiableClassValidator implements ConfigDef.Validator { + + @Override + public void ensureValid(String name, Object value) { + if (value == null) { + // The value will be null if the class couldn't be found; no point in performing follow-up validation + return; + } + + Class cls = (Class) value; + try { + Object o = cls.getDeclaredConstructor().newInstance(); + if (o instanceof Closeable) { + Utils.closeQuietly((Closeable) o, o + " (instantiated for preflight validation)"); + } + } catch (NoSuchMethodException e) { + throw new ConfigException(name, cls.getName(), "Could not find a public no-argument constructor for class" + (e.getMessage() != null ? ": " + e.getMessage() : "")); + } catch (ReflectiveOperationException | RuntimeException e) { + throw new ConfigException(name, cls.getName(), "Could not instantiate class" + (e.getMessage() != null ? ": " + e.getMessage() : "")); + } + } + + @Override + public String toString() { + return "A class with a public, no-argument constructor"; + } +} \ No newline at end of file diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java index 88646112ffcaf..491cd7801e127 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java @@ -16,7 +16,12 @@ */ package org.apache.kafka.connect.integration; +import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.data.Schema; +import org.apache.kafka.connect.data.SchemaAndValue; +import org.apache.kafka.connect.errors.ConnectException; +import org.apache.kafka.connect.storage.Converter; import org.apache.kafka.connect.storage.StringConverter; import org.apache.kafka.connect.transforms.Filter; import org.apache.kafka.connect.transforms.predicates.RecordIsTombstone; @@ -296,6 +301,55 @@ public void testConnectorHasMissingConverterClass() throws InterruptedException ); } + @Test + public void testConnectorHasInvalidConverterClassType() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(KEY_CONVERTER_CLASS_CONFIG, MonitorableSinkConnector.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a converter with a class of the wrong type is specified", + 0 + ); + } + + @Test + public void testConnectorHasAbstractConverter() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(KEY_CONVERTER_CLASS_CONFIG, AbstractTestConverter.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when an abstract converter class is specified" + ); + } + + @Test + public void testConnectorHasConverterWithNoSuitableConstructor() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(KEY_CONVERTER_CLASS_CONFIG, TestConverterWithPrivateConstructor.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a converter class with no suitable constructor is specified" + ); + } + + @Test + public void testConnectorHasConverterThatThrowsExceptionOnInstantiation() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(KEY_CONVERTER_CLASS_CONFIG, TestConverterWithConstructorThatThrowsException.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a converter class that throws an exception on instantiation is specified" + ); + } + @Test public void testConnectorHasMissingHeaderConverterClass() throws InterruptedException { Map config = defaultSinkConnectorProps(); @@ -309,6 +363,42 @@ public void testConnectorHasMissingHeaderConverterClass() throws InterruptedExce ); } + public static abstract class TestConverter implements Converter { + + @Override + public void configure(Map configs, boolean isKey) { + } + + @Override + public byte[] fromConnectData(String topic, Schema schema, Object value) { + return null; + } + + @Override + public SchemaAndValue toConnectData(String topic, byte[] value) { + return null; + } + + @Override + public ConfigDef config() { + return null; + } + } + + public static abstract class AbstractTestConverter extends TestConverter { + } + + public static class TestConverterWithPrivateConstructor extends TestConverter { + private TestConverterWithPrivateConstructor() { + } + } + + public static class TestConverterWithConstructorThatThrowsException extends TestConverter { + public TestConverterWithConstructorThatThrowsException() { + throw new ConnectException("whoops"); + } + } + private Map defaultSourceConnectorProps() { // setup up props for the source connector Map props = new HashMap<>(); From ac9caf8993605c9e7ca591e4e305e18026717cdf Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Mon, 28 Aug 2023 12:56:07 -0400 Subject: [PATCH 02/15] KAFKA-13328: Add preflight validation for connector header converter classes --- .../connect/runtime/ConnectorConfig.java | 7 +- .../ConnectorValidationIntegrationTest.java | 75 ++++++++++++++++++- 2 files changed, 79 insertions(+), 3 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java index 6902de2e0e5d4..655f3be87ce9d 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java @@ -29,6 +29,7 @@ import org.apache.kafka.connect.runtime.isolation.PluginDesc; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.storage.Converter; +import org.apache.kafka.connect.storage.HeaderConverter; import org.apache.kafka.connect.transforms.Transformation; import org.apache.kafka.connect.transforms.predicates.Predicate; import org.apache.kafka.connect.util.ConcreteSubClassValidator; @@ -104,6 +105,10 @@ public class ConnectorConfig extends AbstractConfig { // The Connector config should not have a default for the header converter, since the absence of a config property means that // the worker config settings should be used. Thus, we set the default to null here. public static final String HEADER_CONVERTER_CLASS_DEFAULT = null; + private static final ConfigDef.Validator HEADER_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( + ConcreteSubClassValidator.forSuperClass(HeaderConverter.class), + new InstantiableClassValidator() + ); public static final String TASKS_MAX_CONFIG = "tasks.max"; private static final String TASKS_MAX_DOC = "Maximum number of tasks to use for this connector."; @@ -194,7 +199,7 @@ public static ConfigDef configDef() { .define(TASKS_MAX_CONFIG, Type.INT, TASKS_MAX_DEFAULT, atLeast(TASKS_MIN_CONFIG), Importance.HIGH, TASKS_MAX_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, TASK_MAX_DISPLAY) .define(KEY_CONVERTER_CLASS_CONFIG, Type.CLASS, null, KEY_CONVERTER_CLASS_VALIDATOR, Importance.LOW, KEY_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, KEY_CONVERTER_CLASS_DISPLAY) .define(VALUE_CONVERTER_CLASS_CONFIG, Type.CLASS, null, VALUE_CONVERTER_CLASS_VALIDATOR, Importance.LOW, VALUE_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, VALUE_CONVERTER_CLASS_DISPLAY) - .define(HEADER_CONVERTER_CLASS_CONFIG, Type.CLASS, HEADER_CONVERTER_CLASS_DEFAULT, Importance.LOW, HEADER_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, HEADER_CONVERTER_CLASS_DISPLAY) + .define(HEADER_CONVERTER_CLASS_CONFIG, Type.CLASS, HEADER_CONVERTER_CLASS_DEFAULT, HEADER_CONVERTER_CLASS_VALIDATOR, Importance.LOW, HEADER_CONVERTER_CLASS_DOC, COMMON_GROUP, ++orderInGroup, Width.SHORT, HEADER_CONVERTER_CLASS_DISPLAY) .define(TRANSFORMS_CONFIG, Type.LIST, Collections.emptyList(), aliasValidator("transformation"), Importance.LOW, TRANSFORMS_DOC, TRANSFORMS_GROUP, ++orderInGroup, Width.LONG, TRANSFORMS_DISPLAY) .define(PREDICATES_CONFIG, Type.LIST, Collections.emptyList(), aliasValidator("predicate"), Importance.LOW, PREDICATES_DOC, PREDICATES_GROUP, ++orderInGroup, Width.LONG, PREDICATES_DISPLAY) .define(CONFIG_RELOAD_ACTION_CONFIG, Type.STRING, CONFIG_RELOAD_ACTION_RESTART, diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java index 491cd7801e127..b3e37c9eeeb80 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java @@ -22,6 +22,7 @@ import org.apache.kafka.connect.data.SchemaAndValue; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.storage.Converter; +import org.apache.kafka.connect.storage.HeaderConverter; import org.apache.kafka.connect.storage.StringConverter; import org.apache.kafka.connect.transforms.Filter; import org.apache.kafka.connect.transforms.predicates.RecordIsTombstone; @@ -32,6 +33,7 @@ import org.junit.Test; import org.junit.experimental.categories.Category; +import java.io.IOException; import java.util.HashMap; import java.util.Map; @@ -363,8 +365,63 @@ public void testConnectorHasMissingHeaderConverterClass() throws InterruptedExce ); } - public static abstract class TestConverter implements Converter { + @Test + public void testConnectorHasInvalidHeaderConverterClassType() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(HEADER_CONVERTER_CLASS_CONFIG, MonitorableSinkConnector.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a header converter with a class of the wrong type is specified" + ); + } + + @Test + public void testConnectorHasAbstractHeaderConverter() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(HEADER_CONVERTER_CLASS_CONFIG, AbstractTestConverter.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when an abstract header converter class is specified" + ); + } + + @Test + public void testConnectorHasHeaderConverterWithNoSuitableConstructor() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(HEADER_CONVERTER_CLASS_CONFIG, TestConverterWithPrivateConstructor.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a header converter class with no suitable constructor is specified" + ); + } + @Test + public void testConnectorHasHeaderConverterThatThrowsExceptionOnInstantiation() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(HEADER_CONVERTER_CLASS_CONFIG, TestConverterWithConstructorThatThrowsException.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a header converter class that throws an exception on instantiation is specified" + ); + } + + public static abstract class TestConverter implements Converter, HeaderConverter { + + // Defined by both Converter and HeaderConverter interfaces + @Override + public ConfigDef config() { + return null; + } + + // Defined by Converter interface @Override public void configure(Map configs, boolean isKey) { } @@ -379,10 +436,24 @@ public SchemaAndValue toConnectData(String topic, byte[] value) { return null; } + // Defined by HeaderConverter interface @Override - public ConfigDef config() { + public void close() throws IOException { + } + + @Override + public void configure(Map configs) { + } + + @Override + public SchemaAndValue toConnectHeader(String topic, String headerKey, byte[] value) { return null; } + + @Override + public byte[] fromConnectHeader(String topic, String headerKey, Schema schema, Object value) { + return new byte[0]; + } } public static abstract class AbstractTestConverter extends TestConverter { From dba0649a4fb7c2dc0c227df4624015e896c2d4e6 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 29 Aug 2023 15:37:48 -0400 Subject: [PATCH 03/15] Add utility to close Object instances that may be AutoCloseable --- .../org/apache/kafka/common/utils/Utils.java | 17 +++++++++++++++++ .../util/InstantiableClassValidator.java | 6 +----- 2 files changed, 18 insertions(+), 5 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 6a0913d3c2da1..ea3207a2ec8b4 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -1097,6 +1097,23 @@ public interface UncheckedCloseable extends AutoCloseable { void close(); } + /** + * Closes {@code maybeCloseable} if it implements the {@link AutoCloseable} interface, + * and if an exception is thrown, it is logged at the WARN level. + * Be cautious when passing method references as an argument. For example: + *

+ * {@code closeQuietly(task::stop, "source task");} + *

+ * Although this method gracefully handles null {@link AutoCloseable} objects, attempts to take a method + * reference from a null object will result in a {@link NullPointerException}. In the example code above, + * it would be the caller's responsibility to ensure that {@code task} was non-null before attempting to + * use a method reference from it. + */ + public static void maybeCloseQuietly(Object maybeCloseable, String name) { + if (maybeCloseable instanceof AutoCloseable) + closeQuietly((AutoCloseable) maybeCloseable, name); + } + /** * Closes {@code closeable} and if an exception is thrown, it is logged at the WARN level. * Be cautious when passing method references as an argument. For example: diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java index bc5f05c9aae82..4af81238c088f 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java @@ -20,8 +20,6 @@ import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.utils.Utils; -import java.io.Closeable; - public class InstantiableClassValidator implements ConfigDef.Validator { @Override @@ -34,9 +32,7 @@ public void ensureValid(String name, Object value) { Class cls = (Class) value; try { Object o = cls.getDeclaredConstructor().newInstance(); - if (o instanceof Closeable) { - Utils.closeQuietly((Closeable) o, o + " (instantiated for preflight validation)"); - } + Utils.maybeCloseQuietly(o, o + " (instantiated for preflight validation"); } catch (NoSuchMethodException e) { throw new ConfigException(name, cls.getName(), "Could not find a public no-argument constructor for class" + (e.getMessage() != null ? ": " + e.getMessage() : "")); } catch (ReflectiveOperationException | RuntimeException e) { From abdc024f776685055e5e817f67421bc1ab0d617b Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Thu, 31 Aug 2023 12:09:56 -0400 Subject: [PATCH 04/15] Update clients/src/main/java/org/apache/kafka/common/utils/Utils.java Co-authored-by: Yash Mayya --- clients/src/main/java/org/apache/kafka/common/utils/Utils.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index ea3207a2ec8b4..6f9778aeb19d8 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -1102,7 +1102,7 @@ public interface UncheckedCloseable extends AutoCloseable { * and if an exception is thrown, it is logged at the WARN level. * Be cautious when passing method references as an argument. For example: *

- * {@code closeQuietly(task::stop, "source task");} + * {@code maybeCloseQuietly(task::stop, "source task");} *

* Although this method gracefully handles null {@link AutoCloseable} objects, attempts to take a method * reference from a null object will result in a {@link NullPointerException}. In the example code above, From 1d10efab9604e762f781eee8db802b4607d52ccc Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Thu, 31 Aug 2023 12:45:28 -0400 Subject: [PATCH 05/15] Add newline to end of file, fix Javadoc code snippet example, extract common logic into Utils::ensureConcrete --- .../org/apache/kafka/common/utils/Utils.java | 27 +++++++++++++++++++ .../connect/runtime/ConnectorConfig.java | 22 ++++----------- .../util/ConcreteSubClassValidator.java | 25 +++++------------ .../util/InstantiableClassValidator.java | 2 +- 4 files changed, 39 insertions(+), 37 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 6f9778aeb19d8..9afcef93494f9 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -16,6 +16,7 @@ */ package org.apache.kafka.common.utils; +import java.lang.reflect.Modifier; import java.nio.BufferUnderflowException; import java.nio.ByteOrder; import java.nio.file.StandardOpenOption; @@ -1582,6 +1583,32 @@ public static String[] enumOptions(Class> enumClass) { .toArray(String[]::new); } + /** + * Ensure that the class is concrete (i.e., not abstract). If it is, throw a {@link ConfigException} + * with a friendly error message suggesting a list of concrete child subclasses (if any are known). + * @param cls the class to check; may not be null + * @param name the name of the type of class to use in the error message; e.g., "Transform", + * "Interceptor", or even just "Class"; may be null + * @throws ConfigException if the class is not concrete + */ + public static void ensureConcrete(Class cls, String name) { + Objects.requireNonNull(cls); + if (isBlank(name)) + name = "Class"; + if (Modifier.isAbstract(cls.getModifiers())) { + String childClassNames = Stream.of(cls.getClasses()) + .filter(cls::isAssignableFrom) + .filter(c -> !Modifier.isAbstract(c.getModifiers())) + .filter(c -> Modifier.isPublic(c.getModifiers())) + .map(Class::getName) + .collect(Collectors.joining(", ")); + String message = Utils.isBlank(childClassNames) ? + name + " is abstract and cannot be created." : + name + " is abstract and cannot be created. Did you mean " + childClassNames + "?"; + throw new ConfigException(name, cls.getName(), message); + } + } + /** * Convert time instant to readable string for logging * @param timestamp the timestamp of the instant to be converted. diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java index 655f3be87ce9d..a9d37ff91f6e3 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java @@ -37,7 +37,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; -import java.lang.reflect.Modifier; import java.util.ArrayList; import java.util.Collections; import java.util.HashSet; @@ -46,7 +45,6 @@ import java.util.Locale; import java.util.Map; import java.util.Set; -import java.util.stream.Collectors; import java.util.stream.Stream; import static org.apache.kafka.common.config.ConfigDef.NonEmptyStringWithoutControlChars.nonEmptyStringWithoutControlChars; @@ -87,7 +85,7 @@ public class ConnectorConfig extends AbstractConfig { public static final String KEY_CONVERTER_CLASS_DOC = WorkerConfig.KEY_CONVERTER_CLASS_DOC; public static final String KEY_CONVERTER_CLASS_DISPLAY = "Key converter class"; private static final ConfigDef.Validator KEY_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( - ConcreteSubClassValidator.forSuperClass(Converter.class), + ConcreteSubClassValidator.forSuperClass(Converter.class, "Key converter"), new InstantiableClassValidator() ); @@ -95,7 +93,7 @@ public class ConnectorConfig extends AbstractConfig { public static final String VALUE_CONVERTER_CLASS_DOC = WorkerConfig.VALUE_CONVERTER_CLASS_DOC; public static final String VALUE_CONVERTER_CLASS_DISPLAY = "Value converter class"; private static final ConfigDef.Validator VALUE_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( - ConcreteSubClassValidator.forSuperClass(Converter.class), + ConcreteSubClassValidator.forSuperClass(Converter.class, "Value converter"), new InstantiableClassValidator() ); @@ -106,7 +104,7 @@ public class ConnectorConfig extends AbstractConfig { // the worker config settings should be used. Thus, we set the default to null here. public static final String HEADER_CONVERTER_CLASS_DEFAULT = null; private static final ConfigDef.Validator HEADER_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( - ConcreteSubClassValidator.forSuperClass(HeaderConverter.class), + ConcreteSubClassValidator.forSuperClass(HeaderConverter.class, "Header converter"), new InstantiableClassValidator() ); @@ -510,18 +508,8 @@ ConfigDef getConfigDefFromConfigProvidingClass(String key, Class cls) { if (cls == null || !baseClass.isAssignableFrom(cls)) { throw new ConfigException(key, String.valueOf(cls), "Not a " + baseClass.getSimpleName()); } - if (Modifier.isAbstract(cls.getModifiers())) { - String childClassNames = Stream.of(cls.getClasses()) - .filter(cls::isAssignableFrom) - .filter(c -> !Modifier.isAbstract(c.getModifiers())) - .filter(c -> Modifier.isPublic(c.getModifiers())) - .map(Class::getName) - .collect(Collectors.joining(", ")); - String message = Utils.isBlank(childClassNames) ? - aliasKind + " is abstract and cannot be created." : - aliasKind + " is abstract and cannot be created. Did you mean " + childClassNames + "?"; - throw new ConfigException(key, String.valueOf(cls), message); - } + Utils.ensureConcrete(cls, aliasKind); + T transformation; try { transformation = Utils.newInstance(cls, baseClass); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java index 07683681e566f..e90dc3488f13b 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java @@ -20,19 +20,17 @@ import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.utils.Utils; -import java.lang.reflect.Modifier; -import java.util.stream.Collectors; -import java.util.stream.Stream; - public class ConcreteSubClassValidator implements ConfigDef.Validator { private final Class expectedSuperClass; + private final String superClassName; - private ConcreteSubClassValidator(Class expectedSuperClass) { + private ConcreteSubClassValidator(Class expectedSuperClass, String superClassName) { this.expectedSuperClass = expectedSuperClass; + this.superClassName = superClassName; } - public static ConcreteSubClassValidator forSuperClass(Class expectedSuperClass) { - return new ConcreteSubClassValidator(expectedSuperClass); + public static ConcreteSubClassValidator forSuperClass(Class expectedSuperClass, String superClassName) { + return new ConcreteSubClassValidator(expectedSuperClass, superClassName); } @Override @@ -47,18 +45,7 @@ public void ensureValid(String name, Object value) { throw new ConfigException(name, String.valueOf(cls), "Not a " + expectedSuperClass.getSimpleName()); } - if (Modifier.isAbstract(cls.getModifiers())) { - String childClassNames = Stream.of(cls.getClasses()) - .filter(cls::isAssignableFrom) - .filter(c -> !Modifier.isAbstract(c.getModifiers())) - .filter(c -> Modifier.isPublic(c.getModifiers())) - .map(Class::getName) - .collect(Collectors.joining(", ")); - String message = Utils.isBlank(childClassNames) ? - "Class is abstract and cannot be created." : - "Class is abstract and cannot be created. Did you mean " + childClassNames + "?"; - throw new ConfigException(name, cls.getName(), message); - } + Utils.ensureConcrete(cls, superClassName); } @Override diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java index 4af81238c088f..ee31522154b49 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java @@ -44,4 +44,4 @@ public void ensureValid(String name, Object value) { public String toString() { return "A class with a public, no-argument constructor"; } -} \ No newline at end of file +} From 735d35fa2669dc488a2775a280ba122f52c741f5 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 5 Sep 2023 15:22:25 -0400 Subject: [PATCH 06/15] Remove warning about method references from Javadocs for Utils::maybeCloseQuietly --- .../main/java/org/apache/kafka/common/utils/Utils.java | 8 -------- 1 file changed, 8 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 9afcef93494f9..37d2365dbc01b 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -1101,14 +1101,6 @@ public interface UncheckedCloseable extends AutoCloseable { /** * Closes {@code maybeCloseable} if it implements the {@link AutoCloseable} interface, * and if an exception is thrown, it is logged at the WARN level. - * Be cautious when passing method references as an argument. For example: - *

- * {@code maybeCloseQuietly(task::stop, "source task");} - *

- * Although this method gracefully handles null {@link AutoCloseable} objects, attempts to take a method - * reference from a null object will result in a {@link NullPointerException}. In the example code above, - * it would be the caller's responsibility to ensure that {@code task} was non-null before attempting to - * use a method reference from it. */ public static void maybeCloseQuietly(Object maybeCloseable, String name) { if (maybeCloseable instanceof AutoCloseable) From a554695029163abddac5126db0544e8b62bc0585 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 5 Sep 2023 15:47:13 -0400 Subject: [PATCH 07/15] Simplify Utils::ensureConcrete, refine exception handling in InstantiableClassValidator::ensureValid --- .../org/apache/kafka/common/utils/Utils.java | 24 ++++++++----------- .../connect/runtime/ConnectorConfig.java | 8 +++---- .../util/ConcreteSubClassValidator.java | 10 ++++---- .../util/InstantiableClassValidator.java | 4 ++-- 4 files changed, 20 insertions(+), 26 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 37d2365dbc01b..623da4cedd53c 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -1578,26 +1578,22 @@ public static String[] enumOptions(Class> enumClass) { /** * Ensure that the class is concrete (i.e., not abstract). If it is, throw a {@link ConfigException} * with a friendly error message suggesting a list of concrete child subclasses (if any are known). - * @param cls the class to check; may not be null - * @param name the name of the type of class to use in the error message; e.g., "Transform", - * "Interceptor", or even just "Class"; may be null + * @param klass the class to check; may not be null * @throws ConfigException if the class is not concrete */ - public static void ensureConcrete(Class cls, String name) { - Objects.requireNonNull(cls); - if (isBlank(name)) - name = "Class"; - if (Modifier.isAbstract(cls.getModifiers())) { - String childClassNames = Stream.of(cls.getClasses()) - .filter(cls::isAssignableFrom) + public static void ensureConcrete(Class klass) { + Objects.requireNonNull(klass); + if (Modifier.isAbstract(klass.getModifiers())) { + String childClassNames = Stream.of(klass.getClasses()) + .filter(klass::isAssignableFrom) .filter(c -> !Modifier.isAbstract(c.getModifiers())) .filter(c -> Modifier.isPublic(c.getModifiers())) .map(Class::getName) .collect(Collectors.joining(", ")); - String message = Utils.isBlank(childClassNames) ? - name + " is abstract and cannot be created." : - name + " is abstract and cannot be created. Did you mean " + childClassNames + "?"; - throw new ConfigException(name, cls.getName(), message); + String message = "This class is abstract and cannot be created."; + if (!Utils.isBlank(childClassNames)) + message += " Did you mean " + childClassNames + "?"; + throw new ConfigException(message); } } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java index a9d37ff91f6e3..e73e9a43c5064 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java @@ -85,7 +85,7 @@ public class ConnectorConfig extends AbstractConfig { public static final String KEY_CONVERTER_CLASS_DOC = WorkerConfig.KEY_CONVERTER_CLASS_DOC; public static final String KEY_CONVERTER_CLASS_DISPLAY = "Key converter class"; private static final ConfigDef.Validator KEY_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( - ConcreteSubClassValidator.forSuperClass(Converter.class, "Key converter"), + ConcreteSubClassValidator.forSuperClass(Converter.class), new InstantiableClassValidator() ); @@ -93,7 +93,7 @@ public class ConnectorConfig extends AbstractConfig { public static final String VALUE_CONVERTER_CLASS_DOC = WorkerConfig.VALUE_CONVERTER_CLASS_DOC; public static final String VALUE_CONVERTER_CLASS_DISPLAY = "Value converter class"; private static final ConfigDef.Validator VALUE_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( - ConcreteSubClassValidator.forSuperClass(Converter.class, "Value converter"), + ConcreteSubClassValidator.forSuperClass(Converter.class), new InstantiableClassValidator() ); @@ -104,7 +104,7 @@ public class ConnectorConfig extends AbstractConfig { // the worker config settings should be used. Thus, we set the default to null here. public static final String HEADER_CONVERTER_CLASS_DEFAULT = null; private static final ConfigDef.Validator HEADER_CONVERTER_CLASS_VALIDATOR = ConfigDef.CompositeValidator.of( - ConcreteSubClassValidator.forSuperClass(HeaderConverter.class, "Header converter"), + ConcreteSubClassValidator.forSuperClass(HeaderConverter.class), new InstantiableClassValidator() ); @@ -508,7 +508,7 @@ ConfigDef getConfigDefFromConfigProvidingClass(String key, Class cls) { if (cls == null || !baseClass.isAssignableFrom(cls)) { throw new ConfigException(key, String.valueOf(cls), "Not a " + baseClass.getSimpleName()); } - Utils.ensureConcrete(cls, aliasKind); + Utils.ensureConcrete(cls); T transformation; try { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java index e90dc3488f13b..af62c0d8d2b57 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java @@ -22,15 +22,13 @@ public class ConcreteSubClassValidator implements ConfigDef.Validator { private final Class expectedSuperClass; - private final String superClassName; - private ConcreteSubClassValidator(Class expectedSuperClass, String superClassName) { + private ConcreteSubClassValidator(Class expectedSuperClass) { this.expectedSuperClass = expectedSuperClass; - this.superClassName = superClassName; } - public static ConcreteSubClassValidator forSuperClass(Class expectedSuperClass, String superClassName) { - return new ConcreteSubClassValidator(expectedSuperClass, superClassName); + public static ConcreteSubClassValidator forSuperClass(Class expectedSuperClass) { + return new ConcreteSubClassValidator(expectedSuperClass); } @Override @@ -45,7 +43,7 @@ public void ensureValid(String name, Object value) { throw new ConfigException(name, String.valueOf(cls), "Not a " + expectedSuperClass.getSimpleName()); } - Utils.ensureConcrete(cls, superClassName); + Utils.ensureConcrete(cls); } @Override diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java index ee31522154b49..796be4ed48339 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/InstantiableClassValidator.java @@ -33,9 +33,9 @@ public void ensureValid(String name, Object value) { try { Object o = cls.getDeclaredConstructor().newInstance(); Utils.maybeCloseQuietly(o, o + " (instantiated for preflight validation"); - } catch (NoSuchMethodException e) { + } catch (NoSuchMethodException | IllegalAccessException e) { throw new ConfigException(name, cls.getName(), "Could not find a public no-argument constructor for class" + (e.getMessage() != null ? ": " + e.getMessage() : "")); - } catch (ReflectiveOperationException | RuntimeException e) { + } catch (ReflectiveOperationException | LinkageError | RuntimeException e) { throw new ConfigException(name, cls.getName(), "Could not instantiate class" + (e.getMessage() != null ? ": " + e.getMessage() : "")); } } From 8dfeca5033e5e49233289fe1227872baf1694159 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 5 Sep 2023 16:09:53 -0400 Subject: [PATCH 08/15] Fix failing unit tests --- .../connect/runtime/ConnectorConfigTest.java | 36 ++++++++----------- 1 file changed, 15 insertions(+), 21 deletions(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java index c113039912e41..6ded50ff7f308 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java @@ -198,13 +198,10 @@ public void abstractTransform() { props.put("connector.class", TestConnector.class.getName()); props.put("transforms", "a"); props.put("transforms.a.type", AbstractTransformation.class.getName()); - try { - new ConnectorConfig(MOCK_PLUGINS, props); - } catch (ConfigException ex) { - assertTrue( - ex.getMessage().contains("Transformation is abstract and cannot be created.") - ); - } + ConfigException ex = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue( + ex.getMessage().contains("This class is abstract and cannot be created.") + ); } @Test public void abstractKeyValueTransform() { @@ -213,19 +210,16 @@ public void abstractKeyValueTransform() { props.put("connector.class", TestConnector.class.getName()); props.put("transforms", "a"); props.put("transforms.a.type", AbstractKeyValueTransformation.class.getName()); - try { - new ConnectorConfig(MOCK_PLUGINS, props); - } catch (ConfigException ex) { - assertTrue( - ex.getMessage().contains("Transformation is abstract and cannot be created.") - ); - assertTrue( - ex.getMessage().contains(AbstractKeyValueTransformation.Key.class.getName()) - ); - assertTrue( - ex.getMessage().contains(AbstractKeyValueTransformation.Value.class.getName()) - ); - } + ConfigException ex = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue( + ex.getMessage().contains("This class is abstract and cannot be created.") + ); + assertTrue( + ex.getMessage().contains(AbstractKeyValueTransformation.Key.class.getName()) + ); + assertTrue( + ex.getMessage().contains(AbstractKeyValueTransformation.Value.class.getName()) + ); } @Test @@ -287,7 +281,7 @@ public void abstractPredicate() { props.put("predicates.my-pred.type", AbstractTestPredicate.class.getName()); props.put("predicates.my-pred.int", "84"); ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); - assertTrue(e.getMessage().contains("Predicate is abstract and cannot be created")); + assertTrue(e.getMessage().contains("This class is abstract and cannot be created")); } private void assertTransformationStageWithPredicate(Map props, boolean expectedNegated) { From 6cf1e131d1df339cd92d2050052eb7d193e05786 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 31 Oct 2023 14:53:53 -0400 Subject: [PATCH 09/15] Change Utils::ensureConcrete to Utils::ensureConcreteSubclass --- .../org/apache/kafka/common/utils/Utils.java | 16 +++++++++++++--- .../kafka/connect/runtime/ConnectorConfig.java | 6 +++--- .../connect/util/ConcreteSubClassValidator.java | 7 +------ .../connect/runtime/ConnectorConfigTest.java | 2 +- 4 files changed, 18 insertions(+), 13 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 623da4cedd53c..a8cd5a57540ca 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -1576,16 +1576,26 @@ public static String[] enumOptions(Class> enumClass) { } /** - * Ensure that the class is concrete (i.e., not abstract). If it is, throw a {@link ConfigException} + * Ensure that the class is concrete (i.e., not abstract), and that it subclasses a given base class. + * If it is abstract or does not subclass the given base class, throw a {@link ConfigException} * with a friendly error message suggesting a list of concrete child subclasses (if any are known). + * @param baseClass the expected superclass; may not be null * @param klass the class to check; may not be null * @throws ConfigException if the class is not concrete */ - public static void ensureConcrete(Class klass) { + public static void ensureConcreteSubclass(Class baseClass, Class klass) { + Objects.requireNonNull(baseClass); Objects.requireNonNull(klass); + + if (!baseClass.isAssignableFrom(klass)) { + String inheritFrom = baseClass.isInterface() ? "implement" : "extend"; + String baseClassType = baseClass.isInterface() ? "interface" : "class"; + throw new ConfigException("Class " + klass + " does not " + inheritFrom + " the " + baseClass.getSimpleName() + " " + baseClassType); + } + if (Modifier.isAbstract(klass.getModifiers())) { String childClassNames = Stream.of(klass.getClasses()) - .filter(klass::isAssignableFrom) + .filter(baseClass::isAssignableFrom) .filter(c -> !Modifier.isAbstract(c.getModifiers())) .filter(c -> Modifier.isPublic(c.getModifiers())) .map(Class::getName) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java index e73e9a43c5064..f33f00e40efa3 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java @@ -505,10 +505,10 @@ protected ConfigDef initialConfigDef() { * @param cls The subclass of the baseclass. */ ConfigDef getConfigDefFromConfigProvidingClass(String key, Class cls) { - if (cls == null || !baseClass.isAssignableFrom(cls)) { - throw new ConfigException(key, String.valueOf(cls), "Not a " + baseClass.getSimpleName()); + if (cls == null) { + throw new ConfigException(key, null, "Not a " + baseClass.getSimpleName()); } - Utils.ensureConcrete(cls); + Utils.ensureConcreteSubclass(baseClass, cls); T transformation; try { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java index af62c0d8d2b57..cbb4dedb126db 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConcreteSubClassValidator.java @@ -17,7 +17,6 @@ package org.apache.kafka.connect.util; import org.apache.kafka.common.config.ConfigDef; -import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.utils.Utils; public class ConcreteSubClassValidator implements ConfigDef.Validator { @@ -39,11 +38,7 @@ public void ensureValid(String name, Object value) { } Class cls = (Class) value; - if (!expectedSuperClass.isAssignableFrom(cls)) { - throw new ConfigException(name, String.valueOf(cls), "Not a " + expectedSuperClass.getSimpleName()); - } - - Utils.ensureConcrete(cls); + Utils.ensureConcreteSubclass(expectedSuperClass, cls); } @Override diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java index 6ded50ff7f308..f4e890fcaf362 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ConnectorConfigTest.java @@ -234,7 +234,7 @@ public void wrongPredicateType() { props.put("predicates", "my-pred"); props.put("predicates.my-pred.type", TestConnector.class.getName()); ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); - assertTrue(e.getMessage().contains("Not a Predicate")); + assertEquals("Class " + TestConnector.class + " does not implement the Predicate interface", e.getMessage()); } @Test From 29c95bcb472e1bf958e0a904055f75c6f005057e Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 29 Aug 2023 10:42:39 -0400 Subject: [PATCH 10/15] KAFKA-13328, KAFKA-13329: Add custom preflight validation support for connector header, key, and value converters --- checkstyle/suppressions.xml | 2 +- .../kafka/connect/runtime/AbstractHerder.java | 177 +++++++++++++++--- .../ConnectorValidationIntegrationTest.java | 90 ++++++++- ...org.apache.kafka.connect.storage.Converter | 4 +- ...ache.kafka.connect.storage.HeaderConverter | 4 +- 5 files changed, 244 insertions(+), 33 deletions(-) diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml index e1e746c755afc..3b890a4f84421 100644 --- a/checkstyle/suppressions.xml +++ b/checkstyle/suppressions.xml @@ -159,7 +159,7 @@ files="(KafkaConfigBackingStore|Values|ConnectMetricsRegistry).java"/> + files="(DistributedHerder|AbstractHerder|RestClient|RestServer|JsonConverter|KafkaConfigBackingStore|FileStreamSourceTask|WorkerSourceTask|TopicAdmin).java"/> validateSourceConnectorConfig(SourceConnector return configDef.validateAll(config); } + private ConfigInfos validateConverterConfig( + Map connectorConfig, + ConfigValue converterConfigValue, + Class converterInterface, + Function configDefAccessor, + String converterName, + String converterProperty, + ConverterType converterType + ) { + String converterClass = connectorConfig.get(converterProperty); + + if (converterClass == null + || converterConfigValue == null + || !converterConfigValue.errorMessages().isEmpty() + ) { + // Either no custom converter was specified, or one was specified but there's a problem with it. + // No need to proceed any further. + return null; + } + + T converterInstance; + try { + converterInstance = Utils.newInstance(converterClass, converterInterface); + } catch (ClassNotFoundException | RuntimeException e) { + log.error("Failed to instantiate {} class {}; this should have been caught by prior validation logic", converterName, converterClass, e); + converterConfigValue.addErrorMessage("Failed to load class " + converterClass + (e.getMessage() != null ? ": " + e.getMessage() : "")); + return null; + } + + try (Utils.UncheckedCloseable close = () -> Utils.maybeCloseQuietly(converterInstance, converterName + " " + converterClass);) { + ConfigDef configDef; + try { + configDef = configDefAccessor.apply(converterInstance); + } catch (RuntimeException e) { + log.error("Failed to load ConfigDef from {} of type {}", converterName, converterClass, e); + converterConfigValue.addErrorMessage("Failed to load ConfigDef from " + converterName + (e.getMessage() != null ? ": " + e.getMessage() : "")); + return null; + } + if (configDef == null) { + log.warn("{}.config() has returned a null ConfigDef; no further preflight config validation for this converter will be performed", converterClass); + // Older versions of Connect didn't do any converter validation. + // Even though converters are technically required to return a non-null ConfigDef object from their config() method, + // we permit this case in order to avoid breaking existing converters that, despite not adhering to this requirement, + // can be used successfully with a connector. + return null; + } + final String converterPrefix = converterProperty + "."; + Map converterConfig = connectorConfig.entrySet().stream() + .filter(e -> e.getKey().startsWith(converterPrefix)) + .collect(Collectors.toMap( + e -> e.getKey().substring(converterPrefix.length()), + Map.Entry::getValue + )); + converterConfig.putIfAbsent(ConverterConfig.TYPE_CONFIG, converterType.getName()); + + List configValues; + try { + configValues = configDef.validate(converterConfig); + } catch (RuntimeException e) { + log.error("Failed to perform custom config validation for {} of type {}", converterName, converterClass, e); + converterConfigValue.addErrorMessage("Failed to perform custom config validation for " + converterName + (e.getMessage() != null ? ": " + e.getMessage() : "")); + return null; + } + + return prefixedConfigInfos(configDef.configKeys(), configValues, converterPrefix); + } + } + + private ConfigInfos validateHeaderConverterConfig(Map connectorConfig, ConfigValue headerConverterConfigValue) { + return validateConverterConfig( + connectorConfig, + headerConverterConfigValue, + HeaderConverter.class, + HeaderConverter::config, + "header converter", + HEADER_CONVERTER_CLASS_CONFIG, + ConverterType.HEADER + ); + } + + private ConfigInfos validateKeyConverterConfig(Map connectorConfig, ConfigValue keyConverterConfigValue) { + return validateConverterConfig( + connectorConfig, + keyConverterConfigValue, + Converter.class, + Converter::config, + "key converter", + KEY_CONVERTER_CLASS_CONFIG, + ConverterType.KEY + ); + } + + private ConfigInfos validateValueConverterConfig(Map connectorConfig, ConfigValue valueConverterConfigValue) { + return validateConverterConfig( + connectorConfig, + valueConverterConfigValue, + Converter.class, + Converter::config, + "value converter", + VALUE_CONVERTER_CLASS_CONFIG, + ConverterType.VALUE + ); + } + @Override public void validateConnectorConfig(Map connectorProps, Callback callback) { validateConnectorConfig(connectorProps, callback, true); @@ -526,8 +637,13 @@ ConfigInfos validateConnectorConfig(Map connectorProps, boolean configKeys.putAll(configDef.configKeys()); allGroups.addAll(configDef.groups()); configValues.addAll(config.configValues()); - ConfigInfos configInfos = generateResult(connType, configKeys, configValues, new ArrayList<>(allGroups)); + // do custom converter-specific validation + ConfigInfos headerConverterConfigInfos = validateHeaderConverterConfig(connectorProps, validatedConnectorConfig.get(HEADER_CONVERTER_CLASS_CONFIG)); + ConfigInfos keyConverterConfigInfos = validateKeyConverterConfig(connectorProps, validatedConnectorConfig.get(KEY_CONVERTER_CLASS_CONFIG)); + ConfigInfos valueConverterConfigInfos = validateValueConverterConfig(connectorProps, validatedConnectorConfig.get(VALUE_CONVERTER_CLASS_CONFIG)); + + ConfigInfos configInfos = generateResult(connType, configKeys, configValues, new ArrayList<>(allGroups)); AbstractConfig connectorConfig = new AbstractConfig(new ConfigDef(), connectorProps, doLog); String connName = connectorProps.get(ConnectorConfig.NAME_CONFIG); ConfigInfos producerConfigInfos = null; @@ -567,7 +683,15 @@ ConfigInfos validateConnectorConfig(Map connectorProps, boolean ConnectorClientConfigRequest.ClientType.CONSUMER, connectorClientConfigOverridePolicy); } - return mergeConfigInfos(connType, configInfos, producerConfigInfos, consumerConfigInfos, adminConfigInfos); + return mergeConfigInfos(connType, + configInfos, + producerConfigInfos, + consumerConfigInfos, + adminConfigInfos, + headerConverterConfigInfos, + keyConverterConfigInfos, + valueConverterConfigInfos + ); } } @@ -593,10 +717,6 @@ private static ConfigInfos validateClientOverrides(String connName, org.apache.kafka.connect.health.ConnectorType connectorType, ConnectorClientConfigRequest.ClientType clientType, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) { - int errorCount = 0; - List configInfoList = new LinkedList<>(); - Map configKeys = configDef.configKeys(); - Set groups = new LinkedHashSet<>(); Map clientConfigs = new HashMap<>(); for (Map.Entry rawClientConfig : connectorConfig.originalsWithPrefix(prefix).entrySet()) { String configName = rawClientConfig.getKey(); @@ -610,27 +730,38 @@ private static ConfigInfos validateClientOverrides(String connName, ConnectorClientConfigRequest connectorClientConfigRequest = new ConnectorClientConfigRequest( connName, connectorType, connectorClass, clientConfigs, clientType); List configValues = connectorClientConfigOverridePolicy.validate(connectorClientConfigRequest); - if (configValues != null) { - for (ConfigValue validatedConfigValue : configValues) { - ConfigKey configKey = configKeys.get(validatedConfigValue.name()); - ConfigKeyInfo configKeyInfo = null; - if (configKey != null) { - if (configKey.group != null) { - groups.add(configKey.group); - } - configKeyInfo = convertConfigKey(configKey, prefix); - } - ConfigValue configValue = new ConfigValue(prefix + validatedConfigValue.name(), validatedConfigValue.value(), - validatedConfigValue.recommendedValues(), validatedConfigValue.errorMessages()); - if (configValue.errorMessages().size() > 0) { - errorCount++; + return prefixedConfigInfos(configDef.configKeys(), configValues, prefix); + } + + private static ConfigInfos prefixedConfigInfos(Map configKeys, List configValues, String prefix) { + int errorCount = 0; + Set groups = new LinkedHashSet<>(); + List configInfos = new ArrayList<>(); + + if (configValues == null) { + return new ConfigInfos("", 0, new ArrayList<>(groups), configInfos); + } + + for (ConfigValue validatedConfigValue : configValues) { + ConfigKey configKey = configKeys.get(validatedConfigValue.name()); + ConfigKeyInfo configKeyInfo = null; + if (configKey != null) { + if (configKey.group != null) { + groups.add(configKey.group); } - ConfigValueInfo configValueInfo = convertConfigValue(configValue, configKey != null ? configKey.type : null); - configInfoList.add(new ConfigInfo(configKeyInfo, configValueInfo)); + configKeyInfo = convertConfigKey(configKey, prefix); + } + + ConfigValue configValue = new ConfigValue(prefix + validatedConfigValue.name(), validatedConfigValue.value(), + validatedConfigValue.recommendedValues(), validatedConfigValue.errorMessages()); + if (configValue.errorMessages().size() > 0) { + errorCount++; } + ConfigValueInfo configValueInfo = convertConfigValue(configValue, configKey != null ? configKey.type : null); + configInfos.add(new ConfigInfo(configKeyInfo, configValueInfo)); } - return new ConfigInfos(connectorClass.toString(), errorCount, new ArrayList<>(groups), configInfoList); + return new ConfigInfos("", errorCount, new ArrayList<>(groups), configInfos); } // public for testing diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java index b3e37c9eeeb80..9fa939d1614f6 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java @@ -324,7 +324,8 @@ public void testConnectorHasAbstractConverter() throws InterruptedException { config.get(CONNECTOR_CLASS_CONFIG), config, 1, - "Connector config should fail preflight validation when an abstract converter class is specified" + "Connector config should fail preflight validation when an abstract converter class is specified", + 0 ); } @@ -336,7 +337,8 @@ public void testConnectorHasConverterWithNoSuitableConstructor() throws Interrup config.get(CONNECTOR_CLASS_CONFIG), config, 1, - "Connector config should fail preflight validation when a converter class with no suitable constructor is specified" + "Connector config should fail preflight validation when a converter class with no suitable constructor is specified", + 0 ); } @@ -348,7 +350,35 @@ public void testConnectorHasConverterThatThrowsExceptionOnInstantiation() throws config.get(CONNECTOR_CLASS_CONFIG), config, 1, - "Connector config should fail preflight validation when a converter class that throws an exception on instantiation is specified" + "Connector config should fail preflight validation when a converter class that throws an exception on instantiation is specified", + 0 + ); + } + + @Test + public void testConnectorHasMisconfiguredConverter() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(KEY_CONVERTER_CLASS_CONFIG, TestConverterWithSinglePropertyConfigDef.class.getName()); + config.put(KEY_CONVERTER_CLASS_CONFIG + "." + TestConverterWithSinglePropertyConfigDef.BOOLEAN_PROPERTY_NAME, "notaboolean"); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a converter fails custom validation", + 0 + ); + } + + @Test + public void testConnectorHasConverterWithNoConfigDef() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(KEY_CONVERTER_CLASS_CONFIG, TestConverterWithNoConfigDef.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 0, + "Connector config should not fail preflight validation even when a converter provides a null ConfigDef", + 0 ); } @@ -373,7 +403,8 @@ public void testConnectorHasInvalidHeaderConverterClassType() throws Interrupted config.get(CONNECTOR_CLASS_CONFIG), config, 1, - "Connector config should fail preflight validation when a header converter with a class of the wrong type is specified" + "Connector config should fail preflight validation when a header converter with a class of the wrong type is specified", + 0 ); } @@ -385,7 +416,8 @@ public void testConnectorHasAbstractHeaderConverter() throws InterruptedExceptio config.get(CONNECTOR_CLASS_CONFIG), config, 1, - "Connector config should fail preflight validation when an abstract header converter class is specified" + "Connector config should fail preflight validation when an abstract header converter class is specified", + 0 ); } @@ -397,7 +429,8 @@ public void testConnectorHasHeaderConverterWithNoSuitableConstructor() throws In config.get(CONNECTOR_CLASS_CONFIG), config, 1, - "Connector config should fail preflight validation when a header converter class with no suitable constructor is specified" + "Connector config should fail preflight validation when a header converter class with no suitable constructor is specified", + 0 ); } @@ -409,7 +442,35 @@ public void testConnectorHasHeaderConverterThatThrowsExceptionOnInstantiation() config.get(CONNECTOR_CLASS_CONFIG), config, 1, - "Connector config should fail preflight validation when a header converter class that throws an exception on instantiation is specified" + "Connector config should fail preflight validation when a header converter class that throws an exception on instantiation is specified", + 0 + ); + } + + @Test + public void testConnectorHasMisconfiguredHeaderConverter() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(HEADER_CONVERTER_CLASS_CONFIG, TestConverterWithSinglePropertyConfigDef.class.getName()); + config.put(HEADER_CONVERTER_CLASS_CONFIG + "." + TestConverterWithSinglePropertyConfigDef.BOOLEAN_PROPERTY_NAME, "notaboolean"); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 1, + "Connector config should fail preflight validation when a header converter fails custom validation", + 0 + ); + } + + @Test + public void testConnectorHasHeaderConverterWithNoConfigDef() throws InterruptedException { + Map config = defaultSinkConnectorProps(); + config.put(HEADER_CONVERTER_CLASS_CONFIG, TestConverterWithNoConfigDef.class.getName()); + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation( + config.get(CONNECTOR_CLASS_CONFIG), + config, + 0, + "Connector config should not fail preflight validation even when a header converter provides a null ConfigDef", + 0 ); } @@ -470,6 +531,21 @@ public TestConverterWithConstructorThatThrowsException() { } } + public static class TestConverterWithSinglePropertyConfigDef extends TestConverter { + public static final String BOOLEAN_PROPERTY_NAME = "prop"; + @Override + public ConfigDef config() { + return new ConfigDef().define(BOOLEAN_PROPERTY_NAME, ConfigDef.Type.BOOLEAN, ConfigDef.Importance.HIGH, ""); + } + } + + public static class TestConverterWithNoConfigDef extends TestConverter { + @Override + public ConfigDef config() { + return null; + } + } + private Map defaultSourceConnectorProps() { // setup up props for the source connector Map props = new HashMap<>(); diff --git a/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.Converter b/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.Converter index 6d38aebee3d92..c58e40f243f3c 100644 --- a/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.Converter +++ b/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.Converter @@ -17,4 +17,6 @@ org.apache.kafka.connect.runtime.SampleConverterWithHeaders org.apache.kafka.connect.runtime.ErrorHandlingTaskTest$FaultyConverter org.apache.kafka.connect.runtime.isolation.PluginsTest$TestConverter org.apache.kafka.connect.runtime.isolation.PluginsTest$TestInternalConverter -org.apache.kafka.connect.runtime.isolation.PluginUtilsTest$CollidingConverter \ No newline at end of file +org.apache.kafka.connect.runtime.isolation.PluginUtilsTest$CollidingConverter +org.apache.kafka.connect.integration.ConnectorValidationIntegrationTest$TestConverterWithSinglePropertyConfigDef +org.apache.kafka.connect.integration.ConnectorValidationIntegrationTest$TestConverterWithNoConfigDef \ No newline at end of file diff --git a/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.HeaderConverter b/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.HeaderConverter index a5b008543b179..b14690acafc88 100644 --- a/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.HeaderConverter +++ b/connect/runtime/src/test/resources/META-INF/services/org.apache.kafka.connect.storage.HeaderConverter @@ -17,4 +17,6 @@ org.apache.kafka.connect.runtime.SampleHeaderConverter org.apache.kafka.connect.runtime.ErrorHandlingTaskTest$FaultyConverter org.apache.kafka.connect.runtime.isolation.PluginsTest$TestHeaderConverter org.apache.kafka.connect.runtime.isolation.PluginsTest$TestInternalConverter -org.apache.kafka.connect.runtime.isolation.PluginUtilsTest$CollidingHeaderConverter \ No newline at end of file +org.apache.kafka.connect.runtime.isolation.PluginUtilsTest$CollidingHeaderConverter +org.apache.kafka.connect.integration.ConnectorValidationIntegrationTest$TestConverterWithSinglePropertyConfigDef +org.apache.kafka.connect.integration.ConnectorValidationIntegrationTest$TestConverterWithNoConfigDef \ No newline at end of file From 5b33be21c5cb0cb3a8909dba77a281947fecd17b Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 5 Sep 2023 16:27:31 -0400 Subject: [PATCH 11/15] Change 'transformation' to 'pluginInstance' in EnrichablePlugin::getConfigDefFromConfigProvidingClass --- .../org/apache/kafka/connect/runtime/ConnectorConfig.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java index f33f00e40efa3..6d7f698258d34 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java @@ -510,13 +510,13 @@ ConfigDef getConfigDefFromConfigProvidingClass(String key, Class cls) { } Utils.ensureConcreteSubclass(baseClass, cls); - T transformation; + T pluginInstance; try { - transformation = Utils.newInstance(cls, baseClass); + pluginInstance = Utils.newInstance(cls, baseClass); } catch (Exception e) { throw new ConfigException(key, String.valueOf(cls), "Error getting config definition from " + baseClass.getSimpleName() + ": " + e.getMessage()); } - ConfigDef configDef = config(transformation); + ConfigDef configDef = config(pluginInstance); if (null == configDef) { throw new ConnectException( String.format( From 6144031a08038a3225237fbdb91bc854ecee4c24 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 31 Oct 2023 15:31:03 -0400 Subject: [PATCH 12/15] Switch to try/finally instead of try-with-resources for Converter/HeaderConverter cleanup --- .../java/org/apache/kafka/connect/runtime/AbstractHerder.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java index 147d49fd81777..072c8ff97629a 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java @@ -423,7 +423,7 @@ private ConfigInfos validateConverterConfig( return null; } - try (Utils.UncheckedCloseable close = () -> Utils.maybeCloseQuietly(converterInstance, converterName + " " + converterClass);) { + try { ConfigDef configDef; try { configDef = configDefAccessor.apply(converterInstance); @@ -459,6 +459,8 @@ private ConfigInfos validateConverterConfig( } return prefixedConfigInfos(configDef.configKeys(), configValues, converterPrefix); + } finally { + Utils.maybeCloseQuietly(converterInstance, converterName + " " + converterClass); } } From 3c8157d8a14411f867e39529e39f7c8ab3c49b2f Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 31 Oct 2023 15:58:09 -0400 Subject: [PATCH 13/15] Make AbstractHerder::validateConverterConfig more generic and future-proof --- .../kafka/connect/runtime/AbstractHerder.java | 95 +++++++++++++------ 1 file changed, 65 insertions(+), 30 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java index 072c8ff97629a..175797ef5ee58 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java @@ -79,6 +79,7 @@ import java.util.List; import java.util.Locale; import java.util.Map; +import java.util.Objects; import java.util.Optional; import java.util.Set; import java.util.concurrent.ConcurrentHashMap; @@ -394,73 +395,107 @@ protected Map validateSourceConnectorConfig(SourceConnector return configDef.validateAll(config); } + /** + * General-purpose validation logic for converters that are configured directly + * in a connector config (as opposed to inherited from the worker config). + * @param connectorConfig the configuration for the connector; may not be null + * @param pluginConfigValue the {@link ConfigValue} for the converter property in the connector config; + * may be null, in which case no validation will be performed under the assumption that the + * connector will use inherit the converter settings from the worker + * @param pluginInterface the interface for the plugin type + * (e.g., {@code org.apache.kafka.connect.storage.Converter.class}); + * may not be null + * @param configDefAccessor an accessor that can be used to retrieve a {@link ConfigDef} + * from an instance of the plugin type (e.g., {@code Converter::config}); + * may not be null + * @param pluginName a lowercase, human-readable name for the type of plugin (e.g., {@code "key converter"}); + * may not be null + * @param pluginProperty the property used to define a custom class for the plugin type + * in a connector config (e.g., {@link ConnectorConfig#KEY_CONVERTER_CLASS_CONFIG}); + * may not be null + * @param defaultProperties any default properties to include in the configuration that will be used for + * the plugin; may be null + + * @return a {@link ConfigInfos} object containing validation results for the plugin in the connector config, + * or null if no custom validation was performed (possibly because no custom plugin was defined in the connector + * config) + + * @param the plugin class to perform validation for + */ private ConfigInfos validateConverterConfig( Map connectorConfig, - ConfigValue converterConfigValue, - Class converterInterface, + ConfigValue pluginConfigValue, + Class pluginInterface, Function configDefAccessor, - String converterName, - String converterProperty, - ConverterType converterType + String pluginName, + String pluginProperty, + Map defaultProperties ) { - String converterClass = connectorConfig.get(converterProperty); + Objects.requireNonNull(connectorConfig); + Objects.requireNonNull(pluginInterface); + Objects.requireNonNull(configDefAccessor); + Objects.requireNonNull(pluginName); + Objects.requireNonNull(pluginProperty); + + String pluginClass = connectorConfig.get(pluginProperty); - if (converterClass == null - || converterConfigValue == null - || !converterConfigValue.errorMessages().isEmpty() + if (pluginClass == null + || pluginConfigValue == null + || !pluginConfigValue.errorMessages().isEmpty() ) { // Either no custom converter was specified, or one was specified but there's a problem with it. // No need to proceed any further. return null; } - T converterInstance; + T pluginInstance; try { - converterInstance = Utils.newInstance(converterClass, converterInterface); + pluginInstance = Utils.newInstance(pluginClass, pluginInterface); } catch (ClassNotFoundException | RuntimeException e) { - log.error("Failed to instantiate {} class {}; this should have been caught by prior validation logic", converterName, converterClass, e); - converterConfigValue.addErrorMessage("Failed to load class " + converterClass + (e.getMessage() != null ? ": " + e.getMessage() : "")); + log.error("Failed to instantiate {} class {}; this should have been caught by prior validation logic", pluginName, pluginClass, e); + pluginConfigValue.addErrorMessage("Failed to load class " + pluginClass + (e.getMessage() != null ? ": " + e.getMessage() : "")); return null; } try { ConfigDef configDef; try { - configDef = configDefAccessor.apply(converterInstance); + configDef = configDefAccessor.apply(pluginInstance); } catch (RuntimeException e) { - log.error("Failed to load ConfigDef from {} of type {}", converterName, converterClass, e); - converterConfigValue.addErrorMessage("Failed to load ConfigDef from " + converterName + (e.getMessage() != null ? ": " + e.getMessage() : "")); + log.error("Failed to load ConfigDef from {} of type {}", pluginName, pluginClass, e); + pluginConfigValue.addErrorMessage("Failed to load ConfigDef from " + pluginName + (e.getMessage() != null ? ": " + e.getMessage() : "")); return null; } if (configDef == null) { - log.warn("{}.config() has returned a null ConfigDef; no further preflight config validation for this converter will be performed", converterClass); + log.warn("{}.config() has returned a null ConfigDef; no further preflight config validation for this converter will be performed", pluginClass); // Older versions of Connect didn't do any converter validation. // Even though converters are technically required to return a non-null ConfigDef object from their config() method, // we permit this case in order to avoid breaking existing converters that, despite not adhering to this requirement, // can be used successfully with a connector. return null; } - final String converterPrefix = converterProperty + "."; - Map converterConfig = connectorConfig.entrySet().stream() - .filter(e -> e.getKey().startsWith(converterPrefix)) + final String pluginPrefix = pluginProperty + "."; + Map pluginConfig = connectorConfig.entrySet().stream() + .filter(e -> e.getKey().startsWith(pluginPrefix)) .collect(Collectors.toMap( - e -> e.getKey().substring(converterPrefix.length()), + e -> e.getKey().substring(pluginPrefix.length()), Map.Entry::getValue )); - converterConfig.putIfAbsent(ConverterConfig.TYPE_CONFIG, converterType.getName()); + if (defaultProperties != null) + defaultProperties.forEach(pluginConfig::putIfAbsent); List configValues; try { - configValues = configDef.validate(converterConfig); + configValues = configDef.validate(pluginConfig); } catch (RuntimeException e) { - log.error("Failed to perform custom config validation for {} of type {}", converterName, converterClass, e); - converterConfigValue.addErrorMessage("Failed to perform custom config validation for " + converterName + (e.getMessage() != null ? ": " + e.getMessage() : "")); + log.error("Failed to perform custom config validation for {} of type {}", pluginName, pluginClass, e); + pluginConfigValue.addErrorMessage("Failed to perform custom config validation for " + pluginName + (e.getMessage() != null ? ": " + e.getMessage() : "")); return null; } - return prefixedConfigInfos(configDef.configKeys(), configValues, converterPrefix); + return prefixedConfigInfos(configDef.configKeys(), configValues, pluginPrefix); } finally { - Utils.maybeCloseQuietly(converterInstance, converterName + " " + converterClass); + Utils.maybeCloseQuietly(pluginInstance, pluginName + " " + pluginClass); } } @@ -472,7 +507,7 @@ private ConfigInfos validateHeaderConverterConfig(Map connectorC HeaderConverter::config, "header converter", HEADER_CONVERTER_CLASS_CONFIG, - ConverterType.HEADER + Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.HEADER.getName()) ); } @@ -484,7 +519,7 @@ private ConfigInfos validateKeyConverterConfig(Map connectorConf Converter::config, "key converter", KEY_CONVERTER_CLASS_CONFIG, - ConverterType.KEY + Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.KEY.getName()) ); } @@ -496,7 +531,7 @@ private ConfigInfos validateValueConverterConfig(Map connectorCo Converter::config, "value converter", VALUE_CONVERTER_CLASS_CONFIG, - ConverterType.VALUE + Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.VALUE.getName()) ); } From 4c55f4bc34d9f2fede3db953b01b8dd8218c8e1b Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Wed, 1 May 2024 14:10:33 -0400 Subject: [PATCH 14/15] Address review comments --- .../org/apache/kafka/common/utils/Utils.java | 9 --- .../kafka/connect/runtime/AbstractHerder.java | 75 +++++++++++++------ 2 files changed, 51 insertions(+), 33 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index d9f4679f76fd6..67fbd3368dcff 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -42,10 +42,7 @@ import java.io.StringWriter; import java.lang.reflect.Constructor; import java.lang.reflect.InvocationTargetException; -import java.lang.reflect.Modifier; -import java.nio.BufferUnderflowException; import java.nio.ByteBuffer; -import java.nio.ByteOrder; import java.nio.channels.FileChannel; import java.nio.charset.StandardCharsets; import java.nio.file.FileVisitResult; @@ -55,7 +52,6 @@ import java.nio.file.Paths; import java.nio.file.SimpleFileVisitor; import java.nio.file.StandardCopyOption; -import java.nio.file.StandardOpenOption; import java.nio.file.attribute.BasicFileAttributes; import java.text.DecimalFormat; import java.text.DecimalFormatSymbols; @@ -64,13 +60,11 @@ import java.time.Instant; import java.time.ZoneId; import java.time.format.DateTimeFormatter; -import java.util.AbstractMap; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.Date; -import java.util.EnumSet; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; @@ -78,12 +72,9 @@ import java.util.List; import java.util.Locale; import java.util.Map; -import java.util.Map.Entry; import java.util.Objects; import java.util.Properties; import java.util.Set; -import java.util.SortedSet; -import java.util.TreeSet; import java.util.concurrent.Callable; import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java index dfa3aaf94bb62..174b032ea2488 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java @@ -405,7 +405,8 @@ protected Map validateSourceConnectorConfig(SourceConnector * @param connectorConfig the configuration for the connector; may not be null * @param pluginConfigValue the {@link ConfigValue} for the converter property in the connector config; * may be null, in which case no validation will be performed under the assumption that the - * connector will use inherit the converter settings from the worker + * connector will use inherit the converter settings from the worker. Some errors encountered + * during validation may be {@link ConfigValue#addErrorMessage(String) added} to this object * @param pluginInterface the interface for the plugin type * (e.g., {@code org.apache.kafka.connect.storage.Converter.class}); * may not be null @@ -421,8 +422,8 @@ protected Map validateSourceConnectorConfig(SourceConnector * the plugin; may be null * @return a {@link ConfigInfos} object containing validation results for the plugin in the connector config, - * or null if no custom validation was performed (possibly because no custom plugin was defined in the connector - * config) + * or null if either no custom validation was performed (possibly because no custom plugin was defined in the + * connector config), or if custom validation failed * @param the plugin class to perform validation for */ @@ -433,7 +434,8 @@ private ConfigInfos validateConverterConfig( Function configDefAccessor, String pluginName, String pluginProperty, - Map defaultProperties + Map defaultProperties, + Function reportStage ) { Objects.requireNonNull(connectorConfig); Objects.requireNonNull(pluginInterface); @@ -453,7 +455,8 @@ private ConfigInfos validateConverterConfig( } T pluginInstance; - try { + String stageDescription = "instantiating the connector's " + pluginName + " for validation"; + try (TemporaryStage stage = reportStage.apply(stageDescription)) { pluginInstance = Utils.newInstance(pluginClass, pluginInterface); } catch (ClassNotFoundException | RuntimeException e) { log.error("Failed to instantiate {} class {}; this should have been caught by prior validation logic", pluginName, pluginClass, e); @@ -463,7 +466,8 @@ private ConfigInfos validateConverterConfig( try { ConfigDef configDef; - try { + stageDescription = "retrieving the configuration definition from the connector's " + pluginName; + try (TemporaryStage stage = reportStage.apply(stageDescription)) { configDef = configDefAccessor.apply(pluginInstance); } catch (RuntimeException e) { log.error("Failed to load ConfigDef from {} of type {}", pluginName, pluginClass, e); @@ -479,21 +483,17 @@ private ConfigInfos validateConverterConfig( return null; } final String pluginPrefix = pluginProperty + "."; - Map pluginConfig = connectorConfig.entrySet().stream() - .filter(e -> e.getKey().startsWith(pluginPrefix)) - .collect(Collectors.toMap( - e -> e.getKey().substring(pluginPrefix.length()), - Map.Entry::getValue - )); + Map pluginConfig = Utils.entriesWithPrefix(connectorConfig, pluginPrefix); if (defaultProperties != null) defaultProperties.forEach(pluginConfig::putIfAbsent); List configValues; - try { + stageDescription = "performing config validation for the connector's " + pluginName; + try (TemporaryStage stage = reportStage.apply(stageDescription)) { configValues = configDef.validate(pluginConfig); } catch (RuntimeException e) { - log.error("Failed to perform custom config validation for {} of type {}", pluginName, pluginClass, e); - pluginConfigValue.addErrorMessage("Failed to perform custom config validation for " + pluginName + (e.getMessage() != null ? ": " + e.getMessage() : "")); + log.error("Failed to perform config validation for {} of type {}", pluginName, pluginClass, e); + pluginConfigValue.addErrorMessage("Failed to perform config validation for " + pluginName + (e.getMessage() != null ? ": " + e.getMessage() : "")); return null; } @@ -503,7 +503,11 @@ private ConfigInfos validateConverterConfig( } } - private ConfigInfos validateHeaderConverterConfig(Map connectorConfig, ConfigValue headerConverterConfigValue) { + private ConfigInfos validateHeaderConverterConfig( + Map connectorConfig, + ConfigValue headerConverterConfigValue, + Function reportStage + ) { return validateConverterConfig( connectorConfig, headerConverterConfigValue, @@ -511,11 +515,16 @@ private ConfigInfos validateHeaderConverterConfig(Map connectorC HeaderConverter::config, "header converter", HEADER_CONVERTER_CLASS_CONFIG, - Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.HEADER.getName()) + Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.HEADER.getName()), + reportStage ); } - private ConfigInfos validateKeyConverterConfig(Map connectorConfig, ConfigValue keyConverterConfigValue) { + private ConfigInfos validateKeyConverterConfig( + Map connectorConfig, + ConfigValue keyConverterConfigValue, + Function reportStage + ) { return validateConverterConfig( connectorConfig, keyConverterConfigValue, @@ -523,11 +532,16 @@ private ConfigInfos validateKeyConverterConfig(Map connectorConf Converter::config, "key converter", KEY_CONVERTER_CLASS_CONFIG, - Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.KEY.getName()) + Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.KEY.getName()), + reportStage ); } - private ConfigInfos validateValueConverterConfig(Map connectorConfig, ConfigValue valueConverterConfigValue) { + private ConfigInfos validateValueConverterConfig( + Map connectorConfig, + ConfigValue valueConverterConfigValue, + Function reportStage + ) { return validateConverterConfig( connectorConfig, valueConverterConfigValue, @@ -535,7 +549,8 @@ private ConfigInfos validateValueConverterConfig(Map connectorCo Converter::config, "value converter", VALUE_CONVERTER_CLASS_CONFIG, - Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.VALUE.getName()) + Collections.singletonMap(ConverterConfig.TYPE_CONFIG, ConverterType.VALUE.getName()), + reportStage ); } @@ -711,9 +726,21 @@ ConfigInfos validateConnectorConfig( configValues.addAll(config.configValues()); // do custom converter-specific validation - ConfigInfos headerConverterConfigInfos = validateHeaderConverterConfig(connectorProps, validatedConnectorConfig.get(HEADER_CONVERTER_CLASS_CONFIG)); - ConfigInfos keyConverterConfigInfos = validateKeyConverterConfig(connectorProps, validatedConnectorConfig.get(KEY_CONVERTER_CLASS_CONFIG)); - ConfigInfos valueConverterConfigInfos = validateValueConverterConfig(connectorProps, validatedConnectorConfig.get(VALUE_CONVERTER_CLASS_CONFIG)); + ConfigInfos headerConverterConfigInfos = validateHeaderConverterConfig( + connectorProps, + validatedConnectorConfig.get(HEADER_CONVERTER_CLASS_CONFIG), + reportStage + ); + ConfigInfos keyConverterConfigInfos = validateKeyConverterConfig( + connectorProps, + validatedConnectorConfig.get(KEY_CONVERTER_CLASS_CONFIG), + reportStage + ); + ConfigInfos valueConverterConfigInfos = validateValueConverterConfig( + connectorProps, + validatedConnectorConfig.get(VALUE_CONVERTER_CLASS_CONFIG), + reportStage + ); ConfigInfos configInfos = generateResult(connType, configKeys, configValues, new ArrayList<>(allGroups)); AbstractConfig connectorConfig = new AbstractConfig(new ConfigDef(), connectorProps, doLog); From 637610ae031b13cb700bd3387934dce8872962ee Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Tue, 7 May 2024 11:24:10 -0400 Subject: [PATCH 15/15] Revert unnecessary changes to Utils.java --- .../org/apache/kafka/common/utils/Utils.java | 18 +++++++++--------- 1 file changed, 9 insertions(+), 9 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java index 67fbd3368dcff..f62aee76cbc05 100644 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -16,15 +16,6 @@ */ package org.apache.kafka.common.utils; -import java.lang.reflect.Modifier; -import java.nio.BufferUnderflowException; -import java.nio.ByteOrder; -import java.nio.file.StandardOpenOption; -import java.util.AbstractMap; -import java.util.EnumSet; -import java.util.Map.Entry; -import java.util.SortedSet; -import java.util.TreeSet; import org.apache.kafka.common.KafkaException; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.network.TransferableChannel; @@ -42,7 +33,10 @@ import java.io.StringWriter; import java.lang.reflect.Constructor; import java.lang.reflect.InvocationTargetException; +import java.lang.reflect.Modifier; +import java.nio.BufferUnderflowException; import java.nio.ByteBuffer; +import java.nio.ByteOrder; import java.nio.channels.FileChannel; import java.nio.charset.StandardCharsets; import java.nio.file.FileVisitResult; @@ -52,6 +46,7 @@ import java.nio.file.Paths; import java.nio.file.SimpleFileVisitor; import java.nio.file.StandardCopyOption; +import java.nio.file.StandardOpenOption; import java.nio.file.attribute.BasicFileAttributes; import java.text.DecimalFormat; import java.text.DecimalFormatSymbols; @@ -60,11 +55,13 @@ import java.time.Instant; import java.time.ZoneId; import java.time.format.DateTimeFormatter; +import java.util.AbstractMap; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; import java.util.Collections; import java.util.Date; +import java.util.EnumSet; import java.util.HashMap; import java.util.HashSet; import java.util.Iterator; @@ -72,9 +69,12 @@ import java.util.List; import java.util.Locale; import java.util.Map; +import java.util.Map.Entry; import java.util.Objects; import java.util.Properties; import java.util.Set; +import java.util.SortedSet; +import java.util.TreeSet; import java.util.concurrent.Callable; import java.util.concurrent.atomic.AtomicReference; import java.util.function.BiConsumer;