diff --git a/checkstyle/suppressions.xml b/checkstyle/suppressions.xml
index c3a64742927d5..a561c1339a0a9 100644
--- a/checkstyle/suppressions.xml
+++ b/checkstyle/suppressions.xml
@@ -133,7 +133,7 @@
files="(KafkaConfigBackingStore|Values).java"/>
+ files="(AbstractHerder|DistributedHerder|RestClient|RestServer|JsonConverter|KafkaConfigBackingStore|FileStreamSourceTask|TopicAdmin).java"/>
cls = (Class>) value;
+ try {
+ cls.getDeclaredConstructor().newInstance();
+ } 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";
+ }
+ }
+
+ public static class ConcreteSubClassValidator implements 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();
+ }
+ }
+
public static class ConfigKey {
public final String name;
public final Type type;
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 30555ef4ad261..50ee7ff13c2be 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
@@ -24,6 +24,7 @@
import org.apache.kafka.common.config.ConfigDef.Type;
import org.apache.kafka.common.config.ConfigTransformer;
import org.apache.kafka.common.config.ConfigValue;
+import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.connect.connector.Connector;
import org.apache.kafka.connect.connector.policy.ConnectorClientConfigOverridePolicy;
import org.apache.kafka.connect.connector.policy.ConnectorClientConfigRequest;
@@ -41,9 +42,14 @@
import org.apache.kafka.connect.runtime.rest.errors.BadRequestException;
import org.apache.kafka.connect.source.SourceConnector;
import org.apache.kafka.connect.storage.ConfigBackingStore;
+import org.apache.kafka.connect.storage.ConverterConfig;
+import org.apache.kafka.connect.storage.ConverterType;
+import org.apache.kafka.connect.storage.HeaderConverter;
import org.apache.kafka.connect.storage.StatusBackingStore;
import org.apache.kafka.connect.util.Callback;
import org.apache.kafka.connect.util.ConnectorTaskId;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
import java.io.ByteArrayOutputStream;
import java.io.PrintStream;
@@ -69,6 +75,8 @@
import java.util.regex.Pattern;
import java.util.stream.Collectors;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.HEADER_CONVERTER_CLASS_CONFIG;
+
/**
* Abstract Herder implementation which handles connector/task lifecycle tracking. Extensions
* must invoke the lifecycle hooks appropriately.
@@ -92,6 +100,8 @@
*/
public abstract class AbstractHerder implements Herder, TaskStatus.Listener, ConnectorStatus.Listener {
+ private static final Logger log = LoggerFactory.getLogger(AbstractHerder.class);
+
private final String workerId;
protected final Worker worker;
private final String kafkaClusterId;
@@ -344,12 +354,73 @@ public ConnectorStateInfo.TaskState taskStatus(ConnectorTaskId id) {
status.workerId(), status.trace());
}
- protected Map validateBasicConnectorConfig(Connector connector,
- ConfigDef configDef,
- Map config) {
+ protected Map validateSinkConnectorConfig(ConfigDef configDef, Map config) {
+ return SinkConnectorConfig.validate(configDef.validateAll(config), config);
+ }
+
+ protected Map validateSourceConnectorConfig(ConfigDef configDef, Map config) {
return configDef.validateAll(config);
}
+ private ConfigInfos validateHeaderConverterConfig(Map connectorConfig, ConfigValue headerConverterConfigValue) {
+ String headerConverterClass = connectorConfig.get(HEADER_CONVERTER_CLASS_CONFIG);
+
+ if (headerConverterClass == null
+ || headerConverterConfigValue == null
+ || !headerConverterConfigValue.errorMessages().isEmpty()
+ ) {
+ // Either no custom header converter was specified, or one was specified but there's a problem with it.
+ // No need to proceed any further.
+ return null;
+ }
+
+ HeaderConverter headerConverter;
+ try {
+ headerConverter = Utils.newInstance(headerConverterClass, HeaderConverter.class);
+ } catch (ClassNotFoundException | RuntimeException e) {
+ log.error("Failed to instantiate header converter class {}; this should have been caught by prior validation logic", headerConverterClass, e);
+ headerConverterConfigValue.addErrorMessage("Failed to load class " + headerConverterClass + (e.getMessage() != null ? ": " + e.getMessage() : ""));
+ return null;
+ }
+
+ ConfigDef configDef;
+ try {
+ configDef = headerConverter.config();
+ } catch (RuntimeException e) {
+ log.error("Failed to load ConfigDef from header converter of type {}", headerConverterClass, e);
+ headerConverterConfigValue.addErrorMessage("Failed to load ConfigDef from header converter" + (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", headerConverterClass);
+ // Older versions of Connect didn't do any header converter validation.
+ // Even though header converters are technically required to return a non-null ConfigDef object from HeaderConverter::config,
+ // we permit this case in order to avoid breaking existing header converters that, despite not adhering to this requirement,
+ // can be used successfully with a connector.
+ return null;
+ }
+
+ final String headerConverterPrefix = HEADER_CONVERTER_CLASS_CONFIG + ".";
+ Map headerConverterConfig = connectorConfig.entrySet().stream()
+ .filter(e -> e.getKey().startsWith(headerConverterPrefix))
+ .collect(Collectors.toMap(
+ e -> e.getKey().substring(headerConverterPrefix.length()),
+ Map.Entry::getValue
+ ));
+ headerConverterConfig.put(ConverterConfig.TYPE_CONFIG, ConverterType.HEADER.getName());
+
+ List configValues;
+ try {
+ configValues = configDef.validate(headerConverterConfig);
+ } catch (RuntimeException e) {
+ log.error("Failed to perform custom config validation for header converter of type {}", headerConverterClass, e);
+ headerConverterConfigValue.addErrorMessage("Failed to perform custom config validation for header converter" + (e.getMessage() != null ? ": " + e.getMessage() : ""));
+ return null;
+ }
+
+ return prefixedConfigInfos(configDef.configKeys(), configValues, headerConverterPrefix);
+ }
+
@Override
public void validateConnectorConfig(Map connectorProps, Callback callback) {
validateConnectorConfig(connectorProps, callback, true);
@@ -426,22 +497,18 @@ ConfigInfos validateConnectorConfig(Map connectorProps, boolean
Connector connector = getConnector(connType);
org.apache.kafka.connect.health.ConnectorType connectorType;
ClassLoader savedLoader = plugins().compareAndSwapLoaders(connector);
+ ConfigDef enrichedConfigDef;
+ Map validatedConnectorConfig;
try {
- ConfigDef baseConfigDef;
if (connector instanceof SourceConnector) {
- baseConfigDef = SourceConnectorConfig.configDef();
connectorType = org.apache.kafka.connect.health.ConnectorType.SOURCE;
+ enrichedConfigDef = ConnectorConfig.enrich(plugins(), SourceConnectorConfig.configDef(), connectorProps, false);
+ validatedConnectorConfig = validateSourceConnectorConfig(enrichedConfigDef, connectorProps);
} else {
- baseConfigDef = SinkConnectorConfig.configDef();
- SinkConnectorConfig.validate(connectorProps);
connectorType = org.apache.kafka.connect.health.ConnectorType.SINK;
+ enrichedConfigDef = ConnectorConfig.enrich(plugins(), SinkConnectorConfig.configDef(), connectorProps, false);
+ validatedConnectorConfig = validateSinkConnectorConfig(enrichedConfigDef, connectorProps);
}
- ConfigDef enrichedConfigDef = ConnectorConfig.enrich(plugins(), baseConfigDef, connectorProps, false);
- Map validatedConnectorConfig = validateBasicConnectorConfig(
- connector,
- enrichedConfigDef,
- connectorProps
- );
List configValues = new ArrayList<>(validatedConnectorConfig.values());
Map configKeys = new LinkedHashMap<>(enrichedConfigDef.configKeys());
Set allGroups = new LinkedHashSet<>(enrichedConfigDef.groups());
@@ -468,7 +535,11 @@ 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 header converter-specific validation
+ ConfigInfos headerConverterConfigInfos = validateHeaderConverterConfig(connectorProps, validatedConnectorConfig.get(HEADER_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);
@@ -484,7 +555,7 @@ ConfigInfos validateConnectorConfig(Map connectorProps, boolean
connectorType,
ConnectorClientConfigRequest.ClientType.PRODUCER,
connectorClientConfigOverridePolicy);
- return mergeConfigInfos(connType, configInfos, producerConfigInfos);
+ return mergeConfigInfos(connType, configInfos, producerConfigInfos, headerConverterConfigInfos);
} else {
consumerConfigInfos = validateClientOverrides(connName,
ConnectorConfig.CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX,
@@ -508,7 +579,7 @@ ConfigInfos validateConnectorConfig(Map connectorProps, boolean
}
}
- return mergeConfigInfos(connType, configInfos, producerConfigInfos, consumerConfigInfos, adminConfigInfos);
+ return mergeConfigInfos(connType, configInfos, producerConfigInfos, consumerConfigInfos, adminConfigInfos, headerConverterConfigInfos);
} finally {
Plugins.compareAndSwapLoaders(savedLoader);
}
@@ -536,10 +607,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();
@@ -550,30 +617,42 @@ private static ConfigInfos validateClientOverrides(String connName,
: rawConfigValue;
clientConfigs.put(configName, parsedConfigValue);
}
+
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("", errorCount, 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/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectorConfig.java
index 4ba1ddd6dad45..4c3adcf2ffce8 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,6 +28,8 @@
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.storage.HeaderConverter;
import org.apache.kafka.connect.transforms.Transformation;
import org.apache.kafka.connect.transforms.predicates.Predicate;
import org.slf4j.Logger;
@@ -82,10 +84,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(
+ ConfigDef.ConcreteSubClassValidator.forSuperClass(Converter.class),
+ new ConfigDef.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(
+ ConfigDef.ConcreteSubClassValidator.forSuperClass(Converter.class),
+ new ConfigDef.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;
@@ -93,6 +103,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(
+ ConfigDef.ConcreteSubClassValidator.forSuperClass(HeaderConverter.class),
+ new ConfigDef.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.";
@@ -179,9 +193,9 @@ 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(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(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, 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,
@@ -425,7 +439,10 @@ void enrich(ConfigDef newDef) {
final ConfigDef.Validator typeValidator = ConfigDef.LambdaValidator.with(
(String name, Object value) -> {
validateProps(prefix);
- getConfigDefFromConfigProvidingClass(typeConfig, (Class>) value);
+ // The value will be null if the class couldn't be found; no point in trying to load a ConfigDef for it
+ if (value != null) {
+ getConfigDefFromConfigProvidingClass(typeConfig, (Class>) value);
+ }
},
() -> "valid configs for " + alias + " " + aliasKind.toLowerCase(Locale.ENGLISH));
newDef.define(typeConfig, Type.CLASS, ConfigDef.NO_DEFAULT_VALUE, typeValidator, Importance.HIGH,
@@ -499,13 +516,15 @@ ConfigDef getConfigDefFromConfigProvidingClass(String key, Class> cls) {
aliasKind + " is abstract and cannot be created. Did you mean " + childClassNames + "?";
throw new ConfigException(key, String.valueOf(cls), message);
}
- T transformation;
+ T plugin;
try {
- transformation = Utils.newInstance(cls, baseClass);
+ plugin = Utils.newInstance(cls, baseClass);
} catch (Exception e) {
- throw new ConfigException(key, String.valueOf(cls), "Error getting config definition from " + baseClass.getSimpleName() + ": " + e.getMessage());
+ // Log the entire exception here in order to provide a stack trace that can be useful for debugging classloading issues
+ log.error("Failed to instantiate {} '{}'", baseClass.getSimpleName(), cls, e);
+ throw new ConfigException(key, String.valueOf(cls), "Error instantiating " + baseClass.getSimpleName() + ": " + e.getMessage());
}
- ConfigDef configDef = config(transformation);
+ ConfigDef configDef = config(plugin);
if (null == configDef) {
throw new ConnectException(
String.format(
diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SinkConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SinkConnectorConfig.java
index 93c2cb458ab98..04947687ece11 100644
--- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SinkConnectorConfig.java
+++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SinkConnectorConfig.java
@@ -19,15 +19,18 @@
import org.apache.kafka.common.config.ConfigDef;
import org.apache.kafka.common.config.ConfigDef.Importance;
import org.apache.kafka.common.config.ConfigDef.Type;
-import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.common.config.ConfigException;
+import org.apache.kafka.common.config.ConfigValue;
+import org.apache.kafka.common.utils.Utils;
import org.apache.kafka.connect.runtime.isolation.Plugins;
import org.apache.kafka.connect.sink.SinkTask;
import org.apache.kafka.connect.transforms.util.RegexValidator;
+import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.Map;
+import java.util.function.Consumer;
import java.util.regex.Pattern;
import java.util.stream.Collectors;
@@ -70,19 +73,17 @@ public class SinkConnectorConfig extends ConnectorConfig {
"keys, all error context header keys will start with __connect.errors.";
private static final String DLQ_CONTEXT_HEADERS_ENABLE_DISPLAY = "Enable Error Context Headers";
- static ConfigDef config = ConnectorConfig.configDef()
- .define(TOPICS_CONFIG, ConfigDef.Type.LIST, TOPICS_DEFAULT, ConfigDef.Importance.HIGH, TOPICS_DOC, COMMON_GROUP, 4, ConfigDef.Width.LONG, TOPICS_DISPLAY)
- .define(TOPICS_REGEX_CONFIG, ConfigDef.Type.STRING, TOPICS_REGEX_DEFAULT, new RegexValidator(), ConfigDef.Importance.HIGH, TOPICS_REGEX_DOC, COMMON_GROUP, 4, ConfigDef.Width.LONG, TOPICS_REGEX_DISPLAY)
- .define(DLQ_TOPIC_NAME_CONFIG, ConfigDef.Type.STRING, DLQ_TOPIC_DEFAULT, Importance.MEDIUM, DLQ_TOPIC_NAME_DOC, ERROR_GROUP, 6, ConfigDef.Width.MEDIUM, DLQ_TOPIC_DISPLAY)
- .define(DLQ_TOPIC_REPLICATION_FACTOR_CONFIG, ConfigDef.Type.SHORT, DLQ_TOPIC_REPLICATION_FACTOR_CONFIG_DEFAULT, Importance.MEDIUM, DLQ_TOPIC_REPLICATION_FACTOR_CONFIG_DOC, ERROR_GROUP, 7, ConfigDef.Width.MEDIUM, DLQ_TOPIC_REPLICATION_FACTOR_CONFIG_DISPLAY)
- .define(DLQ_CONTEXT_HEADERS_ENABLE_CONFIG, ConfigDef.Type.BOOLEAN, DLQ_CONTEXT_HEADERS_ENABLE_DEFAULT, Importance.MEDIUM, DLQ_CONTEXT_HEADERS_ENABLE_DOC, ERROR_GROUP, 8, ConfigDef.Width.MEDIUM, DLQ_CONTEXT_HEADERS_ENABLE_DISPLAY);
-
public static ConfigDef configDef() {
- return config;
+ return ConnectorConfig.configDef()
+ .define(TOPICS_CONFIG, ConfigDef.Type.LIST, TOPICS_DEFAULT, ConfigDef.Importance.HIGH, TOPICS_DOC, COMMON_GROUP, 4, ConfigDef.Width.LONG, TOPICS_DISPLAY)
+ .define(TOPICS_REGEX_CONFIG, ConfigDef.Type.STRING, TOPICS_REGEX_DEFAULT, new RegexValidator(), ConfigDef.Importance.HIGH, TOPICS_REGEX_DOC, COMMON_GROUP, 4, ConfigDef.Width.LONG, TOPICS_REGEX_DISPLAY)
+ .define(DLQ_TOPIC_NAME_CONFIG, ConfigDef.Type.STRING, DLQ_TOPIC_DEFAULT, Importance.MEDIUM, DLQ_TOPIC_NAME_DOC, ERROR_GROUP, 6, ConfigDef.Width.MEDIUM, DLQ_TOPIC_DISPLAY)
+ .define(DLQ_TOPIC_REPLICATION_FACTOR_CONFIG, ConfigDef.Type.SHORT, DLQ_TOPIC_REPLICATION_FACTOR_CONFIG_DEFAULT, Importance.MEDIUM, DLQ_TOPIC_REPLICATION_FACTOR_CONFIG_DOC, ERROR_GROUP, 7, ConfigDef.Width.MEDIUM, DLQ_TOPIC_REPLICATION_FACTOR_CONFIG_DISPLAY)
+ .define(DLQ_CONTEXT_HEADERS_ENABLE_CONFIG, ConfigDef.Type.BOOLEAN, DLQ_CONTEXT_HEADERS_ENABLE_DEFAULT, Importance.MEDIUM, DLQ_CONTEXT_HEADERS_ENABLE_DOC, ERROR_GROUP, 8, ConfigDef.Width.MEDIUM, DLQ_CONTEXT_HEADERS_ENABLE_DISPLAY);
}
public SinkConnectorConfig(Plugins plugins, Map props) {
- super(plugins, config, props);
+ super(plugins, configDef(), props);
}
/**
@@ -90,40 +91,87 @@ public SinkConnectorConfig(Plugins plugins, Map props) {
* @param props sink configuration properties
*/
public static void validate(Map props) {
- final boolean hasTopicsConfig = hasTopicsConfig(props);
- final boolean hasTopicsRegexConfig = hasTopicsRegexConfig(props);
- final boolean hasDlqTopicConfig = hasDlqTopicConfig(props);
+ validate(
+ props,
+ error -> {
+ throw new ConfigException(error.property, error.value, error.errorMessage);
+ }
+ );
+ }
+
+ /**
+ * Perform preflight validation for the sink-specific properties for this connector.
+ */
+ public static Map validate(Map validatedConfig, Map props) {
+ validate(props, error -> addErrorMessage(validatedConfig, error));
+ return validatedConfig;
+ }
+
+ private static void validate(Map props, Consumer onError) {
+ final String topicsList = props.get(TOPICS_CONFIG);
+ final String topicsRegex = props.get(TOPICS_REGEX_CONFIG);
+ final String dlqTopic = props.getOrDefault(DLQ_TOPIC_NAME_CONFIG, "").trim();
+ final boolean hasTopicsConfig = !Utils.isBlank(topicsList);
+ final boolean hasTopicsRegexConfig = !Utils.isBlank(topicsRegex);
+ final boolean hasDlqTopicConfig = !Utils.isBlank(dlqTopic);
if (hasTopicsConfig && hasTopicsRegexConfig) {
- throw new ConfigException(SinkTask.TOPICS_CONFIG + " and " + SinkTask.TOPICS_REGEX_CONFIG +
- " are mutually exclusive options, but both are set.");
+ String errorMessage = TOPICS_CONFIG + " and " + TOPICS_REGEX_CONFIG + " are mutually exclusive options, but both are set.";
+ onError.accept(new ConfigError(TOPICS_CONFIG, topicsList, errorMessage));
+ onError.accept(new ConfigError(TOPICS_REGEX_CONFIG, topicsRegex, errorMessage));
}
if (!hasTopicsConfig && !hasTopicsRegexConfig) {
- throw new ConfigException("Must configure one of " +
- SinkTask.TOPICS_CONFIG + " or " + SinkTask.TOPICS_REGEX_CONFIG);
+ String errorMessage = "Must configure one of " + TOPICS_CONFIG + " or " + TOPICS_REGEX_CONFIG;
+ onError.accept(new ConfigError(TOPICS_CONFIG, topicsList, errorMessage));
+ onError.accept(new ConfigError(TOPICS_REGEX_CONFIG, topicsRegex, errorMessage));
}
if (hasDlqTopicConfig) {
- String dlqTopic = props.get(DLQ_TOPIC_NAME_CONFIG).trim();
if (hasTopicsConfig) {
List topics = parseTopicsList(props);
if (topics.contains(dlqTopic)) {
- throw new ConfigException(String.format("The DLQ topic '%s' may not be included in the list of "
- + "topics ('%s=%s') consumed by the connector", dlqTopic, SinkTask.TOPICS_REGEX_CONFIG, topics));
+ String errorMessage = String.format(
+ "The DLQ topic '%s' may not be included in the list of topics ('%s=%s') consumed by the connector",
+ dlqTopic, TOPICS_CONFIG, topics
+ );
+ onError.accept(new ConfigError(TOPICS_CONFIG, topicsList, errorMessage));
}
}
if (hasTopicsRegexConfig) {
- String topicsRegexStr = props.get(SinkTask.TOPICS_REGEX_CONFIG);
- Pattern pattern = Pattern.compile(topicsRegexStr);
+ Pattern pattern = Pattern.compile(topicsRegex);
if (pattern.matcher(dlqTopic).matches()) {
- throw new ConfigException(String.format("The DLQ topic '%s' may not be included in the regex matching the "
- + "topics ('%s=%s') consumed by the connector", dlqTopic, SinkTask.TOPICS_REGEX_CONFIG, topicsRegexStr));
+ String errorMessage = String.format(
+ "The DLQ topic '%s' may not be included in the regex matching the topics ('%s=%s') consumed by the connector",
+ dlqTopic, TOPICS_REGEX_CONFIG, topicsRegex
+ );
+ onError.accept(new ConfigError(TOPICS_REGEX_CONFIG, topicsRegex, errorMessage));
}
}
}
}
+ private static class ConfigError {
+ public final String property;
+ public final Object value;
+ public final String errorMessage;
+
+ public ConfigError(String property, Object value, String errorMessage) {
+ this.property = property;
+ this.value = value;
+ this.errorMessage = errorMessage;
+ }
+ }
+
+ private static void addErrorMessage(Map validatedConfig, ConfigError error) {
+ validatedConfig.computeIfAbsent(
+ error.property,
+ p -> new ConfigValue(error.property, error.value, Collections.emptyList(), new ArrayList<>())
+ ).addErrorMessage(
+ error.errorMessage
+ );
+ }
+
public static boolean hasTopicsConfig(Map props) {
String topicsStr = props.get(TOPICS_CONFIG);
return !Utils.isBlank(topicsStr);
@@ -134,11 +182,6 @@ public static boolean hasTopicsRegexConfig(Map props) {
return !Utils.isBlank(topicsRegexStr);
}
- public static boolean hasDlqTopicConfig(Map props) {
- String dqlTopicStr = props.get(DLQ_TOPIC_NAME_CONFIG);
- return !Utils.isBlank(dqlTopicStr);
- }
-
@SuppressWarnings("unchecked")
public static List parseTopicsList(Map props) {
List topics = (List) ConfigDef.parseType(TOPICS_CONFIG, props.get(TOPICS_CONFIG), Type.LIST);
@@ -170,6 +213,6 @@ public boolean enableErrantRecordReporter() {
}
public static void main(String[] args) {
- System.out.println(config.toHtml(4, config -> "sinkconnectorconfigs_" + config));
+ System.out.println(configDef().toHtml(4, config -> "sinkconnectorconfigs_" + config));
}
}
diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceConnectorConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceConnectorConfig.java
index 7cf5d67899e61..d7f32f41f0367 100644
--- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceConnectorConfig.java
+++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceConnectorConfig.java
@@ -58,7 +58,6 @@ public Object get(String key) {
}
}
- private static ConfigDef config = SourceConnectorConfig.configDef();
private final EnrichedSourceConnectorConfig enrichedSourceConfig;
public static ConfigDef configDef() {
@@ -116,9 +115,9 @@ public static ConfigDef enrich(ConfigDef baseConfigDef, Map prop
}
public SourceConnectorConfig(Plugins plugins, Map props, boolean createTopics) {
- super(plugins, config, props);
+ super(plugins, configDef(), props);
if (createTopics && props.entrySet().stream().anyMatch(e -> e.getKey().startsWith(TOPIC_CREATION_PREFIX))) {
- ConfigDef defaultConfigDef = embedDefaultGroup(config);
+ ConfigDef defaultConfigDef = embedDefaultGroup(configDef());
// This config is only used to set default values for partitions and replication
// factor from the default group and otherwise it remains unused
AbstractConfig defaultGroup = new AbstractConfig(defaultConfigDef, props, false);
@@ -181,6 +180,6 @@ public Map topicCreationOtherConfigs(String group) {
}
public static void main(String[] args) {
- System.out.println(config.toHtml(4, config -> "sourceconnectorconfigs_" + config));
+ System.out.println(configDef().toHtml(4, config -> "sourceconnectorconfigs_" + config));
}
}
diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java
index 5aa327ed61e5c..a47db4a7fb18c 100644
--- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java
+++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java
@@ -29,7 +29,6 @@
import org.apache.kafka.common.utils.ThreadUtils;
import org.apache.kafka.common.utils.Time;
import org.apache.kafka.common.utils.Utils;
-import org.apache.kafka.connect.connector.Connector;
import org.apache.kafka.connect.connector.policy.ConnectorClientConfigOverridePolicy;
import org.apache.kafka.connect.errors.AlreadyExistsException;
import org.apache.kafka.connect.errors.ConnectException;
@@ -57,7 +56,6 @@
import org.apache.kafka.connect.runtime.rest.entities.TaskInfo;
import org.apache.kafka.connect.runtime.rest.errors.BadRequestException;
import org.apache.kafka.connect.runtime.rest.errors.ConnectRestException;
-import org.apache.kafka.connect.sink.SinkConnector;
import org.apache.kafka.connect.storage.ConfigBackingStore;
import org.apache.kafka.connect.storage.StatusBackingStore;
import org.apache.kafka.connect.util.Callback;
@@ -96,6 +94,8 @@
import java.util.concurrent.atomic.AtomicReference;
import java.util.stream.Collectors;
+import static org.apache.kafka.clients.CommonClientConfigs.GROUP_ID_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX;
import static org.apache.kafka.connect.runtime.WorkerConfig.TOPIC_TRACKING_ENABLE_CONFIG;
import static org.apache.kafka.connect.runtime.distributed.ConnectProtocol.CONNECT_PROTOCOL_V0;
import static org.apache.kafka.connect.runtime.distributed.ConnectProtocolCompatibility.EAGER;
@@ -842,22 +842,30 @@ public void deleteConnectorConfig(final String connName, final Callback validateBasicConnectorConfig(Connector connector,
- ConfigDef configDef,
- Map config) {
- Map validatedConfig = super.validateBasicConnectorConfig(connector, configDef, config);
- if (connector instanceof SinkConnector) {
+ protected Map validateSinkConnectorConfig(ConfigDef configDef, Map config) {
+ Map validatedConfig = super.validateSinkConnectorConfig(configDef, config);
+ final String overriddenConsumerGroupIdConfig = CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX + GROUP_ID_CONFIG;
+ if (config.containsKey(overriddenConsumerGroupIdConfig)) {
+ String consumerGroupId = config.get(overriddenConsumerGroupIdConfig);
+ ConfigValue validatedGroupId = validatedConfig.computeIfAbsent(
+ overriddenConsumerGroupIdConfig,
+ p -> new ConfigValue(overriddenConsumerGroupIdConfig, consumerGroupId, Collections.emptyList(), new ArrayList<>())
+ );
+ if (workerGroupId.equals(consumerGroupId)) {
+ validatedGroupId.addErrorMessage("Consumer group " + consumerGroupId +
+ " conflicts with Connect worker group " + workerGroupId);
+ }
+ } else {
ConfigValue validatedName = validatedConfig.get(ConnectorConfig.NAME_CONFIG);
String name = (String) validatedName.value();
if (workerGroupId.equals(SinkUtils.consumerGroupId(name))) {
validatedName.addErrorMessage("Consumer group for sink connector named " + name +
- " conflicts with Connect worker group " + workerGroupId);
+ " conflicts with Connect worker group " + workerGroupId);
}
}
return validatedConfig;
}
-
@Override
public void putConnectorConfig(final String connName, final Map config, final boolean allowReplace,
final Callback> callback) {
diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorClientPolicyIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorClientPolicyIntegrationTest.java
index a0abece17d750..553465a9b8ba7 100644
--- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorClientPolicyIntegrationTest.java
+++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorClientPolicyIntegrationTest.java
@@ -133,7 +133,7 @@ private void assertFailCreateConnector(String policy, Map props)
connect.configureConnector(CONNECTOR_NAME, props);
fail("Shouldn't be able to create connector");
} catch (ConnectRestException e) {
- assertEquals(e.statusCode(), 400);
+ assertEquals(400, e.statusCode());
} finally {
connect.stop();
}
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
new file mode 100644
index 0000000000000..7d5720000fa9c
--- /dev/null
+++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorValidationIntegrationTest.java
@@ -0,0 +1,514 @@
+/*
+ * 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.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.HeaderConverter;
+import org.apache.kafka.connect.storage.StringConverter;
+import org.apache.kafka.connect.transforms.Filter;
+import org.apache.kafka.connect.transforms.predicates.RecordIsTombstone;
+import org.apache.kafka.connect.util.clusters.EmbeddedConnectCluster;
+import org.apache.kafka.test.IntegrationTest;
+import org.junit.AfterClass;
+import org.junit.BeforeClass;
+import org.junit.Test;
+import org.junit.experimental.categories.Category;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.apache.kafka.clients.consumer.ConsumerConfig.GROUP_ID_CONFIG;
+import static org.apache.kafka.connect.integration.MonitorableSourceConnector.TOPIC_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLASS_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.HEADER_CONVERTER_CLASS_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.KEY_CONVERTER_CLASS_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.NAME_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.PREDICATES_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.TASKS_MAX_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.TRANSFORMS_CONFIG;
+import static org.apache.kafka.connect.runtime.ConnectorConfig.VALUE_CONVERTER_CLASS_CONFIG;
+import static org.apache.kafka.connect.runtime.SinkConnectorConfig.DLQ_TOPIC_NAME_CONFIG;
+import static org.apache.kafka.connect.runtime.SinkConnectorConfig.TOPICS_CONFIG;
+import static org.apache.kafka.connect.runtime.SinkConnectorConfig.TOPICS_REGEX_CONFIG;
+import static org.apache.kafka.connect.runtime.SourceConnectorConfig.TOPIC_CREATION_GROUPS_CONFIG;
+
+/**
+ * Integration test for preflight connector config validation
+ */
+@Category(IntegrationTest.class)
+public class ConnectorValidationIntegrationTest {
+
+ private static final String WORKER_GROUP_ID = "connect-worker-group-id";
+
+ // Use a single embedded cluster for all test cases in order to cut down on runtime
+ private static EmbeddedConnectCluster connect;
+
+ @BeforeClass
+ public static void setup() {
+ Map workerProps = new HashMap<>();
+ workerProps.put(GROUP_ID_CONFIG, WORKER_GROUP_ID);
+
+ // build a Connect cluster backed by Kafka and Zk
+ connect = new EmbeddedConnectCluster.Builder()
+ .name("connector-validation-connect-cluster")
+ .workerProps(workerProps)
+ .build();
+ connect.start();
+ }
+
+ @AfterClass
+ public static void close() {
+ if (connect != null) {
+ // stop all Connect, Kafka and Zk threads.
+ Utils.closeQuietly(connect::stop, "Embedded Connect cluster");
+ }
+ }
+
+ @Test
+ public void testSinkConnectorHasNeitherTopicsListNorTopicsRegex() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ config.remove(TOPICS_CONFIG);
+ config.remove(TOPICS_REGEX_CONFIG);
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 2, // One error each for topics list and topics regex
+ "Sink connector config should fail preflight validation when neither topics list nor topics regex are provided"
+ );
+ }
+
+ @Test
+ public void testSinkConnectorHasBothTopicsListAndTopicsRegex() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ config.put(TOPICS_CONFIG, "t1");
+ config.put(TOPICS_REGEX_CONFIG, "r.*");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 2, // One error each for topics list and topics regex
+ "Sink connector config should fail preflight validation when both topics list and topics regex are provided"
+ );
+ }
+
+ @Test
+ public void testSinkConnectorDeadLetterQueueTopicInTopicsList() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ config.put(TOPICS_CONFIG, "t1");
+ config.put(DLQ_TOPIC_NAME_CONFIG, "t1");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Sink connector config should fail preflight validation when DLQ topic is included in topics list"
+ );
+ }
+
+ @Test
+ public void testSinkConnectorDeadLetterQueueTopicMatchesTopicsRegex() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ config.put(TOPICS_REGEX_CONFIG, "r.*");
+ config.put(DLQ_TOPIC_NAME_CONFIG, "ruh.roh");
+ config.remove(TOPICS_CONFIG);
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Sink connector config should fail preflight validation when DLQ topic matches topics regex"
+ );
+ }
+
+ @Test
+ public void testSinkConnectorDefaultGroupIdConflictsWithWorkerGroupId() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ // Combined with the logic in SinkUtils::consumerGroupId, this should conflict with the worker group ID
+ config.put(NAME_CONFIG, "worker-group-id");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Sink connector config should fail preflight validation when default consumer group ID conflicts with Connect worker group ID"
+ );
+ }
+
+ @Test
+ public void testSinkConnectorOverriddenGroupIdConflictsWithWorkerGroupId() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ config.put(CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX + GROUP_ID_CONFIG, WORKER_GROUP_ID);
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Sink connector config should fail preflight validation when overridden consumer group ID conflicts with Connect worker group ID"
+ );
+ }
+
+ @Test
+ public void testSourceConnectorHasDuplicateTopicCreationGroups() throws InterruptedException {
+ Map config = defaultSourceConnectorProps();
+ config.put(TOPIC_CREATION_GROUPS_CONFIG, "g1, g2, g1");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Source connector config should fail preflight validation when the same topic creation group is specified multiple times"
+ );
+ }
+
+ @Test
+ // TODO: Is this actually necessary? Should we permit the same SMT to be applied multiple times?
+ public void testConnectorHasDuplicateTransformations() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ String transformName = "t";
+ config.put(TRANSFORMS_CONFIG, transformName + ", " + transformName);
+ config.put(TRANSFORMS_CONFIG + "." + transformName + ".type", Filter.class.getName());
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when the same transformation is specified multiple times"
+ );
+ }
+
+ @Test
+ public void testConnectorHasMissingTransformClass() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ String transformName = "t";
+ config.put(TRANSFORMS_CONFIG, transformName);
+ config.put(TRANSFORMS_CONFIG + "." + transformName + ".type", "WheresTheFruit");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when a transformation with a class not found on the worker is specified"
+ );
+ }
+
+ @Test
+ public void testConnectorHasInvalidTransformClass() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ String transformName = "t";
+ config.put(TRANSFORMS_CONFIG, transformName);
+ config.put(TRANSFORMS_CONFIG + "." + transformName + ".type", MonitorableSinkConnector.class.getName());
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when a transformation with a class of the wrong type is specified"
+ );
+ }
+
+ @Test
+ public void testConnectorHasNegationForUndefinedPredicate() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ String transformName = "t";
+ config.put(TRANSFORMS_CONFIG, transformName);
+ config.put(TRANSFORMS_CONFIG + "." + transformName + ".type", Filter.class.getName());
+ config.put(TRANSFORMS_CONFIG + "." + transformName + ".negate", "true");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when an undefined predicate is negated"
+ );
+ }
+
+ @Test
+ public void testConnectorHasDuplicatePredicates() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ String predicateName = "p";
+ config.put(PREDICATES_CONFIG, predicateName + ", " + predicateName);
+ config.put(PREDICATES_CONFIG + "." + predicateName + ".type", RecordIsTombstone.class.getName());
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when the same predicate is specified multiple times"
+ );
+ }
+
+ @Test
+ public void testConnectorHasMissingPredicateClass() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ String predicateName = "p";
+ config.put(PREDICATES_CONFIG, predicateName);
+ config.put(PREDICATES_CONFIG + "." + predicateName + ".type", "WheresTheFruit");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when a predicate with a class not found on the worker is specified"
+ );
+ }
+
+ @Test
+ public void testConnectorHasInvalidPredicateClass() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ String predicateName = "p";
+ config.put(PREDICATES_CONFIG, predicateName);
+ config.put(PREDICATES_CONFIG + "." + predicateName + ".type", MonitorableSinkConnector.class.getName());
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when a predicate with a class of the wrong type is specified"
+ );
+ }
+
+ @Test
+ public void testConnectorHasMissingConverterClass() throws InterruptedException {
+ Map config = defaultSinkConnectorProps();
+ config.put(KEY_CONVERTER_CLASS_CONFIG, "WheresTheFruit");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when a converter with a class not found on the worker is specified"
+ );
+ }
+
+ @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"
+ );
+ }
+
+ @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();
+ config.put(HEADER_CONVERTER_CLASS_CONFIG, "WheresTheFruit");
+ connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(
+ config.get(CONNECTOR_CLASS_CONFIG),
+ config,
+ 1,
+ "Connector config should fail preflight validation when a header converter with a class not found on the worker is specified"
+ );
+ }
+
+ @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"
+ );
+ }
+
+ @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"
+ );
+ }
+
+ @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"
+ );
+ }
+
+ public static class TestConverter implements Converter, HeaderConverter {
+ @Override
+ public void close() {
+ }
+
+ @Override
+ public void configure(Map configs) {
+ }
+
+ @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 SchemaAndValue toConnectHeader(String topic, String headerKey, byte[] value) {
+ return null;
+ }
+
+ @Override
+ public byte[] fromConnectHeader(String topic, String headerKey, Schema schema, Object 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");
+ }
+ }
+
+ 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<>();
+ props.put(NAME_CONFIG, "source-connector");
+ props.put(CONNECTOR_CLASS_CONFIG, MonitorableSourceConnector.class.getSimpleName());
+ props.put(TASKS_MAX_CONFIG, "1");
+ props.put(TOPIC_CONFIG, "t1");
+ props.put(KEY_CONVERTER_CLASS_CONFIG, StringConverter.class.getName());
+ props.put(VALUE_CONVERTER_CLASS_CONFIG, StringConverter.class.getName());
+ return props;
+ }
+
+ private Map defaultSinkConnectorProps() {
+ // setup up props for the sink connector
+ Map props = new HashMap<>();
+ props.put(NAME_CONFIG, "sink-connector");
+ props.put(CONNECTOR_CLASS_CONFIG, MonitorableSinkConnector.class.getSimpleName());
+ props.put(TASKS_MAX_CONFIG, "1");
+ props.put(TOPICS_CONFIG, "t1");
+ props.put(KEY_CONVERTER_CLASS_CONFIG, StringConverter.class.getName());
+ props.put(VALUE_CONVERTER_CLASS_CONFIG, StringConverter.class.getName());
+ return props;
+ }
+
+}
diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSourceConnector.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSourceConnector.java
index afd93257e584b..390927cee183b 100644
--- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSourceConnector.java
+++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSourceConnector.java
@@ -114,8 +114,8 @@ public void start(Map props) {
taskId = props.get("task.id");
connectorName = props.get("connector.name");
topicName = props.getOrDefault(TOPIC_CONFIG, "sequential-topic");
- throughput = Long.valueOf(props.getOrDefault("throughput", "-1"));
- batchSize = Integer.valueOf(props.getOrDefault("messages.per.poll", "1"));
+ throughput = Long.parseLong(props.getOrDefault("throughput", "-1"));
+ batchSize = Integer.parseInt(props.getOrDefault("messages.per.poll", "1"));
taskHandle = RuntimeHandles.get().connectorHandle(connectorName).taskHandle(taskId);
Map offset = Optional.ofNullable(
context.offsetStorageReader().offset(Collections.singletonMap("task.id", taskId)))
diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java
index 8c8d00da1b7cd..98d16a0f28fbd 100644
--- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java
+++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/AbstractHerderTest.java
@@ -19,7 +19,6 @@
import org.apache.kafka.clients.CommonClientConfigs;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.common.config.ConfigDef;
-import org.apache.kafka.common.config.ConfigException;
import org.apache.kafka.common.config.ConfigTransformer;
import org.apache.kafka.common.config.ConfigValue;
import org.apache.kafka.common.config.SaslConfigs;
@@ -420,7 +419,15 @@ public void testConfigValidationInvalidTopics() {
config.put(SinkConnectorConfig.TOPICS_CONFIG, "topic1,topic2");
config.put(SinkConnectorConfig.TOPICS_REGEX_CONFIG, "topic.*");
- assertThrows(ConfigException.class, () -> herder.validateConnectorConfig(config, false));
+ ConfigInfos validation = herder.validateConnectorConfig(config, false);
+
+ ConfigInfo topicsListInfo = findInfo(validation, SinkConnectorConfig.TOPICS_CONFIG);
+ assertNotNull(topicsListInfo);
+ assertEquals(1, topicsListInfo.configValue().errors().size());
+
+ ConfigInfo topicsRegexInfo = findInfo(validation, SinkConnectorConfig.TOPICS_REGEX_CONFIG);
+ assertNotNull(topicsRegexInfo);
+ assertEquals(1, topicsRegexInfo.configValue().errors().size());
verifyAll();
}
@@ -435,7 +442,11 @@ public void testConfigValidationTopicsWithDlq() {
config.put(SinkConnectorConfig.TOPICS_CONFIG, "topic1");
config.put(SinkConnectorConfig.DLQ_TOPIC_NAME_CONFIG, "topic1");
- assertThrows(ConfigException.class, () -> herder.validateConnectorConfig(config, false));
+ ConfigInfos validation = herder.validateConnectorConfig(config, false);
+
+ ConfigInfo topicsListInfo = findInfo(validation, SinkConnectorConfig.TOPICS_CONFIG);
+ assertNotNull(topicsListInfo);
+ assertEquals(1, topicsListInfo.configValue().errors().size());
verifyAll();
}
@@ -450,12 +461,15 @@ public void testConfigValidationTopicsRegexWithDlq() {
config.put(SinkConnectorConfig.TOPICS_REGEX_CONFIG, "topic.*");
config.put(SinkConnectorConfig.DLQ_TOPIC_NAME_CONFIG, "topic1");
- assertThrows(ConfigException.class, () -> herder.validateConnectorConfig(config, false));
+ ConfigInfos validation = herder.validateConnectorConfig(config, false);
+
+ ConfigInfo topicsRegexInfo = findInfo(validation, SinkConnectorConfig.TOPICS_REGEX_CONFIG);
+ assertNotNull(topicsRegexInfo);
+ assertEquals(1, topicsRegexInfo.configValue().errors().size());
verifyAll();
}
- @SuppressWarnings({"rawtypes", "unchecked"})
@Test()
public void testConfigValidationTransformsExtendResults() throws Throwable {
AbstractHerder herder = createConfigValidationHerder(TestSourceConnector.class, noneConnectorClientConfigOverridePolicy);
@@ -492,7 +506,7 @@ public void testConfigValidationTransformsExtendResults() throws Throwable {
"Transforms: xformB"
);
assertEquals(expectedGroups, result.groups());
- assertEquals(2, result.errorCount());
+ assertEquals(1, result.errorCount());
Map infos = result.values().stream()
.collect(Collectors.toMap(info -> info.configKey().name(), Function.identity()));
assertEquals(22, infos.size());
@@ -552,7 +566,7 @@ public void testConfigValidationPredicatesExtendResults() {
"Predicates: predY"
);
assertEquals(expectedGroups, result.groups());
- assertEquals(2, result.errorCount());
+ assertEquals(1, result.errorCount());
Map infos = result.values().stream()
.collect(Collectors.toMap(info -> info.configKey().name(), Function.identity()));
assertEquals(24, infos.size());
diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java
index 245bb7544396d..94cdfc237cfd2 100644
--- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java
+++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/DistributedHerderTest.java
@@ -802,7 +802,7 @@ public void testConnectorNameConflictsWithWorkerGroupId() throws Exception {
// CONN2 creation should fail because the worker group id (connect-test-group) conflicts with
// the consumer group id we would use for this sink
Map validatedConfigs =
- herder.validateBasicConnectorConfig(connectorMock, ConnectorConfig.configDef(), config);
+ herder.validateSinkConnectorConfig(ConnectorConfig.configDef(), config);
ConfigValue nameConfig = validatedConfigs.get(ConnectorConfig.NAME_CONFIG);
assertNotNull(nameConfig.errorMessages());
diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedConnectClusterAssertions.java b/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedConnectClusterAssertions.java
index edd99c8042cc1..1fbc8ce841143 100644
--- a/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedConnectClusterAssertions.java
+++ b/connect/runtime/src/test/java/org/apache/kafka/connect/util/clusters/EmbeddedConnectClusterAssertions.java
@@ -45,7 +45,7 @@ public class EmbeddedConnectClusterAssertions {
private static final Logger log = LoggerFactory.getLogger(EmbeddedConnectClusterAssertions.class);
public static final long WORKER_SETUP_DURATION_MS = TimeUnit.SECONDS.toMillis(60);
- public static final long VALIDATION_DURATION_MS = TimeUnit.SECONDS.toMillis(30);
+ public static final long VALIDATION_DURATION_MS = TimeUnit.SECONDS.toMillis(5);
public static final long CONNECTOR_SETUP_DURATION_MS = TimeUnit.SECONDS.toMillis(30);
private static final long CONNECT_INTERNAL_TOPIC_UPDATES_DURATION_MS = TimeUnit.SECONDS.toMillis(60);
@@ -243,7 +243,7 @@ public void assertExactlyNumErrorsOnConnectorConfigValidation(String connectorCl
connectorClass,
connConfig,
numErrors,
- (actual, expected) -> actual == expected
+ Number::equals
).orElse(false),
VALIDATION_DURATION_MS,
"Didn't meet the exact requested number of validation errors: " + numErrors);