From 643e74e730ab86f724ac1bead403c22d6fc478a2 Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Mon, 11 May 2020 15:36:26 +0100 Subject: [PATCH 01/13] KIP-585/KAFKA-9673: Filter and Conditional SMTs * Add Predicate interface * Add Filter SMT * Add the predicate implementations defined in the KIP. * Create abstraction in ConnectorConfig for configuring Transformations and Connectors with the "alias prefix" mechanism * Add tests and fix existing tests. --- .../transforms/predicates/Predicate.java | 48 +++ .../connect/runtime/ConnectorConfig.java | 368 +++++++++++++----- .../runtime/PredicatedTransformation.java | 67 ++++ .../isolation/DelegatingClassLoader.java | 10 + .../runtime/isolation/PluginScanResult.java | 8 + .../connect/runtime/isolation/Plugins.java | 5 + .../connect/runtime/AbstractHerderTest.java | 106 ++++- .../connect/runtime/ConnectorConfigTest.java | 195 +++++++++- .../runtime/PredicatedTransformationTest.java | 126 ++++++ .../kafka/connect/transforms/Filter.java | 51 +++ .../transforms/predicates/HasHeaderKey.java | 54 +++ .../predicates/RecordIsTombstone.java | 48 +++ .../predicates/TopicNameMatches.java | 72 ++++ .../predicates/HasHeaderKeyTest.java | 99 +++++ .../predicates/TopicNameMatchesTest.java | 65 ++++ 15 files changed, 1208 insertions(+), 114 deletions(-) create mode 100644 connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java create mode 100644 connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java create mode 100644 connect/runtime/src/test/java/org/apache/kafka/connect/runtime/PredicatedTransformationTest.java create mode 100644 connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java create mode 100644 connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java create mode 100644 connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java create mode 100644 connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java create mode 100644 connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java create mode 100644 connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java diff --git a/connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java b/connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java new file mode 100644 index 0000000000000..d50efd3752b00 --- /dev/null +++ b/connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java @@ -0,0 +1,48 @@ +/* + * 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.transforms.predicates; + +import org.apache.kafka.common.Configurable; +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.connect.connector.ConnectRecord; + +/** + *

A predicate on records. + * Predicates can be used to conditionally apply a {@link org.apache.kafka.connect.transforms.Transformation} + * by configuring the transformation's {@code predicate} (and {@code negate}) configuration parameters. + * In particular, the {@code Filter} transformation can be conditionally applied in order to filter + * certain records from further processing. + * + *

Implementations of this interface must be public and have a public constructor with no parameters. + * + * @param The type of record. + */ +public interface Predicate> extends Configurable, AutoCloseable { + + /** + * Configuration specification for this predicate. + */ + ConfigDef config(); + + /** + * Returns whether the given record satisfies this predicate. + */ + boolean test(R record); + + @Override + void close(); +} \ No newline at end of file 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 1cf93a4fd6466..49f2b28cc23a2 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,9 @@ import org.apache.kafka.connect.runtime.isolation.PluginDesc; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.transforms.Transformation; +import org.apache.kafka.connect.transforms.predicates.Predicate; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import java.lang.reflect.Modifier; import java.util.ArrayList; @@ -38,6 +41,7 @@ import java.util.List; import java.util.Locale; import java.util.Map; +import java.util.Set; import java.util.stream.Collectors; import java.util.stream.Stream; @@ -57,8 +61,11 @@ *

*/ public class ConnectorConfig extends AbstractConfig { + private static final Logger log = LoggerFactory.getLogger(ConnectorConfig.class); + protected static final String COMMON_GROUP = "Common"; protected static final String TRANSFORMS_GROUP = "Transforms"; + protected static final String PREDICATES_GROUP = "Predicates"; protected static final String ERROR_GROUP = "Error Handling"; public static final String NAME_CONFIG = "name"; @@ -98,6 +105,10 @@ public class ConnectorConfig extends AbstractConfig { private static final String TRANSFORMS_DOC = "Aliases for the transformations to be applied to records."; private static final String TRANSFORMS_DISPLAY = "Transforms"; + public static final String PREDICATES_CONFIG = "predicates"; + private static final String PREDICATES_DOC = "Aliases for the predicates used by transformations."; + private static final String PREDICATES_DISPLAY = "Predicates"; + public static final String CONFIG_RELOAD_ACTION_CONFIG = "config.action.reload"; private static final String CONFIG_RELOAD_ACTION_DOC = "The action that Connect should take on the connector when changes in external " + @@ -170,21 +181,8 @@ public static ConfigDef configDef() { .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(TRANSFORMS_CONFIG, Type.LIST, Collections.emptyList(), ConfigDef.CompositeValidator.of(new ConfigDef.NonNullValidator(), new ConfigDef.Validator() { - @SuppressWarnings("unchecked") - @Override - public void ensureValid(String name, Object value) { - final List transformAliases = (List) value; - if (transformAliases.size() > new HashSet<>(transformAliases).size()) { - throw new ConfigException(name, value, "Duplicate alias provided."); - } - } - - @Override - public String toString() { - return "unique transformation aliases"; - } - }), Importance.LOW, TRANSFORMS_DOC, TRANSFORMS_GROUP, ++orderInGroup, Width.LONG, TRANSFORMS_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, in(CONFIG_RELOAD_ACTION_NONE, CONFIG_RELOAD_ACTION_RESTART), Importance.LOW, CONFIG_RELOAD_ACTION_DOC, COMMON_GROUP, ++orderInGroup, Width.MEDIUM, CONFIG_RELOAD_ACTION_DISPLAY) @@ -201,6 +199,24 @@ public String toString() { ERRORS_LOG_INCLUDE_MESSAGES_DOC, ERROR_GROUP, ++orderInErrorGroup, Width.SHORT, ERRORS_LOG_INCLUDE_MESSAGES_DISPLAY); } + private static ConfigDef.CompositeValidator aliasValidator(String kind) { + return ConfigDef.CompositeValidator.of(new ConfigDef.NonNullValidator(), new ConfigDef.Validator() { + @SuppressWarnings("unchecked") + @Override + public void ensureValid(String name, Object value) { + final List aliases = (List) value; + if (aliases.size() > new HashSet<>(aliases).size()) { + throw new ConfigException(name, value, "Duplicate alias provided."); + } + } + + @Override + public String toString() { + return "unique " + kind + " aliases"; + } + }); + } + public ConnectorConfig(Plugins plugins) { this(plugins, new HashMap()); } @@ -257,12 +273,25 @@ public > List> transformations() { final List> transformations = new ArrayList<>(transformAliases.size()); for (String alias : transformAliases) { final String prefix = TRANSFORMS_CONFIG + "." + alias + "."; + try { @SuppressWarnings("unchecked") final Transformation transformation = getClass(prefix + "type").asSubclass(Transformation.class) .getDeclaredConstructor().newInstance(); - transformation.configure(originalsWithPrefix(prefix)); - transformations.add(transformation); + Map configs = originalsWithPrefix(prefix); + Object predicateAlias = configs.remove("predicate"); + Object negate = configs.remove("negate"); + transformation.configure(configs); + if (predicateAlias != null) { + String predicatePrefix = "predicates." + predicateAlias + "."; + @SuppressWarnings("unchecked") + Predicate predicate = getClass(predicatePrefix + "type").asSubclass(Predicate.class) + .getDeclaredConstructor().newInstance(); + predicate.configure(originalsWithPrefix(predicatePrefix)); + transformations.add(new PredicatedTransformation<>(predicate, negate == null ? false : Boolean.parseBoolean(negate.toString()), transformation)); + } else { + transformations.add(transformation); + } } catch (Exception e) { throw new ConnectException(e); } @@ -276,116 +305,247 @@ public > List> transformations() { *

* {@code requireFullConfig} specifies whether required config values that are missing should cause an exception to be thrown. */ + @SuppressWarnings({"rawtypes", "unchecked"}) public static ConfigDef enrich(Plugins plugins, ConfigDef baseConfigDef, Map props, boolean requireFullConfig) { - Object transformAliases = ConfigDef.parseType(TRANSFORMS_CONFIG, props.get(TRANSFORMS_CONFIG), Type.LIST); - if (!(transformAliases instanceof List)) { - return baseConfigDef; - } - ConfigDef newDef = new ConfigDef(baseConfigDef); - LinkedHashSet uniqueTransformAliases = new LinkedHashSet<>((List) transformAliases); - for (Object o : uniqueTransformAliases) { - if (!(o instanceof String)) { - throw new ConfigException("Item in " + TRANSFORMS_CONFIG + " property is not of " - + "type String"); + new EnrichablePlugin>("transformation", TRANSFORMS_CONFIG, TRANSFORMS_GROUP, (Class) Transformation.class, + props, requireFullConfig) { + @SuppressWarnings("rawtypes") + @Override + protected Set>> plugins() { + return (Set) plugins.transformations(); } - String alias = (String) o; - final String prefix = TRANSFORMS_CONFIG + "." + alias + "."; - final String group = TRANSFORMS_GROUP + ": " + alias; - int orderInGroup = 0; - - final String transformationTypeConfig = prefix + "type"; - final ConfigDef.Validator typeValidator = new ConfigDef.Validator() { - @Override - public void ensureValid(String name, Object value) { - getConfigDefFromTransformation(transformationTypeConfig, (Class) value); - } - }; - newDef.define(transformationTypeConfig, Type.CLASS, ConfigDef.NO_DEFAULT_VALUE, typeValidator, Importance.HIGH, - "Class for the '" + alias + "' transformation.", group, orderInGroup++, Width.LONG, "Transformation type for " + alias, - Collections.emptyList(), new TransformationClassRecommender(plugins)); - final ConfigDef transformationConfigDef; - try { - final String className = props.get(transformationTypeConfig); - final Class cls = (Class) ConfigDef.parseType(transformationTypeConfig, className, Type.CLASS); - transformationConfigDef = getConfigDefFromTransformation(transformationTypeConfig, cls); - } catch (ConfigException e) { - if (requireFullConfig) { - throw e; - } else { - continue; + + @Override + protected ConfigDef initialConfigDef() { + // All Transformations get these config parameters implicitly + return super.initialConfigDef() + .define("predicate", Type.STRING, "", Importance.MEDIUM, + "The alias of a predicate used to determine whether to apply this transformation.") + .define("negate", Type.BOOLEAN, false, Importance.MEDIUM, + "Whether the configured predicate should be negated."); + } + + @Override + protected Stream> configDefsForClass(String typeConfig) { + return super.configDefsForClass(typeConfig) + .filter(entry -> { + // The implicit parameters mask any from the transformer with the same name + if ("predicate".equals(entry.getValue()) || "negate".equals(entry.getValue())) { + log.warn("Transformer config " + entry.getValue() + " is masked by implicit config of that name"); + return false; + } else { + return true; + } + }); + } + + @Override + protected ConfigDef config(Transformation transformation) { + return transformation.config(); + } + + @Override + protected void validateProps(String prefix) { + if (props.containsKey(prefix + "negate") && + !props.containsKey(prefix + "predicate")) { + throw new ConfigException("Config '" + prefix + "negate' provided but there is no config '" + prefix + "predicate' to be negated."); } } + }.enrich(newDef); - newDef.embed(prefix, group, orderInGroup, transformationConfigDef); - } + new EnrichablePlugin>("predicate", PREDICATES_CONFIG, PREDICATES_GROUP, (Class) Predicate.class, props, requireFullConfig) { + @Override + protected Set>> plugins() { + return (Set) plugins.predicates(); + } + @Override + protected ConfigDef config(Predicate predicate) { + return predicate.config(); + } + }.enrich(newDef); return newDef; } /** - * Return {@link ConfigDef} from {@code transformationCls}, which is expected to be a non-null {@code Class}, - * by instantiating it and invoking {@link Transformation#config()}. + * An abstraction over "enrichable plugins" ({@link Transformation}s and {@link Predicate}s) used for computing the + * contribution to a Connectors ConfigDef. + * + * This is not entirely elegant because + * although they basically use the same "alias prefix" configuration idiom there are some differences. + * The abstract method pattern is used to cope with this. + * @param The type of plugin (either {@code Transformation} or {@code Predicate}). */ - static ConfigDef getConfigDefFromTransformation(String key, Class transformationCls) { - if (transformationCls == null || !Transformation.class.isAssignableFrom(transformationCls)) { - throw new ConfigException(key, String.valueOf(transformationCls), "Not a Transformation"); - } - if (Modifier.isAbstract(transformationCls.getModifiers())) { - String childClassNames = Stream.of(transformationCls.getClasses()) - .filter(transformationCls::isAssignableFrom) - .filter(c -> !Modifier.isAbstract(c.getModifiers())) - .filter(c -> Modifier.isPublic(c.getModifiers())) - .map(Class::getName) - .collect(Collectors.joining(", ")); - String message = childClassNames.trim().isEmpty() ? - "Transformation is abstract and cannot be created." : - "Transformation is abstract and cannot be created. Did you mean " + childClassNames + "?"; - throw new ConfigException(key, String.valueOf(transformationCls), message); + static abstract class EnrichablePlugin { + + private final String aliasKind; + private final String aliasConfig; + private final String aliasGroup; + private final Class baseClass; + private final Map props; + private final boolean requireFullConfig; + + public EnrichablePlugin( + String aliasKind, + String aliasConfig, String aliasGroup, Class baseClass, + Map props, boolean requireFullConfig) { + this.aliasKind = aliasKind; + this.aliasConfig = aliasConfig; + this.aliasGroup = aliasGroup; + this.baseClass = baseClass; + this.props = props; + this.requireFullConfig = requireFullConfig; } - Transformation transformation; - try { - transformation = transformationCls.asSubclass(Transformation.class).getConstructor().newInstance(); - } catch (Exception e) { - ConfigException exception = new ConfigException(key, String.valueOf(transformationCls), "Error getting config definition from Transformation: " + e.getMessage()); - exception.initCause(e); - throw exception; + + /** Add the configs for this alias to the given {@code ConfigDef}. */ + void enrich(ConfigDef newDef) { + Object aliases = ConfigDef.parseType(aliasConfig, props.get(aliasConfig), Type.LIST); + if (!(aliases instanceof List)) { + return; + } + + LinkedHashSet uniqueAliases = new LinkedHashSet<>((List) aliases); + for (Object o : uniqueAliases) { + if (!(o instanceof String)) { + throw new ConfigException("Item in " + aliasConfig + " property is not of " + + "type String"); + } + String alias = (String) o; + final String prefix = aliasConfig + "." + alias + "."; + final String group = aliasGroup + ": " + alias; + int orderInGroup = 0; + + final String typeConfig = prefix + "type"; + final ConfigDef.Validator typeValidator = new ConfigDef.Validator() { + @Override + public void ensureValid(String name, Object value) { + validateProps(prefix); + getConfigDefFromConfigProvidingClass(typeConfig, (Class) value); + } + }; + newDef.define(typeConfig, Type.CLASS, ConfigDef.NO_DEFAULT_VALUE, typeValidator, Importance.HIGH, + "Class for the '" + alias + "' " + aliasKind + ".", group, orderInGroup++, Width.LONG, + baseClass.getSimpleName() + " type for " + alias, + Collections.emptyList(), new ClassRecommender()); + + final ConfigDef configDef = populateConfigDef(typeConfig); + if (configDef == null) continue; + newDef.embed(prefix, group, orderInGroup, configDef); + } } - ConfigDef configDef = transformation.config(); - if (null == configDef) { - throw new ConnectException( - String.format( - "%s.config() must return a ConfigDef that is not null.", - transformationCls.getName() - ) - ); + + /** Subclasses can add extra validation of the {@link #props}. */ + protected void validateProps(String prefix) { } + + /** + * Populates the ConfigDef according to the configs returned from {@code configs()} method of class + * named in the {@code ...type} parameter of the {@code props}. + */ + protected ConfigDef populateConfigDef(String typeConfig) { + final ConfigDef configDef = initialConfigDef(); + try { + configDefsForClass(typeConfig) + .forEach(entry -> configDef.define(entry.getValue())); + + } catch (ConfigException e) { + if (requireFullConfig) { + throw e; + } else { + return null; + } + } + return configDef; } - return configDef; - } - /** - * Recommend bundled transformations. - */ - static final class TransformationClassRecommender implements ConfigDef.Recommender { - private final Plugins plugins; + /** + * Return a stream of configs provided by the {@code configs()} method of class + * named in the {@code ...type} parameter of the {@code props}. + */ + protected Stream> configDefsForClass(String typeConfig) { + final Class cls = (Class) ConfigDef.parseType(typeConfig, props.get(typeConfig), Type.CLASS); + return getConfigDefFromConfigProvidingClass(typeConfig, cls) + .configKeys().entrySet().stream(); + } - TransformationClassRecommender(Plugins plugins) { - this.plugins = plugins; + /** Get an initial ConfigDef */ + protected ConfigDef initialConfigDef() { + return new ConfigDef(); } - @Override - public List validValues(String name, Map parsedConfig) { - List transformationPlugins = new ArrayList<>(); - for (PluginDesc plugin : plugins.transformations()) { - transformationPlugins.add(plugin.pluginClass()); + /** + * Return {@link ConfigDef} from {@code cls}, which is expected to be a non-null {@code Class}, + * by instantiating it and invoking {@link #config(T)}. + * @param key + * @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 (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 aliasKind = this.aliasKind.substring(0, 1).toUpperCase(Locale.ENGLISH) + this.aliasKind.substring(1); + String message = childClassNames.trim().isEmpty() ? + 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); } - return Collections.unmodifiableList(transformationPlugins); + T transformation; + try { + transformation = cls.asSubclass(baseClass).getConstructor().newInstance(); + } catch (Exception e) { + throw new ConfigException(key, String.valueOf(cls), "Error getting config definition from " + baseClass.getSimpleName() + ": " + e.getMessage()); + } + ConfigDef configDef = config(transformation); + if (null == configDef) { + throw new ConnectException( + String.format( + "%s.config() must return a ConfigDef that is not null.", + cls.getName() + ) + ); + } + return configDef; } - @Override - public boolean visible(String name, Map parsedConfig) { - return true; + /** + * Get the ConfigDef from the given entity. + * This is necessary because there's no abstraction across {@link Transformation#config()} and + * {@link Predicate#config()}. + */ + protected abstract ConfigDef config(T t); + + /** + * The transformation or predicate plugins (as appropriate for T) to be used + * for the {@link ClassRecommender}. + */ + protected abstract Set> plugins(); + + /** + * Recommend bundled transformations or predicates. + */ + final class ClassRecommender implements ConfigDef.Recommender { + + @Override + public List validValues(String name, Map parsedConfig) { + List result = new ArrayList<>(); + for (PluginDesc plugin : plugins()) { + result.add(plugin.pluginClass()); + } + return Collections.unmodifiableList(result); + } + + @Override + public boolean visible(String name, Map parsedConfig) { + return true; + } } } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java new file mode 100644 index 0000000000000..117608a1e1ec8 --- /dev/null +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java @@ -0,0 +1,67 @@ +/* + * 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.runtime; + +import java.util.Map; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.common.utils.Utils; +import org.apache.kafka.connect.connector.ConnectRecord; +import org.apache.kafka.connect.transforms.Transformation; +import org.apache.kafka.connect.transforms.predicates.Predicate; + +/** + * Decorator for a {@link Transformation} which applies the delegate only when a + * {@link Predicate} is true (or false, according to {@code negate}). + * @param + */ +class PredicatedTransformation> implements Transformation { + + /*test*/ final Predicate predicate; + /*test*/ final Transformation delegate; + /*test*/ final boolean negate; + + PredicatedTransformation(Predicate predicate, boolean negate, Transformation delegate) { + this.predicate = predicate; + this.negate = negate; + this.delegate = delegate; + } + + @Override + public void configure(Map configs) { + + } + + @Override + public R apply(R record) { + if (negate ^ predicate.test(record)) { + return delegate.apply(record); + } + return record; + } + + @Override + public ConfigDef config() { + return null; + } + + @Override + public void close() { + Utils.closeQuietly(delegate, "predicated"); + Utils.closeQuietly(predicate, "predicate"); + } +} diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java index c4714c2f0f2ac..23f22e346584b 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java @@ -24,6 +24,7 @@ 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.reflections.Configuration; import org.reflections.Reflections; import org.reflections.ReflectionsException; @@ -72,6 +73,7 @@ public class DelegatingClassLoader extends URLClassLoader { private final SortedSet> converters; private final SortedSet> headerConverters; private final SortedSet> transformations; + private final SortedSet> predicates; private final SortedSet> configProviders; private final SortedSet> restExtensions; private final SortedSet> connectorClientConfigPolicies; @@ -92,6 +94,7 @@ public DelegatingClassLoader(List pluginPaths, ClassLoader parent) { this.converters = new TreeSet<>(); this.headerConverters = new TreeSet<>(); this.transformations = new TreeSet<>(); + this.predicates = new TreeSet<>(); this.configProviders = new TreeSet<>(); this.restExtensions = new TreeSet<>(); this.connectorClientConfigPolicies = new TreeSet<>(); @@ -121,6 +124,10 @@ public Set> transformations() { return transformations; } + public Set> predicates() { + return predicates; + } + public Set> configProviders() { return configProviders; } @@ -269,6 +276,8 @@ private void scanUrlsAndAddPlugins( headerConverters.addAll(plugins.headerConverters()); addPlugins(plugins.transformations(), loader); transformations.addAll(plugins.transformations()); + addPlugins(plugins.predicates(), loader); + predicates.addAll(plugins.predicates()); addPlugins(plugins.configProviders(), loader); configProviders.addAll(plugins.configProviders()); addPlugins(plugins.restExtensions(), loader); @@ -329,6 +338,7 @@ private PluginScanResult scanPluginPath( getPluginDesc(reflections, Converter.class, loader), getPluginDesc(reflections, HeaderConverter.class, loader), getPluginDesc(reflections, Transformation.class, loader), + getPluginDesc(reflections, Predicate.class, loader), getServiceLoaderPluginDesc(ConfigProvider.class, loader), getServiceLoaderPluginDesc(ConnectRestExtension.class, loader), getServiceLoaderPluginDesc(ConnectorClientConfigOverridePolicy.class, loader) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginScanResult.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginScanResult.java index e64a96c6f00a0..ac42dbb98a1fd 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginScanResult.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginScanResult.java @@ -23,6 +23,7 @@ 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 java.util.Arrays; import java.util.Collection; @@ -33,6 +34,7 @@ public class PluginScanResult { private final Collection> converters; private final Collection> headerConverters; private final Collection> transformations; + private final Collection> predicates; private final Collection> configProviders; private final Collection> restExtensions; private final Collection> connectorClientConfigPolicies; @@ -44,6 +46,7 @@ public PluginScanResult( Collection> converters, Collection> headerConverters, Collection> transformations, + Collection> predicates, Collection> configProviders, Collection> restExtensions, Collection> connectorClientConfigPolicies @@ -52,6 +55,7 @@ public PluginScanResult( this.converters = converters; this.headerConverters = headerConverters; this.transformations = transformations; + this.predicates = predicates; this.configProviders = configProviders; this.restExtensions = restExtensions; this.connectorClientConfigPolicies = connectorClientConfigPolicies; @@ -76,6 +80,10 @@ public Collection> transformations() { return transformations; } + public Collection> predicates() { + return predicates; + } + public Collection> configProviders() { return configProviders; } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java index d068a03ceecde..d507059eacc8c 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/Plugins.java @@ -33,6 +33,7 @@ import org.apache.kafka.connect.storage.ConverterType; 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; import org.slf4j.LoggerFactory; @@ -167,6 +168,10 @@ public Set> transformations() { return delegatingLoader.transformations(); } + public Set> predicates() { + return delegatingLoader.predicates(); + } + public Set> configProviders() { return delegatingLoader.configProviders(); } 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 97bba9468639b..844114b3c7bfa 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 @@ -42,6 +42,7 @@ import org.apache.kafka.connect.storage.ConfigBackingStore; import org.apache.kafka.connect.storage.StatusBackingStore; import org.apache.kafka.connect.transforms.Transformation; +import org.apache.kafka.connect.transforms.predicates.Predicate; import org.apache.kafka.connect.util.ConnectorTaskId; import org.easymock.Capture; import org.easymock.EasyMock; @@ -295,7 +296,8 @@ public void testConfigValidationMissingName() { assertEquals(2, result.errorCount()); Map infos = result.values().stream() .collect(Collectors.toMap(info -> info.configKey().name(), Function.identity())); - assertEquals(16, infos.size()); + // Base connector config has 14 fields, connector's configs add 2 + assertEquals(17, infos.size()); // Missing name should generate an error assertEquals(ConnectorConfig.NAME_CONFIG, infos.get(ConnectorConfig.NAME_CONFIG).configValue().name()); @@ -360,7 +362,7 @@ public void testConfigValidationTransformsExtendResults() { assertEquals(2, result.errorCount()); Map infos = result.values().stream() .collect(Collectors.toMap(info -> info.configKey().name(), Function.identity())); - assertEquals(19, infos.size()); + assertEquals(22, infos.size()); // Should get 2 type fields from the transforms, first adds its own config since it has a valid class assertEquals("transforms.xformA.type", infos.get("transforms.xformA.type").configValue().name()); @@ -373,6 +375,78 @@ public void testConfigValidationTransformsExtendResults() { verifyAll(); } + @Test() + public void testConfigValidationPredicatesExtendResults() { + AbstractHerder herder = createConfigValidationHerder(TestSourceConnector.class, noneConnectorClientConfigOverridePolicy); + + // 2 transform aliases defined -> 2 plugin lookups + Set> transformations = new HashSet<>(); + transformations.add(new PluginDesc(SampleTransformation.class, "1.0", classLoader)); + EasyMock.expect(plugins.transformations()).andReturn(transformations).times(1); + + Set> predicates = new HashSet<>(); + predicates.add(new PluginDesc(SamplePredicate.class, "1.0", classLoader)); + EasyMock.expect(plugins.predicates()).andReturn(predicates).times(2); + + replayAll(); + + // Define 2 transformations. One has a class defined and so can get embedded configs, the other is missing + // class info that should generate an error. + Map config = new HashMap<>(); + config.put(ConnectorConfig.CONNECTOR_CLASS_CONFIG, TestSourceConnector.class.getName()); + config.put(ConnectorConfig.NAME_CONFIG, "connector-name"); + config.put(ConnectorConfig.TRANSFORMS_CONFIG, "xformA"); + config.put(ConnectorConfig.TRANSFORMS_CONFIG + ".xformA.type", SampleTransformation.class.getName()); + config.put(ConnectorConfig.TRANSFORMS_CONFIG + ".xformA.predicate", "predX"); + config.put(ConnectorConfig.PREDICATES_CONFIG, "predX,predY"); + config.put(ConnectorConfig.PREDICATES_CONFIG + ".predX.type", SamplePredicate.class.getName()); + config.put("required", "value"); // connector required config + ConfigInfos result = herder.validateConnectorConfig(config); + assertEquals(herder.connectorTypeForClass(config.get(ConnectorConfig.CONNECTOR_CLASS_CONFIG)), ConnectorType.SOURCE); + + // We expect there to be errors due to the missing name and .... Note that these assertions depend heavily on + // the config fields for SourceConnectorConfig, but we expect these to change rarely. + assertEquals(TestSourceConnector.class.getName(), result.name()); + // Each transform also gets its own group + List expectedGroups = Arrays.asList( + ConnectorConfig.COMMON_GROUP, + ConnectorConfig.TRANSFORMS_GROUP, + ConnectorConfig.PREDICATES_GROUP, + ConnectorConfig.ERROR_GROUP, + SourceConnectorConfig.TOPIC_CREATION_GROUP, + "Transforms: xformA", + "Predicates: predX", + "Predicates: predY" + ); + assertEquals(expectedGroups, result.groups()); + assertEquals(2, result.errorCount()); + Map infos = result.values().stream() + .collect(Collectors.toMap(info -> info.configKey().name(), Function.identity())); + assertEquals(24, infos.size()); + // Should get 2 type fields from the transforms, first adds its own config since it has a valid class + assertEquals("transforms.xformA.type", + infos.get("transforms.xformA.type").configValue().name()); + assertTrue(infos.get("transforms.xformA.type").configValue().errors().isEmpty()); + assertEquals("transforms.xformA.subconfig", + infos.get("transforms.xformA.subconfig").configValue().name()); + assertEquals("transforms.xformA.predicate", + infos.get("transforms.xformA.predicate").configValue().name()); + assertTrue(infos.get("transforms.xformA.predicate").configValue().errors().isEmpty()); + assertEquals("transforms.xformA.negate", + infos.get("transforms.xformA.negate").configValue().name()); + assertTrue(infos.get("transforms.xformA.negate").configValue().errors().isEmpty()); + assertEquals("predicates.predX.type", + infos.get("predicates.predX.type").configValue().name()); + assertEquals("predicates.predX.predconfig", + infos.get("predicates.predX.predconfig").configValue().name()); + assertEquals("predicates.predY.type", + infos.get("predicates.predY.type").configValue().name()); + assertFalse( + infos.get("predicates.predY.type").configValue().errors().isEmpty()); + + verifyAll(); + } + @Test() public void testConfigValidationPrincipalOnlyOverride() { AbstractHerder herder = createConfigValidationHerder(TestSourceConnector.class, new PrincipalConnectorClientConfigOverridePolicy()); @@ -402,8 +476,8 @@ public void testConfigValidationPrincipalOnlyOverride() { ); assertEquals(expectedGroups, result.groups()); assertEquals(1, result.errorCount()); - // Base connector config has 13 fields, connector's configs add 2, and 2 producer overrides - assertEquals(18, result.values().size()); + // Base connector config has 14 fields, connector's configs add 2, and 2 producer overrides + assertEquals(19, result.values().size()); assertTrue(result.values().stream().anyMatch( configInfo -> ackConfigKey.equals(configInfo.configValue().name()) && !configInfo.configValue().errors().isEmpty())); assertTrue(result.values().stream().anyMatch( @@ -564,6 +638,30 @@ public void close() { } } + public static class SamplePredicate> implements Predicate { + + @Override + public ConfigDef config() { + return new ConfigDef() + .define("predconfig", ConfigDef.Type.STRING, "default", ConfigDef.Importance.LOW, "docs"); + } + + @Override + public boolean test(R record) { + return false; + } + + @Override + public void close() { + + } + + @Override + public void configure(Map configs) { + + } + } + // We need to use a real class here due to some issue with mocking java.lang.Class private abstract class BogusSourceConnector extends SourceConnector { } 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 48d494767228e..aeb3be0eedaba 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 @@ -16,6 +16,12 @@ */ package org.apache.kafka.connect.runtime; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; + import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.connect.connector.ConnectRecord; @@ -23,14 +29,10 @@ import org.apache.kafka.connect.runtime.isolation.PluginDesc; import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.transforms.Transformation; +import org.apache.kafka.connect.transforms.predicates.Predicate; +import org.junit.Ignore; import org.junit.Test; -import java.util.Collections; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Set; - import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -214,6 +216,187 @@ public void abstractKeyValueTransform() { } } + + + @Test(expected = ConfigException.class) + public void wrongPredicateType() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.predicate", "my-pred"); + props.put("predicates", "my-pred"); + props.put("predicates.my-pred.type", TestConnector.class.getName()); + new ConnectorConfig(MOCK_PLUGINS, props); + } + + @Test + public void singleConditionalTransform() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.predicate", "my-pred"); + props.put("transforms.a.negate", "true"); + props.put("predicates", "my-pred"); + props.put("predicates.my-pred.type", TestPredicate.class.getName()); + props.put("predicates.my-pred.int", "84"); + assertPredicatedTransform(props, true); + } + + @Test + public void predicateNegationDefaultsToFalse() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.predicate", "my-pred"); + props.put("predicates", "my-pred"); + props.put("predicates.my-pred.type", TestPredicate.class.getName()); + props.put("predicates.my-pred.int", "84"); + assertPredicatedTransform(props, false); + } + + @Test(expected = ConfigException.class) + public void abstractPredicate() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.predicate", "my-pred"); + props.put("predicates", "my-pred"); + props.put("predicates.my-pred.type", AbstractTestPredicate.class.getName()); + props.put("predicates.my-pred.int", "84"); + assertPredicatedTransform(props, false); + } + + private void assertPredicatedTransform(Map props, boolean expectedNegated) { + final ConnectorConfig config = new ConnectorConfig(MOCK_PLUGINS, props); + final List> transformations = config.transformations(); + assertEquals(1, transformations.size()); + assertTrue(transformations.get(0) instanceof PredicatedTransformation); + PredicatedTransformation predicated = (PredicatedTransformation) transformations.get(0); + + assertEquals(expectedNegated, predicated.negate); + + assertTrue(predicated.delegate instanceof ConnectorConfigTest.SimpleTransformation); + assertEquals(42, ((SimpleTransformation) predicated.delegate).magicNumber); + + assertTrue(predicated.predicate instanceof ConnectorConfigTest.TestPredicate); + assertEquals(84, ((TestPredicate) predicated.predicate).param); + + predicated.close(); + + assertEquals(0, ((SimpleTransformation) predicated.delegate).magicNumber); + assertEquals(0, ((TestPredicate) predicated.predicate).param); + } + + @Test + public void misconfiguredPredicate() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.predicate", "my-pred"); + props.put("transforms.a.negate", "true"); + props.put("predicates", "my-pred"); + props.put("predicates.my-pred.type", TestPredicate.class.getName()); + props.put("predicates.my-pred.int", "79"); + try { + new ConnectorConfig(MOCK_PLUGINS, props); + fail(); + } catch (ConfigException e) { + assertTrue(e.getMessage().contains("Value must be at least 80")); + } + } + + @Ignore("Is this really an error. There's no actual need for the predicates config (unlike transforms where it defines the order).") + @Test(expected = ConfigException.class) + public void missingPredicates() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.predicate", "my-pred"); + //props.put("predicates", "my-pred"); + props.put("predicates.my-pred.type", TestPredicate.class.getName()); + props.put("predicates.my-pred.int", "84"); + new ConnectorConfig(MOCK_PLUGINS, props); + } + + @Test(expected = ConfigException.class) + public void missingPredicateConfig() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.predicate", "my-pred"); + props.put("predicates", "my-pred"); + //props.put("predicates.my-pred.type", TestPredicate.class.getName()); + //props.put("predicates.my-pred.int", "84"); + new ConnectorConfig(MOCK_PLUGINS, props); + } + + @Test(expected = ConfigException.class) + public void negatedButNoPredicate() { + Map props = new HashMap<>(); + props.put("name", "test"); + props.put("connector.class", TestConnector.class.getName()); + props.put("transforms", "a"); + props.put("transforms.a.type", SimpleTransformation.class.getName()); + props.put("transforms.a.magic.number", "42"); + props.put("transforms.a.negate", "true"); + new ConnectorConfig(MOCK_PLUGINS, props); + } + + static class TestPredicate> implements Predicate { + + int param; + + public TestPredicate() { } + + @Override + public ConfigDef config() { + return new ConfigDef().define("int", ConfigDef.Type.INT, 80, ConfigDef.Range.atLeast(80), ConfigDef.Importance.MEDIUM, + "A test parameter"); + } + + @Override + public boolean test(R record) { + return false; + } + + @Override + public void close() { + param = 0; + } + + @Override + public void configure(Map configs) { + param = Integer.parseInt((String) configs.get("int")); + } + } + + static abstract class AbstractTestPredicate> implements Predicate { + + public AbstractTestPredicate() { } + + } + public static abstract class AbstractTransformation> implements Transformation { } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/PredicatedTransformationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/PredicatedTransformationTest.java new file mode 100644 index 0000000000000..75542e7352881 --- /dev/null +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/PredicatedTransformationTest.java @@ -0,0 +1,126 @@ +/* + * 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.runtime; + +import java.util.Map; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.connect.source.SourceRecord; +import org.apache.kafka.connect.transforms.Transformation; +import org.apache.kafka.connect.transforms.predicates.Predicate; +import org.junit.Test; + +import static java.util.Collections.singletonMap; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertTrue; + +public class PredicatedTransformationTest { + + private final SourceRecord initial = new SourceRecord(singletonMap("initial", 1), null, null, null, null); + private final SourceRecord transformed = new SourceRecord(singletonMap("transformed", 2), null, null, null, null); + + @Test + public void apply() { + applyAndAssert(true, false, transformed); + applyAndAssert(true, true, initial); + applyAndAssert(false, false, initial); + applyAndAssert(false, true, transformed); + } + + private void applyAndAssert(boolean predicateResult, boolean negate, + SourceRecord expectedResult) { + class TestTransformation implements Transformation { + + private boolean closed = false; + private SourceRecord transformedRecord; + + private TestTransformation(SourceRecord transformedRecord) { + this.transformedRecord = transformedRecord; + } + + @Override + public SourceRecord apply(SourceRecord record) { + return transformedRecord; + } + + @Override + public ConfigDef config() { + return null; + } + + @Override + public void close() { + closed = true; + } + + @Override + public void configure(Map configs) { + + } + + private void assertClosed() { + assertTrue("Transformer should be closed", closed); + } + } + + class TestPredicate implements Predicate { + + private boolean testResult; + private boolean closed = false; + + private TestPredicate(boolean testResult) { + this.testResult = testResult; + } + + @Override + public ConfigDef config() { + return null; + } + + @Override + public boolean test(SourceRecord record) { + return testResult; + } + + @Override + public void close() { + closed = true; + } + + @Override + public void configure(Map configs) { + + } + + private void assertClosed() { + assertTrue("Predicate should be closed", closed); + } + } + TestPredicate predicate = new TestPredicate(predicateResult); + TestTransformation predicatedTransform = new TestTransformation(transformed); + PredicatedTransformation pt = new PredicatedTransformation<>( + predicate, + negate, + predicatedTransform); + + assertEquals(expectedResult, pt.apply(initial)); + + pt.close(); + predicate.assertClosed(); + predicatedTransform.assertClosed(); + } +} diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java new file mode 100644 index 0000000000000..0b5a8218b54c3 --- /dev/null +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.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.transforms; + +import java.util.Map; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.connect.connector.ConnectRecord; + +/** + * Drops all records, filtering them from subsequent transformations in the chain. + * This is intended to be used conditionally to filter out records matching (or not matching) + * a particular {@link org.apache.kafka.connect.transforms.predicates.Predicate}. + * @param The type of record. + */ +public class Filter> implements Transformation { + + @Override + public R apply(R record) { + return null; + } + + @Override + public ConfigDef config() { + return null; + } + + @Override + public void close() { + + } + + @Override + public void configure(Map configs) { + + } +} diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java new file mode 100644 index 0000000000000..4750d050ea56a --- /dev/null +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java @@ -0,0 +1,54 @@ +/* + * 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.transforms.predicates; + +import java.util.Map; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.connect.connector.ConnectRecord; + +/** + * A predicate which is true for records with at least one header with the configured name. + * @param The type of connect record. + */ +public class HasHeaderKey> implements Predicate { + + private static final String NAME_CONFIG_KEY = "name"; + private String name; + + @Override + public ConfigDef config() { + return new ConfigDef().define(NAME_CONFIG_KEY, ConfigDef.Type.STRING, null, + new ConfigDef.NonEmptyString(), ConfigDef.Importance.MEDIUM, + "The header name."); + } + + @Override + public boolean test(R record) { + return record.headers().allWithName(name).hasNext(); + } + + @Override + public void close() { + + } + + @Override + public void configure(Map configs) { + this.name = (String) configs.get(NAME_CONFIG_KEY); + } +} diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java new file mode 100644 index 0000000000000..64f40cf7d827f --- /dev/null +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java @@ -0,0 +1,48 @@ +/* + * 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.transforms.predicates; + +import java.util.Map; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.connect.connector.ConnectRecord; + +/** + * A predicate which is true for records which are tombstones (i.e. have null key). + * @param The type of connect record. + */ +public class RecordIsTombstone> implements Predicate { + @Override + public ConfigDef config() { + return new ConfigDef(); + } + + @Override + public boolean test(R record) { + return record.key() == null; + } + + @Override + public void close() { + + } + + @Override + public void configure(Map configs) { + + } +} diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java new file mode 100644 index 0000000000000..a904e1c5a23c3 --- /dev/null +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java @@ -0,0 +1,72 @@ +/* + * 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.transforms.predicates; + +import java.util.Map; +import java.util.regex.Pattern; +import java.util.regex.PatternSyntaxException; + +import org.apache.kafka.common.config.ConfigDef; +import org.apache.kafka.common.config.ConfigException; +import org.apache.kafka.connect.connector.ConnectRecord; + +/** + * A predicate which is true for records with a topic name that matches the configured regular expression. + * @param The type of connect record. + */ +public class TopicNameMatches> implements Predicate { + + public static final String PATTERN_CONFIG_KEY = "pattern"; + private Pattern pattern; + + @Override + public ConfigDef config() { + return new ConfigDef().define(PATTERN_CONFIG_KEY, ConfigDef.Type.STRING, null, + new ConfigDef.Validator() { + @Override + public void ensureValid(String name, Object value) { + if (value != null) { + compile(name, value); + } + } + }, ConfigDef.Importance.MEDIUM, + "A Java regular expression for matching against the name of a record's topic."); + } + + private Pattern compile(String name, Object value) { + try { + return Pattern.compile((String) value); + } catch (PatternSyntaxException e) { + throw new ConfigException(name, value, "entry must be a Java-compatible regular expression: " + e.getMessage()); + } + } + + @Override + public boolean test(R record) { + return record.topic() != null && pattern.matcher(record.topic()).matches(); + } + + @Override + public void close() { + + } + + @Override + public void configure(Map configs) { + this.pattern = compile(PATTERN_CONFIG_KEY, configs.get(PATTERN_CONFIG_KEY)); + } +} diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java new file mode 100644 index 0000000000000..b2ad5450e6948 --- /dev/null +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java @@ -0,0 +1,99 @@ +/* + * 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.transforms.predicates; + +import java.util.Arrays; +import java.util.Collections; +import java.util.List; +import java.util.stream.Collectors; + +import org.apache.kafka.common.config.ConfigValue; +import org.apache.kafka.connect.data.Schema; +import org.apache.kafka.connect.header.Header; +import org.apache.kafka.connect.source.SourceRecord; +import org.junit.Test; + +import static java.util.Collections.singletonList; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class HasHeaderKeyTest { + + @Test + public void testConfig() { + HasHeaderKey predicate = new HasHeaderKey<>(); + predicate.config().validate(Collections.singletonMap("name", "foo")); + + List configs = predicate.config().validate(Collections.singletonMap("name", "")); + assertEquals(singletonList("Invalid value for configuration name: String must be non-empty"), configs.get(0).errorMessages()); + } + + @Test + public void testTest() { + HasHeaderKey predicate = new HasHeaderKey<>(); + predicate.configure(Collections.singletonMap("name", "foo")); + + assertTrue(predicate.test(recordWithHeaders("foo"))); + assertTrue(predicate.test(recordWithHeaders("foo", "bar"))); + assertTrue(predicate.test(recordWithHeaders("bar", "foo", "bar", "foo"))); + assertFalse(predicate.test(recordWithHeaders("bar"))); + assertFalse(predicate.test(recordWithHeaders("bar", "bar"))); + assertFalse(predicate.test(recordWithHeaders())); + assertFalse(predicate.test(new SourceRecord(null, null, null, null, null))); + + } + + private SourceRecord recordWithHeaders(String... headers) { + return new SourceRecord(null, null, null, null, null, null, null, null, null, + Arrays.stream(headers).map(header -> new TestHeader(header)).collect(Collectors.toList())); + } + + private static class TestHeader implements Header { + + private final String key; + + public TestHeader(String key) { + this.key = key; + } + + @Override + public String key() { + return key; + } + + @Override + public Schema schema() { + return null; + } + + @Override + public Object value() { + return null; + } + + @Override + public Header with(Schema schema, Object value) { + return null; + } + + @Override + public Header rename(String key) { + return null; + } + } +} diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java new file mode 100644 index 0000000000000..500f2b5cd6e31 --- /dev/null +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java @@ -0,0 +1,65 @@ +/* + * 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.transforms.predicates; + +import java.util.Collections; +import java.util.List; + +import org.apache.kafka.common.config.ConfigValue; +import org.apache.kafka.connect.source.SourceRecord; +import org.junit.Test; + +import static java.util.Collections.singletonList; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class TopicNameMatchesTest { + + @Test + public void testConfig() { + TopicNameMatches predicate = new TopicNameMatches<>(); + predicate.config().validate(Collections.singletonMap("pattern", "my-prefix-.*")); + + List configs = predicate.config().validate(Collections.singletonMap("pattern", "*")); + assertEquals(singletonList("Invalid value * for configuration pattern: " + + "entry must be a Java-compatible regular expression: " + + "Dangling meta character '*' near index 0" + System.lineSeparator() + + "*" + System.lineSeparator() + + "^"), + configs.get(0).errorMessages()); + } + + @Test + public void testTest() { + TopicNameMatches predicate = new TopicNameMatches<>(); + predicate.configure(Collections.singletonMap("pattern", "my-prefix-.*")); + + assertTrue(predicate.test(recordWithTopicName("my-prefix-"))); + assertTrue(predicate.test(recordWithTopicName("my-prefix-foo"))); + assertFalse(predicate.test(recordWithTopicName("x-my-prefix-"))); + assertFalse(predicate.test(recordWithTopicName("x-my-prefix-foo"))); + assertFalse(predicate.test(recordWithTopicName("your-prefix-"))); + assertFalse(predicate.test(recordWithTopicName("your-prefix-foo"))); + assertFalse(predicate.test(new SourceRecord(null, null, null, null, null))); + + } + + private SourceRecord recordWithTopicName(String topicName) { + return new SourceRecord(null, null, topicName, null, null); + } +} From 2cfe4e1d3a92afc4c5188b8b5f8ac3732c24b527 Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Wed, 20 May 2020 16:54:37 +0100 Subject: [PATCH 02/13] Use Predicates config group --- .../org/apache/kafka/connect/runtime/AbstractHerderTest.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) 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 844114b3c7bfa..bad2254cfe94a 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 @@ -292,7 +292,7 @@ public void testConfigValidationMissingName() { // the config fields for SourceConnectorConfig, but we expect these to change rarely. assertEquals(TestSourceConnector.class.getName(), result.name()); assertEquals(Arrays.asList(ConnectorConfig.COMMON_GROUP, ConnectorConfig.TRANSFORMS_GROUP, - ConnectorConfig.ERROR_GROUP, SourceConnectorConfig.TOPIC_CREATION_GROUP), result.groups()); + ConnectorConfig.PREDICATES_GROUP, ConnectorConfig.ERROR_GROUP, SourceConnectorConfig.TOPIC_CREATION_GROUP), result.groups()); assertEquals(2, result.errorCount()); Map infos = result.values().stream() .collect(Collectors.toMap(info -> info.configKey().name(), Function.identity())); @@ -353,6 +353,7 @@ public void testConfigValidationTransformsExtendResults() { List expectedGroups = Arrays.asList( ConnectorConfig.COMMON_GROUP, ConnectorConfig.TRANSFORMS_GROUP, + ConnectorConfig.PREDICATES_GROUP, ConnectorConfig.ERROR_GROUP, SourceConnectorConfig.TOPIC_CREATION_GROUP, "Transforms: xformA", @@ -471,6 +472,7 @@ public void testConfigValidationPrincipalOnlyOverride() { List expectedGroups = Arrays.asList( ConnectorConfig.COMMON_GROUP, ConnectorConfig.TRANSFORMS_GROUP, + ConnectorConfig.PREDICATES_GROUP, ConnectorConfig.ERROR_GROUP, SourceConnectorConfig.TOPIC_CREATION_GROUP ); From 0f0e1383f29d293441b70e69a28e077f0c343ce6 Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Fri, 22 May 2020 07:06:09 +0100 Subject: [PATCH 03/13] Some review comments --- .../transforms/predicates/HasHeaderKey.java | 17 ++++++---- .../predicates/RecordIsTombstone.java | 9 ++++-- .../predicates/TopicNameMatches.java | 31 ++++++++++--------- 3 files changed, 34 insertions(+), 23 deletions(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java index 4750d050ea56a..6501b401ef7a8 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java @@ -16,10 +16,13 @@ */ package org.apache.kafka.connect.transforms.predicates; +import java.util.Iterator; import java.util.Map; import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.connect.connector.ConnectRecord; +import org.apache.kafka.connect.header.Header; +import org.apache.kafka.connect.transforms.util.SimpleConfig; /** * A predicate which is true for records with at least one header with the configured name. @@ -27,19 +30,21 @@ */ public class HasHeaderKey> implements Predicate { - private static final String NAME_CONFIG_KEY = "name"; + private static final String NAME_CONFIG = "name"; + private static final ConfigDef CONFIG_DEF = new ConfigDef().define(NAME_CONFIG, ConfigDef.Type.STRING, null, + new ConfigDef.NonEmptyString(), ConfigDef.Importance.MEDIUM, + "The header name."); private String name; @Override public ConfigDef config() { - return new ConfigDef().define(NAME_CONFIG_KEY, ConfigDef.Type.STRING, null, - new ConfigDef.NonEmptyString(), ConfigDef.Importance.MEDIUM, - "The header name."); + return CONFIG_DEF; } @Override public boolean test(R record) { - return record.headers().allWithName(name).hasNext(); + Iterator
headerIterator = record.headers().allWithName(name); + return headerIterator != null && headerIterator.hasNext(); } @Override @@ -49,6 +54,6 @@ public void close() { @Override public void configure(Map configs) { - this.name = (String) configs.get(NAME_CONFIG_KEY); + this.name = new SimpleConfig(config(), configs).getString(NAME_CONFIG); } } diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java index 64f40cf7d827f..7d5c47e66c1bc 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java @@ -22,18 +22,21 @@ import org.apache.kafka.connect.connector.ConnectRecord; /** - * A predicate which is true for records which are tombstones (i.e. have null key). + * A predicate which is true for records which are tombstones (i.e. have null value). * @param The type of connect record. */ public class RecordIsTombstone> implements Predicate { + + private static final ConfigDef CONFIG_DEF = new ConfigDef(); + @Override public ConfigDef config() { - return new ConfigDef(); + return CONFIG_DEF; } @Override public boolean test(R record) { - return record.key() == null; + return record.value() == null; } @Override diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java index a904e1c5a23c3..2962a397c5125 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java @@ -23,6 +23,7 @@ import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.connect.connector.ConnectRecord; +import org.apache.kafka.connect.transforms.util.SimpleConfig; /** * A predicate which is true for records with a topic name that matches the configured regular expression. @@ -30,26 +31,27 @@ */ public class TopicNameMatches> implements Predicate { - public static final String PATTERN_CONFIG_KEY = "pattern"; + private static final String PATTERN_CONFIG = "pattern"; + private static final ConfigDef CONFIG_DEF = new ConfigDef().define(PATTERN_CONFIG, ConfigDef.Type.STRING, null, + new ConfigDef.Validator() { + @Override + public void ensureValid(String name, Object value) { + if (value instanceof String) { + compile(name, (String) value); + } + } + }, ConfigDef.Importance.MEDIUM, + "A Java regular expression for matching against the name of a record's topic."); private Pattern pattern; @Override public ConfigDef config() { - return new ConfigDef().define(PATTERN_CONFIG_KEY, ConfigDef.Type.STRING, null, - new ConfigDef.Validator() { - @Override - public void ensureValid(String name, Object value) { - if (value != null) { - compile(name, value); - } - } - }, ConfigDef.Importance.MEDIUM, - "A Java regular expression for matching against the name of a record's topic."); + return CONFIG_DEF; } - private Pattern compile(String name, Object value) { + private static Pattern compile(String name, String value) { try { - return Pattern.compile((String) value); + return Pattern.compile(value); } catch (PatternSyntaxException e) { throw new ConfigException(name, value, "entry must be a Java-compatible regular expression: " + e.getMessage()); } @@ -67,6 +69,7 @@ public void close() { @Override public void configure(Map configs) { - this.pattern = compile(PATTERN_CONFIG_KEY, configs.get(PATTERN_CONFIG_KEY)); + SimpleConfig simpleConfig = new SimpleConfig(config(), configs); + this.pattern = compile(PATTERN_CONFIG, simpleConfig.getString(PATTERN_CONFIG)); } } From 6477fe9cd5b537f515695a0099cd1df49bf711b1 Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Fri, 22 May 2020 07:07:07 +0100 Subject: [PATCH 04/13] fixup --- .../main/java/org/apache/kafka/connect/transforms/Filter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java index 0b5a8218b54c3..c805c84bf3e81 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java @@ -36,7 +36,7 @@ public R apply(R record) { @Override public ConfigDef config() { - return null; + return new ConfigDef(); } @Override From 460c5f826a8272d6fba06ec1926ef3c8e5fc3e2a Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Fri, 22 May 2020 07:21:22 +0100 Subject: [PATCH 05/13] Comment --- .../java/org/apache/kafka/common/utils/Utils.java | 12 ++++++++++++ .../kafka/connect/runtime/ConnectorConfig.java | 9 ++++----- 2 files changed, 16 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 7b47a1a3553cf..7ee81b3b8f219 100755 --- a/clients/src/main/java/org/apache/kafka/common/utils/Utils.java +++ b/clients/src/main/java/org/apache/kafka/common/utils/Utils.java @@ -342,6 +342,18 @@ public static Class loadClass(String klass, Class base) thro return Class.forName(klass, true, Utils.getContextOrKafkaClassLoader()).asSubclass(base); } + /** + * Cast {@code klass} to {@code base} and instantiate it. + * @param klass The class to instantiate + * @param base A know baseclass of klass. + * @param the type of the base class + * @throws ClassCastException If {@code klass} is not a subclass of {@code base}. + * @return the new instance. + */ + public static T newInstance(Class klass, Class base) { + return Utils.newInstance(klass.asSubclass(base)); + } + /** * Construct a new object using a class name and parameters. * 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 49f2b28cc23a2..d9d8d98e9d601 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 @@ -22,6 +22,7 @@ import org.apache.kafka.common.config.ConfigDef.Type; import org.apache.kafka.common.config.ConfigDef.Width; import org.apache.kafka.common.config.ConfigException; +import org.apache.kafka.common.utils.Utils; import org.apache.kafka.connect.connector.ConnectRecord; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.runtime.errors.ToleranceType; @@ -276,8 +277,7 @@ public > List> transformations() { try { @SuppressWarnings("unchecked") - final Transformation transformation = getClass(prefix + "type").asSubclass(Transformation.class) - .getDeclaredConstructor().newInstance(); + final Transformation transformation = Utils.newInstance(getClass(prefix + "type"), Transformation.class); Map configs = originalsWithPrefix(prefix); Object predicateAlias = configs.remove("predicate"); Object negate = configs.remove("negate"); @@ -285,8 +285,7 @@ public > List> transformations() { if (predicateAlias != null) { String predicatePrefix = "predicates." + predicateAlias + "."; @SuppressWarnings("unchecked") - Predicate predicate = getClass(predicatePrefix + "type").asSubclass(Predicate.class) - .getDeclaredConstructor().newInstance(); + Predicate predicate = Utils.newInstance(getClass(predicatePrefix + "type"), Predicate.class); predicate.configure(originalsWithPrefix(predicatePrefix)); transformations.add(new PredicatedTransformation<>(predicate, negate == null ? false : Boolean.parseBoolean(negate.toString()), transformation)); } else { @@ -499,7 +498,7 @@ ConfigDef getConfigDefFromConfigProvidingClass(String key, Class cls) { } T transformation; try { - transformation = cls.asSubclass(baseClass).getConstructor().newInstance(); + transformation = Utils.newInstance(cls, baseClass); } catch (Exception e) { throw new ConfigException(key, String.valueOf(cls), "Error getting config definition from " + baseClass.getSimpleName() + ": " + e.getMessage()); } From fba43c8fef6a14c7a73639bf3358aa1c52673c8e Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Fri, 22 May 2020 08:48:39 +0100 Subject: [PATCH 06/13] tidy and fix a couple of tests --- .../connect/runtime/ConnectorConfig.java | 31 +++++++++++-------- .../runtime/PredicatedTransformation.java | 15 ++++++--- .../connect/runtime/ConnectorConfigTest.java | 4 +-- 3 files changed, 30 insertions(+), 20 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 d9d8d98e9d601..4468e6ba857a6 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 @@ -279,8 +279,8 @@ public > List> transformations() { @SuppressWarnings("unchecked") final Transformation transformation = Utils.newInstance(getClass(prefix + "type"), Transformation.class); Map configs = originalsWithPrefix(prefix); - Object predicateAlias = configs.remove("predicate"); - Object negate = configs.remove("negate"); + Object predicateAlias = configs.remove(PredicatedTransformation.PREDICATE_CONFIG); + Object negate = configs.remove(PredicatedTransformation.NEGATE_CONFIG); transformation.configure(configs); if (predicateAlias != null) { String predicatePrefix = "predicates." + predicateAlias + "."; @@ -307,7 +307,7 @@ public > List> transformations() { @SuppressWarnings({"rawtypes", "unchecked"}) public static ConfigDef enrich(Plugins plugins, ConfigDef baseConfigDef, Map props, boolean requireFullConfig) { ConfigDef newDef = new ConfigDef(baseConfigDef); - new EnrichablePlugin>("transformation", TRANSFORMS_CONFIG, TRANSFORMS_GROUP, (Class) Transformation.class, + new EnrichablePlugin>("Transformation", TRANSFORMS_CONFIG, TRANSFORMS_GROUP, (Class) Transformation.class, props, requireFullConfig) { @SuppressWarnings("rawtypes") @Override @@ -320,9 +320,9 @@ protected Set>> plugins() { protected ConfigDef initialConfigDef() { // All Transformations get these config parameters implicitly return super.initialConfigDef() - .define("predicate", Type.STRING, "", Importance.MEDIUM, + .define(PredicatedTransformation.PREDICATE_CONFIG, Type.STRING, "", Importance.MEDIUM, "The alias of a predicate used to determine whether to apply this transformation.") - .define("negate", Type.BOOLEAN, false, Importance.MEDIUM, + .define(PredicatedTransformation.NEGATE_CONFIG, Type.BOOLEAN, false, Importance.MEDIUM, "Whether the configured predicate should be negated."); } @@ -331,8 +331,10 @@ protected Stream> configDefsForClass(Stri return super.configDefsForClass(typeConfig) .filter(entry -> { // The implicit parameters mask any from the transformer with the same name - if ("predicate".equals(entry.getValue()) || "negate".equals(entry.getValue())) { - log.warn("Transformer config " + entry.getValue() + " is masked by implicit config of that name"); + if (PredicatedTransformation.PREDICATE_CONFIG.equals(entry.getValue()) + || PredicatedTransformation.NEGATE_CONFIG.equals(entry.getValue())) { + log.warn("Transformer config {} is masked by implicit config of that name", + entry.getValue()); return false; } else { return true; @@ -347,14 +349,18 @@ protected ConfigDef config(Transformation transformation) { @Override protected void validateProps(String prefix) { - if (props.containsKey(prefix + "negate") && - !props.containsKey(prefix + "predicate")) { - throw new ConfigException("Config '" + prefix + "negate' provided but there is no config '" + prefix + "predicate' to be negated."); + String prefixedNegate = prefix + PredicatedTransformation.NEGATE_CONFIG; + String prefixedPredicate = prefix + PredicatedTransformation.PREDICATE_CONFIG; + if (props.containsKey(prefixedNegate) && + !props.containsKey(prefixedPredicate)) { + throw new ConfigException("Config '" + prefixedNegate + "' was provided " + + "but there is no config '" + prefixedPredicate + "' defining a predicate to be negated."); } } }.enrich(newDef); - new EnrichablePlugin>("predicate", PREDICATES_CONFIG, PREDICATES_GROUP, (Class) Predicate.class, props, requireFullConfig) { + new EnrichablePlugin>("Predicate", PREDICATES_CONFIG, PREDICATES_GROUP, + (Class) Predicate.class, props, requireFullConfig) { @Override protected Set>> plugins() { return (Set) plugins.predicates(); @@ -425,7 +431,7 @@ public void ensureValid(String name, Object value) { } }; newDef.define(typeConfig, Type.CLASS, ConfigDef.NO_DEFAULT_VALUE, typeValidator, Importance.HIGH, - "Class for the '" + alias + "' " + aliasKind + ".", group, orderInGroup++, Width.LONG, + "Class for the '" + alias + "' " + aliasKind.toLowerCase(Locale.ENGLISH) + ".", group, orderInGroup++, Width.LONG, baseClass.getSimpleName() + " type for " + alias, Collections.emptyList(), new ClassRecommender()); @@ -490,7 +496,6 @@ ConfigDef getConfigDefFromConfigProvidingClass(String key, Class cls) { .filter(c -> Modifier.isPublic(c.getModifiers())) .map(Class::getName) .collect(Collectors.joining(", ")); - String aliasKind = this.aliasKind.substring(0, 1).toUpperCase(Locale.ENGLISH) + this.aliasKind.substring(1); String message = childClassNames.trim().isEmpty() ? aliasKind + " is abstract and cannot be created." : aliasKind + " is abstract and cannot be created. Did you mean " + childClassNames + "?"; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java index 117608a1e1ec8..8977867215799 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java @@ -21,6 +21,7 @@ import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.utils.Utils; import org.apache.kafka.connect.connector.ConnectRecord; +import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.transforms.Transformation; import org.apache.kafka.connect.transforms.predicates.Predicate; @@ -31,9 +32,11 @@ */ class PredicatedTransformation> implements Transformation { - /*test*/ final Predicate predicate; - /*test*/ final Transformation delegate; - /*test*/ final boolean negate; + static final String PREDICATE_CONFIG = "predicate"; + static final String NEGATE_CONFIG = "negate"; + /*test*/ Predicate predicate; + /*test*/ Transformation delegate; + /*test*/ boolean negate; PredicatedTransformation(Predicate predicate, boolean negate, Transformation delegate) { this.predicate = predicate; @@ -43,7 +46,8 @@ class PredicatedTransformation> implements Transforma @Override public void configure(Map configs) { - + throw new ConnectException(PredicatedTransformation.class.getName() + ".configure() " + + "should never be called directly."); } @Override @@ -56,7 +60,8 @@ public R apply(R record) { @Override public ConfigDef config() { - return null; + throw new ConnectException(PredicatedTransformation.class.getName() + ".config() " + + "should never be called directly."); } @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 aeb3be0eedaba..8cdc34544a37d 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 @@ -363,7 +363,7 @@ public void negatedButNoPredicate() { new ConnectorConfig(MOCK_PLUGINS, props); } - static class TestPredicate> implements Predicate { + public static class TestPredicate> implements Predicate { int param; @@ -391,7 +391,7 @@ public void configure(Map configs) { } } - static abstract class AbstractTestPredicate> implements Predicate { + public static abstract class AbstractTestPredicate> implements Predicate { public AbstractTestPredicate() { } From f55e14616ee9698d24d0d1e6534cbe477333d7f0 Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Fri, 22 May 2020 09:51:04 +0100 Subject: [PATCH 07/13] Numerous fixes and an integration test. --- .../connect/runtime/ConnectorConfig.java | 1 - .../runtime/PredicatedTransformation.java | 9 + .../isolation/DelegatingClassLoader.java | 1 + .../connect/tools/TransformationDoc.java | 4 +- .../connect/integration/ConnectorHandle.java | 24 +- .../integration/MonitorableSinkConnector.java | 2 +- .../MonitorableSourceConnector.java | 5 +- .../kafka/connect/integration/TaskHandle.java | 20 +- .../TransformationIntegrationTest.java | 297 ++++++++++++++++++ .../kafka/connect/transforms/Filter.java | 7 +- .../transforms/predicates/HasHeaderKey.java | 7 + .../predicates/RecordIsTombstone.java | 5 + .../predicates/TopicNameMatches.java | 36 +-- 13 files changed, 383 insertions(+), 35 deletions(-) create mode 100644 connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.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 4468e6ba857a6..6dcc8be0f4761 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 @@ -315,7 +315,6 @@ protected Set>> plugins() { return (Set) plugins.transformations(); } - @Override protected ConfigDef initialConfigDef() { // All Transformations get these config parameters implicitly diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java index 8977867215799..5315da5510981 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java @@ -69,4 +69,13 @@ public void close() { Utils.closeQuietly(delegate, "predicated"); Utils.closeQuietly(predicate, "predicate"); } + + @Override + public String toString() { + return "PredicatedTransformation{" + + "predicate=" + predicate + + ", delegate=" + delegate + + ", negate=" + negate + + '}'; + } } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java index 23f22e346584b..86a53c75ced58 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/DelegatingClassLoader.java @@ -412,6 +412,7 @@ private void addAllAliases() { addAliases(converters); addAliases(headerConverters); addAliases(transformations); + addAliases(predicates); addAliases(restExtensions); addAliases(connectorClientConfigPolicies); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/tools/TransformationDoc.java b/connect/runtime/src/main/java/org/apache/kafka/connect/tools/TransformationDoc.java index 77d7728e2434b..3809d393455c9 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/tools/TransformationDoc.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/tools/TransformationDoc.java @@ -19,6 +19,7 @@ import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.connect.transforms.Cast; import org.apache.kafka.connect.transforms.ExtractField; +import org.apache.kafka.connect.transforms.Filter; import org.apache.kafka.connect.transforms.Flatten; import org.apache.kafka.connect.transforms.HoistField; import org.apache.kafka.connect.transforms.InsertField; @@ -60,7 +61,8 @@ private DocInfo(String transformationName, String overview, ConfigDef configDef) new DocInfo(RegexRouter.class.getName(), RegexRouter.OVERVIEW_DOC, RegexRouter.CONFIG_DEF), new DocInfo(Flatten.class.getName(), Flatten.OVERVIEW_DOC, Flatten.CONFIG_DEF), new DocInfo(Cast.class.getName(), Cast.OVERVIEW_DOC, Cast.CONFIG_DEF), - new DocInfo(TimestampConverter.class.getName(), TimestampConverter.OVERVIEW_DOC, TimestampConverter.CONFIG_DEF) + new DocInfo(TimestampConverter.class.getName(), TimestampConverter.OVERVIEW_DOC, TimestampConverter.CONFIG_DEF), + new DocInfo(Filter.class.getName(), Filter.OVERVIEW_DOC, Filter.CONFIG_DEF) ); private static void printTransformationHtml(PrintStream out, DocInfo docInfo) { diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java index a4a461288301b..70cd99f4ba41c 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java @@ -16,19 +16,21 @@ */ package org.apache.kafka.connect.integration; -import org.apache.kafka.connect.errors.DataException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import java.util.Collection; import java.util.List; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; +import java.util.function.Consumer; import java.util.stream.Collectors; import java.util.stream.IntStream; +import org.apache.kafka.connect.errors.DataException; +import org.apache.kafka.connect.sink.SinkRecord; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + /** * A handle to a connector executing in a Connect cluster. */ @@ -57,7 +59,19 @@ public ConnectorHandle(String connectorName) { * @return a non-null {@link TaskHandle} */ public TaskHandle taskHandle(String taskId) { - return taskHandles.computeIfAbsent(taskId, k -> new TaskHandle(this, taskId)); + return taskHandle(taskId, null); + } + + /** + * Get or create a task handle for a given task id. The task need not be created when this method is called. If the + * handle is called before the task is created, the task will bind to the handle once it starts (or restarts). + * + * @param taskId the task id + * @param consumer A callback invoked when a sink task processes a record. + * @return a non-null {@link TaskHandle} + */ + public TaskHandle taskHandle(String taskId, Consumer consumer) { + return taskHandles.computeIfAbsent(taskId, k -> new TaskHandle(this, taskId, consumer)); } /** diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java index 05b2dfdd8864d..f82556012cf43 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/MonitorableSinkConnector.java @@ -125,7 +125,7 @@ public void open(Collection partitions) { @Override public void put(Collection records) { for (SinkRecord rec : records) { - taskHandle.record(); + taskHandle.record(rec); TopicPartition tp = cachedTopicPartitions .computeIfAbsent(rec.topic(), v -> new HashMap<>()) .computeIfAbsent(rec.kafkaPartition(), v -> new TopicPartition(rec.topic(), rec.kafkaPartition())); 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 fbd763a928035..c11b9bd79d0fb 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 @@ -20,6 +20,7 @@ import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.connect.connector.Task; import org.apache.kafka.connect.data.Schema; +import org.apache.kafka.connect.header.ConnectHeaders; import org.apache.kafka.connect.runtime.TestSourceConnector; import org.apache.kafka.connect.source.SourceRecord; import org.apache.kafka.connect.source.SourceTask; @@ -138,7 +139,9 @@ public List poll() { Schema.STRING_SCHEMA, "key-" + taskId + "-" + seqno, Schema.STRING_SCHEMA, - "value-" + taskId + "-" + seqno)) + "value-" + taskId + "-" + seqno, + null, + new ConnectHeaders().addLong("header-" + seqno, seqno))) .collect(Collectors.toList()); } return null; diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java index 1159cb8fcc714..f04c139b132bd 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java @@ -16,15 +16,17 @@ */ package org.apache.kafka.connect.integration; -import org.apache.kafka.connect.errors.DataException; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; +import java.util.function.Consumer; import java.util.stream.IntStream; +import org.apache.kafka.connect.errors.DataException; +import org.apache.kafka.connect.sink.SinkRecord; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + /** * A handle to an executing task in a worker. Use this class to record progress, for example: number of records seen * by the task using so far, or waiting for partitions to be assigned to the task. @@ -37,22 +39,26 @@ public class TaskHandle { private final ConnectorHandle connectorHandle; private final AtomicInteger partitionsAssigned = new AtomicInteger(0); private final StartAndStopCounter startAndStopCounter = new StartAndStopCounter(); + private final Consumer consumer; private CountDownLatch recordsRemainingLatch; private CountDownLatch recordsToCommitLatch; private int expectedRecords = -1; private int expectedCommits = -1; - public TaskHandle(ConnectorHandle connectorHandle, String taskId) { - log.info("Created task {} for connector {}", taskId, connectorHandle); + public TaskHandle(ConnectorHandle connectorHandle, String taskId, Consumer consumer) { this.taskId = taskId; this.connectorHandle = connectorHandle; + this.consumer = consumer; } /** * Record a message arrival at the task and the connector overall. */ - public void record() { + public void record(SinkRecord rec) { + if (consumer != null) { + consumer.accept(rec); + } if (recordsRemainingLatch != null) { recordsRemainingLatch.countDown(); } diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java new file mode 100644 index 0000000000000..dabf296a32a02 --- /dev/null +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java @@ -0,0 +1,297 @@ +/* + * 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 java.util.HashMap; +import java.util.Map; +import java.util.Properties; +import java.util.concurrent.TimeUnit; + +import org.apache.kafka.clients.consumer.ConsumerRecord; +import org.apache.kafka.connect.storage.StringConverter; +import org.apache.kafka.connect.transforms.Filter; +import org.apache.kafka.connect.transforms.predicates.HasHeaderKey; +import org.apache.kafka.connect.transforms.predicates.RecordIsTombstone; +import org.apache.kafka.connect.transforms.predicates.TopicNameMatches; +import org.apache.kafka.connect.util.clusters.EmbeddedConnectCluster; +import org.apache.kafka.test.IntegrationTest; +import org.junit.After; +import org.junit.Before; +import org.junit.Test; +import org.junit.experimental.categories.Category; + +import static java.util.Collections.singletonMap; +import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLASS_CONFIG; +import static org.apache.kafka.connect.runtime.ConnectorConfig.KEY_CONVERTER_CLASS_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.TOPICS_CONFIG; +import static org.apache.kafka.connect.runtime.WorkerConfig.OFFSET_COMMIT_INTERVAL_MS_CONFIG; +import static org.apache.kafka.test.TestUtils.waitForCondition; +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertNotNull; + +/** + * An integration test for connectors with transformations + */ +@Category(IntegrationTest.class) +public class TransformationIntegrationTest { + + private static final int NUM_RECORDS_PRODUCED = 2000; + private static final int NUM_TOPIC_PARTITIONS = 3; + private static final long RECORD_TRANSFER_DURATION_MS = TimeUnit.SECONDS.toMillis(30); + private static final long OBSERVED_RECORDS_DURATION_MS = TimeUnit.SECONDS.toMillis(60); + private static final int NUM_TASKS = 3; + private static final int NUM_WORKERS = 3; + private static final String CONNECTOR_NAME = "simple-conn"; + private static final String SINK_CONNECTOR_CLASS_NAME = MonitorableSinkConnector.class.getSimpleName(); + private static final String SOURCE_CONNECTOR_CLASS_NAME = MonitorableSourceConnector.class.getSimpleName(); + + private EmbeddedConnectCluster connect; + private ConnectorHandle connectorHandle; + + @Before + public void setup() { + // setup Connect worker properties + Map exampleWorkerProps = new HashMap<>(); + exampleWorkerProps.put(OFFSET_COMMIT_INTERVAL_MS_CONFIG, String.valueOf(5_000)); + + // setup Kafka broker properties + Properties exampleBrokerProps = new Properties(); + exampleBrokerProps.put("auto.create.topics.enable", "false"); + + // build a Connect cluster backed by Kafka and Zk + connect = new EmbeddedConnectCluster.Builder() + .name("connect-cluster") + .numWorkers(NUM_WORKERS) + .numBrokers(1) + .workerProps(exampleWorkerProps) + .brokerProps(exampleBrokerProps) + .build(); + + // start the clusters + connect.start(); + + // get a handle to the connector + connectorHandle = RuntimeHandles.get().connectorHandle(CONNECTOR_NAME); + } + + @After + public void close() { + // delete connector handle + RuntimeHandles.get().deleteConnector(CONNECTOR_NAME); + + // stop all Connect, Kafka and Zk threads. + connect.stop(); + } + + /** + * Test the {@link Filter} transformer with a + * {@link TopicNameMatches} predicate on a sink connector. + */ + @Test + public void testFilterOnTopicNameWithSinkConnector() throws Exception { + Map observedRecords = observeRecords(); + + // create test topics + String fooTopic = "foo-topic"; + String barTopic = "bar-topic"; + int numFooRecords = NUM_RECORDS_PRODUCED; + int numBarRecords = NUM_RECORDS_PRODUCED; + connect.kafka().createTopic(fooTopic, NUM_TOPIC_PARTITIONS); + connect.kafka().createTopic(barTopic, NUM_TOPIC_PARTITIONS); + + // setup up props for the sink connector + Map props = new HashMap<>(); + props.put("name", CONNECTOR_NAME); + props.put(CONNECTOR_CLASS_CONFIG, SINK_CONNECTOR_CLASS_NAME); + props.put(TASKS_MAX_CONFIG, String.valueOf(NUM_TASKS)); + props.put(TOPICS_CONFIG, String.join(",", fooTopic, barTopic)); + props.put(KEY_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(VALUE_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(TRANSFORMS_CONFIG, "filter"); + props.put(TRANSFORMS_CONFIG + ".filter.type", Filter.class.getSimpleName()); + props.put(TRANSFORMS_CONFIG + ".filter.predicate", "barPredicate"); + props.put(PREDICATES_CONFIG, "barPredicate"); + props.put(PREDICATES_CONFIG + ".barPredicate.type", TopicNameMatches.class.getSimpleName()); + props.put(PREDICATES_CONFIG + ".barPredicate.pattern", "bar-.*"); + + // expect all records to be consumed by the connector + connectorHandle.expectedRecords(numFooRecords); + + // expect all records to be consumed by the connector + connectorHandle.expectedCommits(numFooRecords); + + // start a sink connector + connect.configureConnector(CONNECTOR_NAME, props); + + // produce some messages into source topic partitions + for (int i = 0; i < numBarRecords; i++) { + connect.kafka().produce(barTopic, i % NUM_TOPIC_PARTITIONS, "key", "simple-message-value-" + i); + } + for (int i = 0; i < numFooRecords; i++) { + connect.kafka().produce(fooTopic, i % NUM_TOPIC_PARTITIONS, "key", "simple-message-value-" + i); + } + + // consume all records from the source topic or fail, to ensure that they were correctly produced. + assertEquals("Unexpected number of records consumed", numFooRecords, + connect.kafka().consume(numFooRecords, RECORD_TRANSFER_DURATION_MS, fooTopic).count()); + assertEquals("Unexpected number of records consumed", numFooRecords, + connect.kafka().consume(numBarRecords, RECORD_TRANSFER_DURATION_MS, barTopic).count()); + + // wait for the connector tasks to consume all records. + connectorHandle.awaitRecords(RECORD_TRANSFER_DURATION_MS); + + // wait for the connector tasks to commit all records. + connectorHandle.awaitCommits(RECORD_TRANSFER_DURATION_MS); + + // Assert that we didn't see any baz + Map expectedRecordCounts = singletonMap(fooTopic, Long.valueOf(numFooRecords)); + assertObservedRecords(observedRecords, expectedRecordCounts); + + // delete connector + connect.deleteConnector(CONNECTOR_NAME); + } + + private void assertObservedRecords(Map observedRecords, Map expectedRecordCounts) throws InterruptedException { + waitForCondition(() -> expectedRecordCounts.equals(observedRecords), + OBSERVED_RECORDS_DURATION_MS, + () -> "The observed records should be " + expectedRecordCounts + " but was " + observedRecords); + } + + private Map observeRecords() { + Map observedRecords = new HashMap<>(); + // record all the record we see + connectorHandle.taskHandle(CONNECTOR_NAME + "-0", + record -> observedRecords.compute(record.topic(), + (key, value) -> value == null ? 1 : value + 1)); + return observedRecords; + } + + /** + * Test the {@link Filter} transformer with a + * {@link RecordIsTombstone} predicate on a sink connector. + */ + @Test + public void testFilterOnTombstonesWithSinkConnector() throws Exception { + Map observedRecords = observeRecords(); + + // create test topics + String topic = "foo-topic"; + int numRecords = NUM_RECORDS_PRODUCED; + connect.kafka().createTopic(topic, NUM_TOPIC_PARTITIONS); + + // setup up props for the sink connector + Map props = new HashMap<>(); + props.put("name", CONNECTOR_NAME); + props.put(CONNECTOR_CLASS_CONFIG, SINK_CONNECTOR_CLASS_NAME); + props.put(TASKS_MAX_CONFIG, String.valueOf(NUM_TASKS)); + props.put(TOPICS_CONFIG, String.join(",", topic)); + props.put(KEY_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(VALUE_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(TRANSFORMS_CONFIG, "filter"); + props.put(TRANSFORMS_CONFIG + ".filter.type", Filter.class.getSimpleName()); + props.put(TRANSFORMS_CONFIG + ".filter.predicate", "barPredicate"); + props.put(PREDICATES_CONFIG, "barPredicate"); + props.put(PREDICATES_CONFIG + ".barPredicate.type", RecordIsTombstone.class.getSimpleName()); + + // expect only half the records to be consumed by the connector + connectorHandle.expectedCommits(numRecords); + connectorHandle.expectedRecords(numRecords / 2); + + // start a sink connector + connect.configureConnector(CONNECTOR_NAME, props); + + // produce some messages into source topic partitions + for (int i = 0; i < numRecords; i++) { + connect.kafka().produce(topic, i % NUM_TOPIC_PARTITIONS, "key", i % 2 == 0 ? "simple-message-value-" + i : null); + } + + // consume all records from the source topic or fail, to ensure that they were correctly produced. + assertEquals("Unexpected number of records consumed", numRecords, + connect.kafka().consume(numRecords, RECORD_TRANSFER_DURATION_MS, topic).count()); + + // wait for the connector tasks to consume all records. + connectorHandle.awaitRecords(RECORD_TRANSFER_DURATION_MS); + + // wait for the connector tasks to commit all records. + connectorHandle.awaitCommits(RECORD_TRANSFER_DURATION_MS); + + Map expectedRecordCounts = singletonMap(topic, Long.valueOf(numRecords / 2)); + assertObservedRecords(observedRecords, expectedRecordCounts); + + // delete connector + connect.deleteConnector(CONNECTOR_NAME); + } + + /** + * Test the {@link Filter} transformer with a + * {@link HasHeaderKey} predicate on a source connector. + */ + @Test + public void testFilterOnHasHeaderKeyWithSourceConnector() throws Exception { + // create test topic + connect.kafka().createTopic("test-topic", NUM_TOPIC_PARTITIONS); + + // setup up props for the sink connector + Map props = new HashMap<>(); + props.put("name", CONNECTOR_NAME); + props.put(CONNECTOR_CLASS_CONFIG, SOURCE_CONNECTOR_CLASS_NAME); + props.put(TASKS_MAX_CONFIG, String.valueOf(NUM_TASKS)); + props.put("topic", "test-topic"); + props.put("throughput", String.valueOf(500)); + props.put(KEY_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(VALUE_CONVERTER_CLASS_CONFIG, StringConverter.class.getName()); + props.put(TRANSFORMS_CONFIG, "filter"); + props.put(TRANSFORMS_CONFIG + ".filter.type", Filter.class.getSimpleName()); + props.put(TRANSFORMS_CONFIG + ".filter.predicate", "headerPredicate"); + props.put(TRANSFORMS_CONFIG + ".filter.negate", "true"); + props.put(PREDICATES_CONFIG, "headerPredicate"); + props.put(PREDICATES_CONFIG + ".headerPredicate.type", HasHeaderKey.class.getSimpleName()); + props.put(PREDICATES_CONFIG + ".headerPredicate.name", "header-8"); + + // expect all records to be produced by the connector + connectorHandle.expectedRecords(NUM_RECORDS_PRODUCED); + + // expect all records to be produced by the connector + connectorHandle.expectedCommits(NUM_RECORDS_PRODUCED); + + // validate the intended connector configuration, a valid config + connect.assertions().assertExactlyNumErrorsOnConnectorConfigValidation(SOURCE_CONNECTOR_CLASS_NAME, props, 0, + "Validating connector configuration produced an unexpected number or errors."); + + // start a source connector + connect.configureConnector(CONNECTOR_NAME, props); + + // wait for the connector tasks to produce enough records + connectorHandle.awaitRecords(RECORD_TRANSFER_DURATION_MS); + + // wait for the connector tasks to commit enough records + connectorHandle.awaitCommits(RECORD_TRANSFER_DURATION_MS); + + // consume all records from the source topic or fail, to ensure that they were correctly produced + for (ConsumerRecord record : connect.kafka().consume(1, RECORD_TRANSFER_DURATION_MS, "test-topic")) { + assertNotNull("Expected header to exist", + record.headers().lastHeader("header-8")); + } + + // delete connector + connect.deleteConnector(CONNECTOR_NAME); + } +} diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java index c805c84bf3e81..d7fb54eaca9db 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/Filter.java @@ -29,6 +29,11 @@ */ public class Filter> implements Transformation { + public static final String OVERVIEW_DOC = "Drops all records, filtering them from subsequent transformations in the chain. " + + "This is intended to be used conditionally to filter out records matching (or not matching) " + + "a particular Predicate."; + public static final ConfigDef CONFIG_DEF = new ConfigDef(); + @Override public R apply(R record) { return null; @@ -36,7 +41,7 @@ public R apply(R record) { @Override public ConfigDef config() { - return new ConfigDef(); + return CONFIG_DEF; } @Override diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java index 6501b401ef7a8..4b541b24b7bbc 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java @@ -56,4 +56,11 @@ public void close() { public void configure(Map configs) { this.name = new SimpleConfig(config(), configs).getString(NAME_CONFIG); } + + @Override + public String toString() { + return "HasHeaderKey{" + + "name='" + name + '\'' + + '}'; + } } diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java index 7d5c47e66c1bc..e39591c02f727 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/RecordIsTombstone.java @@ -48,4 +48,9 @@ public void close() { public void configure(Map configs) { } + + @Override + public String toString() { + return "RecordIsTombstone{}"; + } } diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java index 2962a397c5125..439da3b1b52aa 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java @@ -23,6 +23,7 @@ import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.connect.connector.ConnectRecord; +import org.apache.kafka.connect.transforms.util.RegexValidator; import org.apache.kafka.connect.transforms.util.SimpleConfig; /** @@ -32,15 +33,8 @@ public class TopicNameMatches> implements Predicate { private static final String PATTERN_CONFIG = "pattern"; - private static final ConfigDef CONFIG_DEF = new ConfigDef().define(PATTERN_CONFIG, ConfigDef.Type.STRING, null, - new ConfigDef.Validator() { - @Override - public void ensureValid(String name, Object value) { - if (value instanceof String) { - compile(name, (String) value); - } - } - }, ConfigDef.Importance.MEDIUM, + private static final ConfigDef CONFIG_DEF = new ConfigDef().define(PATTERN_CONFIG, ConfigDef.Type.STRING, ".*", + new RegexValidator(), ConfigDef.Importance.MEDIUM, "A Java regular expression for matching against the name of a record's topic."); private Pattern pattern; @@ -49,14 +43,6 @@ public ConfigDef config() { return CONFIG_DEF; } - private static Pattern compile(String name, String value) { - try { - return Pattern.compile(value); - } catch (PatternSyntaxException e) { - throw new ConfigException(name, value, "entry must be a Java-compatible regular expression: " + e.getMessage()); - } - } - @Override public boolean test(R record) { return record.topic() != null && pattern.matcher(record.topic()).matches(); @@ -70,6 +56,20 @@ public void close() { @Override public void configure(Map configs) { SimpleConfig simpleConfig = new SimpleConfig(config(), configs); - this.pattern = compile(PATTERN_CONFIG, simpleConfig.getString(PATTERN_CONFIG)); + Pattern result; + String value = simpleConfig.getString(PATTERN_CONFIG); + try { + result = Pattern.compile(value); + } catch (PatternSyntaxException e) { + throw new ConfigException(PATTERN_CONFIG, value, "entry must be a Java-compatible regular expression: " + e.getMessage()); + } + this.pattern = result; + } + + @Override + public String toString() { + return "TopicNameMatches{" + + "pattern=" + pattern + + '}'; } } From 49b362cd4144fbe2dc10114d3c1a4095c4c30bdb Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Wed, 27 May 2020 11:06:28 +0100 Subject: [PATCH 08/13] Review comments --- .../kafka/connect/runtime/ConnectorConfig.java | 12 ++++++------ .../connect/runtime/PredicatedTransformation.java | 6 +++--- .../connect/runtime/isolation/PluginUtils.java | 2 +- .../kafka/connect/integration/ConnectorHandle.java | 10 +++++----- .../kafka/connect/integration/TaskHandle.java | 10 +++++----- .../kafka/connect/runtime/ConnectorConfigTest.java | 14 ++++++-------- .../connect/runtime/isolation/PluginUtilsTest.java | 7 +++++++ 7 files changed, 33 insertions(+), 28 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 6dcc8be0f4761..37a710cbfd33e 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 @@ -159,6 +159,7 @@ public class ConnectorConfig extends AbstractConfig { public static final String CONNECTOR_CLIENT_PRODUCER_OVERRIDES_PREFIX = "producer.override."; public static final String CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX = "consumer.override."; public static final String CONNECTOR_CLIENT_ADMIN_OVERRIDES_PREFIX = "admin.override."; + public static final String PREDICATES_PREFIX = "predicates."; private final EnrichedConnectorConfig enrichedConfig; private static class EnrichedConnectorConfig extends AbstractConfig { @@ -283,7 +284,7 @@ public > List> transformations() { Object negate = configs.remove(PredicatedTransformation.NEGATE_CONFIG); transformation.configure(configs); if (predicateAlias != null) { - String predicatePrefix = "predicates." + predicateAlias + "."; + String predicatePrefix = PREDICATES_PREFIX + predicateAlias + "."; @SuppressWarnings("unchecked") Predicate predicate = Utils.newInstance(getClass(predicatePrefix + "type"), Predicate.class); predicate.configure(originalsWithPrefix(predicatePrefix)); @@ -422,13 +423,12 @@ void enrich(ConfigDef newDef) { int orderInGroup = 0; final String typeConfig = prefix + "type"; - final ConfigDef.Validator typeValidator = new ConfigDef.Validator() { - @Override - public void ensureValid(String name, Object value) { + final ConfigDef.Validator typeValidator = ConfigDef.LambdaValidator.with( + (String name, Object value) -> { validateProps(prefix); getConfigDefFromConfigProvidingClass(typeConfig, (Class) value); - } - }; + }, + () -> "valid configs for " + alias + " " + aliasKind.toLowerCase(Locale.ENGLISH)); newDef.define(typeConfig, Type.CLASS, ConfigDef.NO_DEFAULT_VALUE, typeValidator, Importance.HIGH, "Class for the '" + alias + "' " + aliasKind.toLowerCase(Locale.ENGLISH) + ".", group, orderInGroup++, Width.LONG, baseClass.getSimpleName() + " type for " + alias, diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java index 5315da5510981..d61772f9f890c 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/PredicatedTransformation.java @@ -34,9 +34,9 @@ class PredicatedTransformation> implements Transforma static final String PREDICATE_CONFIG = "predicate"; static final String NEGATE_CONFIG = "negate"; - /*test*/ Predicate predicate; - /*test*/ Transformation delegate; - /*test*/ boolean negate; + Predicate predicate; + Transformation delegate; + boolean negate; PredicatedTransformation(Predicate predicate, boolean negate, Transformation delegate) { this.predicate = predicate; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginUtils.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginUtils.java index 36feac50d1824..5edd7657c072f 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginUtils.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/isolation/PluginUtils.java @@ -128,7 +128,7 @@ public class PluginUtils { // added to the WHITELIST), then this base interface or class needs to be excluded in the // regular expression pattern private static final Pattern WHITELIST = Pattern.compile("^org\\.apache\\.kafka\\.(?:connect\\.(?:" - + "transforms\\.(?!Transformation$).*" + + "transforms\\.(?!Transformation|predicates\\.Predicate$).*" + "|json\\..*" + "|file\\..*" + "|mirror\\..*" diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java index 70cd99f4ba41c..06bc37352e455 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/ConnectorHandle.java @@ -16,6 +16,11 @@ */ package org.apache.kafka.connect.integration; +import org.apache.kafka.connect.errors.DataException; +import org.apache.kafka.connect.sink.SinkRecord; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.util.Collection; import java.util.List; import java.util.Map; @@ -26,11 +31,6 @@ import java.util.stream.Collectors; import java.util.stream.IntStream; -import org.apache.kafka.connect.errors.DataException; -import org.apache.kafka.connect.sink.SinkRecord; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - /** * A handle to a connector executing in a Connect cluster. */ diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java index f04c139b132bd..a13c3304f3f3d 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java @@ -16,17 +16,17 @@ */ package org.apache.kafka.connect.integration; +import org.apache.kafka.connect.errors.DataException; +import org.apache.kafka.connect.sink.SinkRecord; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; + import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicInteger; import java.util.function.Consumer; import java.util.stream.IntStream; -import org.apache.kafka.connect.errors.DataException; -import org.apache.kafka.connect.sink.SinkRecord; -import org.slf4j.Logger; -import org.slf4j.LoggerFactory; - /** * A handle to an executing task in a worker. Use this class to record progress, for example: number of records seen * by the task using so far, or waiting for partitions to be assigned to the task. 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 8cdc34544a37d..5d78a00a9dfe4 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 @@ -16,12 +16,6 @@ */ package org.apache.kafka.connect.runtime; -import java.util.Collections; -import java.util.HashMap; -import java.util.List; -import java.util.Map; -import java.util.Set; - import org.apache.kafka.common.config.ConfigDef; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.connect.connector.ConnectRecord; @@ -33,6 +27,12 @@ import org.junit.Ignore; import org.junit.Test; +import java.util.Collections; +import java.util.HashMap; +import java.util.List; +import java.util.Map; +import java.util.Set; + import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -216,8 +216,6 @@ public void abstractKeyValueTransform() { } } - - @Test(expected = ConfigException.class) public void wrongPredicateType() { Map props = new HashMap<>(); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginUtilsTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginUtilsTest.java index c406ead57e66f..1be3e42bab4fe 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginUtilsTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/isolation/PluginUtilsTest.java @@ -99,6 +99,9 @@ public void testConnectFrameworkClasses() { assertFalse(PluginUtils.shouldLoadInIsolation( "org.apache.kafka.connect.transforms.Transformation") ); + assertFalse(PluginUtils.shouldLoadInIsolation( + "org.apache.kafka.connect.transforms.predicates.Predicate") + ); assertFalse(PluginUtils.shouldLoadInIsolation( "org.apache.kafka.connect.storage.Converter") ); @@ -128,6 +131,10 @@ public void testAllowedConnectFrameworkClasses() { assertTrue(PluginUtils.shouldLoadInIsolation( "org.apache.kafka.connect.transforms.ExtractField$Key") ); + assertTrue(PluginUtils.shouldLoadInIsolation("org.apache.kafka.connect.transforms.predicates.")); + assertTrue(PluginUtils.shouldLoadInIsolation( + "org.apache.kafka.connect.transforms.predicates.TopicNameMatches") + ); assertTrue(PluginUtils.shouldLoadInIsolation("org.apache.kafka.connect.json.")); assertTrue(PluginUtils.shouldLoadInIsolation( "org.apache.kafka.connect.json.JsonConverter") From e76a0b4d4601310bd080e812668466997ac61268 Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Wed, 27 May 2020 11:08:18 +0100 Subject: [PATCH 09/13] Review comments --- .../apache/kafka/connect/integration/TaskHandle.java | 10 +++++++--- 1 file changed, 7 insertions(+), 3 deletions(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java index a13c3304f3f3d..b63c08bf36811 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TaskHandle.java @@ -52,12 +52,16 @@ public TaskHandle(ConnectorHandle connectorHandle, String taskId, Consumer Date: Wed, 27 May 2020 18:09:00 +0100 Subject: [PATCH 10/13] Apply suggestions from code review Co-authored-by: Randall Hauch --- .../kafka/connect/transforms/predicates/Predicate.java | 7 ++++++- .../kafka/connect/transforms/predicates/HasHeaderKey.java | 2 +- 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java b/connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java index d50efd3752b00..cc38fc049b5aa 100644 --- a/connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java +++ b/connect/api/src/main/java/org/apache/kafka/connect/transforms/predicates/Predicate.java @@ -35,14 +35,19 @@ public interface Predicate> extends Configurable, Aut /** * Configuration specification for this predicate. + * + * @return the configuration definition for this predicate; never null */ ConfigDef config(); /** * Returns whether the given record satisfies this predicate. + * + * @param record the record to evaluate; may not be null + * @return true if the predicate matches, or false otherwise */ boolean test(R record); @Override void close(); -} \ No newline at end of file +} diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java index 4b541b24b7bbc..ee519aaf2d612 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java @@ -31,7 +31,7 @@ public class HasHeaderKey> implements Predicate { private static final String NAME_CONFIG = "name"; - private static final ConfigDef CONFIG_DEF = new ConfigDef().define(NAME_CONFIG, ConfigDef.Type.STRING, null, + private static final ConfigDef CONFIG_DEF = new ConfigDef().define(NAME_CONFIG, ConfigDef.Type.STRING, ConfigDef.NO_DEFAULT_VALUE, new ConfigDef.NonEmptyString(), ConfigDef.Importance.MEDIUM, "The header name."); private String name; From 03ce156772cc890cf0abf6529e960d0c5809dfbe Mon Sep 17 00:00:00 2001 From: Randall Hauch Date: Wed, 27 May 2020 17:39:53 -0500 Subject: [PATCH 11/13] KAFKA-9673: Fix a few unit and integration tests --- .../TransformationIntegrationTest.java | 26 ++++++++++++++++++- .../predicates/TopicNameMatchesTest.java | 9 +++---- 2 files changed, 28 insertions(+), 7 deletions(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java index dabf296a32a02..ad9c02db17150 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java @@ -33,6 +33,8 @@ import org.junit.Before; import org.junit.Test; import org.junit.experimental.categories.Category; +import org.slf4j.Logger; +import org.slf4j.LoggerFactory; import static java.util.Collections.singletonMap; import static org.apache.kafka.connect.runtime.ConnectorConfig.CONNECTOR_CLASS_CONFIG; @@ -53,11 +55,13 @@ @Category(IntegrationTest.class) public class TransformationIntegrationTest { + private static final Logger log = LoggerFactory.getLogger(TransformationIntegrationTest.class); + private static final int NUM_RECORDS_PRODUCED = 2000; private static final int NUM_TOPIC_PARTITIONS = 3; private static final long RECORD_TRANSFER_DURATION_MS = TimeUnit.SECONDS.toMillis(30); private static final long OBSERVED_RECORDS_DURATION_MS = TimeUnit.SECONDS.toMillis(60); - private static final int NUM_TASKS = 3; + private static final int NUM_TASKS = 1; private static final int NUM_WORKERS = 3; private static final String CONNECTOR_NAME = "simple-conn"; private static final String SINK_CONNECTOR_CLASS_NAME = MonitorableSinkConnector.class.getSimpleName(); @@ -107,6 +111,8 @@ public void close() { */ @Test public void testFilterOnTopicNameWithSinkConnector() throws Exception { + assertConnectReady(); + Map observedRecords = observeRecords(); // create test topics @@ -140,6 +146,7 @@ public void testFilterOnTopicNameWithSinkConnector() throws Exception { // start a sink connector connect.configureConnector(CONNECTOR_NAME, props); + assertConnectorRunning(); // produce some messages into source topic partitions for (int i = 0; i < numBarRecords; i++) { @@ -169,6 +176,17 @@ public void testFilterOnTopicNameWithSinkConnector() throws Exception { connect.deleteConnector(CONNECTOR_NAME); } + private void assertConnectReady() throws InterruptedException { + connect.assertions().assertExactlyNumBrokersAreUp(1, "Brokers did not start in time."); + connect.assertions().assertExactlyNumWorkersAreUp(NUM_WORKERS, "Worker did not start in time."); + log.info("Completed startup of {} Kafka brokers and {} Connect workers", 1, NUM_WORKERS); + } + + private void assertConnectorRunning() throws InterruptedException { + connect.assertions().assertConnectorAndAtLeastNumTasksAreRunning(CONNECTOR_NAME, NUM_TASKS, + "Connector tasks did not start in time."); + } + private void assertObservedRecords(Map observedRecords, Map expectedRecordCounts) throws InterruptedException { waitForCondition(() -> expectedRecordCounts.equals(observedRecords), OBSERVED_RECORDS_DURATION_MS, @@ -190,6 +208,8 @@ record -> observedRecords.compute(record.topic(), */ @Test public void testFilterOnTombstonesWithSinkConnector() throws Exception { + assertConnectReady(); + Map observedRecords = observeRecords(); // create test topics @@ -217,6 +237,7 @@ public void testFilterOnTombstonesWithSinkConnector() throws Exception { // start a sink connector connect.configureConnector(CONNECTOR_NAME, props); + assertConnectorRunning(); // produce some messages into source topic partitions for (int i = 0; i < numRecords; i++) { @@ -246,6 +267,8 @@ public void testFilterOnTombstonesWithSinkConnector() throws Exception { */ @Test public void testFilterOnHasHeaderKeyWithSourceConnector() throws Exception { + assertConnectReady(); + // create test topic connect.kafka().createTopic("test-topic", NUM_TOPIC_PARTITIONS); @@ -278,6 +301,7 @@ public void testFilterOnHasHeaderKeyWithSourceConnector() throws Exception { // start a source connector connect.configureConnector(CONNECTOR_NAME, props); + assertConnectorRunning(); // wait for the connector tasks to produce enough records connectorHandle.awaitRecords(RECORD_TRANSFER_DURATION_MS); diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java index 500f2b5cd6e31..6b2fcb9d98afa 100644 --- a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java @@ -36,12 +36,9 @@ public void testConfig() { predicate.config().validate(Collections.singletonMap("pattern", "my-prefix-.*")); List configs = predicate.config().validate(Collections.singletonMap("pattern", "*")); - assertEquals(singletonList("Invalid value * for configuration pattern: " + - "entry must be a Java-compatible regular expression: " + - "Dangling meta character '*' near index 0" + System.lineSeparator() + - "*" + System.lineSeparator() + - "^"), - configs.get(0).errorMessages()); + List errorMsgs = configs.get(0).errorMessages(); + assertEquals(1, errorMsgs.size()); + assertTrue(errorMsgs.get(0).contains("Invalid regex")); } @Test From 60cd39eb0e17cdf7e3fa52886391e774460123c7 Mon Sep 17 00:00:00 2001 From: Randall Hauch Date: Wed, 27 May 2020 18:28:06 -0500 Subject: [PATCH 12/13] KAFKA-9673: Added some tests, made others a bit more robust, and added more validation to TopicNameMatches --- .../connect/runtime/ConnectorConfigTest.java | 56 ++++++++++--------- .../transforms/predicates/HasHeaderKey.java | 3 +- .../predicates/TopicNameMatches.java | 6 +- .../predicates/HasHeaderKeyTest.java | 24 ++++++++ .../predicates/TopicNameMatchesTest.java | 33 ++++++++++- 5 files changed, 92 insertions(+), 30 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 5d78a00a9dfe4..c35663b761d84 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 @@ -24,7 +24,6 @@ import org.apache.kafka.connect.runtime.isolation.Plugins; import org.apache.kafka.connect.transforms.Transformation; import org.apache.kafka.connect.transforms.predicates.Predicate; -import org.junit.Ignore; import org.junit.Test; import java.util.Collections; @@ -34,6 +33,7 @@ import java.util.Set; import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; import static org.junit.Assert.fail; @@ -83,41 +83,45 @@ public void noTransforms() { new ConnectorConfig(MOCK_PLUGINS, props); } - @Test(expected = ConfigException.class) + @Test public void danglingTransformAlias() { Map props = new HashMap<>(); props.put("name", "test"); props.put("connector.class", TestConnector.class.getName()); props.put("transforms", "dangler"); - new ConnectorConfig(MOCK_PLUGINS, props); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("Not a Transformation")); } - @Test(expected = ConfigException.class) + @Test public void emptyConnectorName() { Map props = new HashMap<>(); props.put("name", ""); props.put("connector.class", TestConnector.class.getName()); - new ConnectorConfig(MOCK_PLUGINS, props); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("String may not be empty")); } - @Test(expected = ConfigException.class) + @Test public void wrongTransformationType() { Map props = new HashMap<>(); props.put("name", "test"); props.put("connector.class", TestConnector.class.getName()); props.put("transforms", "a"); props.put("transforms.a.type", "uninstantiable"); - new ConnectorConfig(MOCK_PLUGINS, props); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("Class uninstantiable could not be found")); } - @Test(expected = ConfigException.class) + @Test public void unconfiguredTransform() { Map props = new HashMap<>(); props.put("name", "test"); props.put("connector.class", TestConnector.class.getName()); props.put("transforms", "a"); props.put("transforms.a.type", SimpleTransformation.class.getName()); - new ConnectorConfig(MOCK_PLUGINS, props); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("Missing required configuration \"transforms.a.magic.number\" which")); } @Test @@ -128,12 +132,8 @@ public void misconfiguredTransform() { props.put("transforms", "a"); props.put("transforms.a.type", SimpleTransformation.class.getName()); props.put("transforms.a.magic.number", "40"); - try { - new ConnectorConfig(MOCK_PLUGINS, props); - fail(); - } catch (ConfigException e) { - assertTrue(e.getMessage().contains("Value must be at least 42")); - } + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("Value must be at least 42")); } @Test @@ -216,7 +216,7 @@ public void abstractKeyValueTransform() { } } - @Test(expected = ConfigException.class) + @Test public void wrongPredicateType() { Map props = new HashMap<>(); props.put("name", "test"); @@ -227,7 +227,8 @@ public void wrongPredicateType() { props.put("transforms.a.predicate", "my-pred"); props.put("predicates", "my-pred"); props.put("predicates.my-pred.type", TestConnector.class.getName()); - new ConnectorConfig(MOCK_PLUGINS, props); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("Not a Predicate")); } @Test @@ -261,7 +262,7 @@ public void predicateNegationDefaultsToFalse() { assertPredicatedTransform(props, false); } - @Test(expected = ConfigException.class) + @Test public void abstractPredicate() { Map props = new HashMap<>(); props.put("name", "test"); @@ -273,7 +274,8 @@ public void abstractPredicate() { props.put("predicates", "my-pred"); props.put("predicates.my-pred.type", AbstractTestPredicate.class.getName()); props.put("predicates.my-pred.int", "84"); - assertPredicatedTransform(props, false); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("Predicate is abstract and cannot be created")); } private void assertPredicatedTransform(Map props, boolean expectedNegated) { @@ -318,9 +320,8 @@ public void misconfiguredPredicate() { } } - @Ignore("Is this really an error. There's no actual need for the predicates config (unlike transforms where it defines the order).") - @Test(expected = ConfigException.class) - public void missingPredicates() { + @Test + public void missingPredicateAliasProperty() { Map props = new HashMap<>(); props.put("name", "test"); props.put("connector.class", TestConnector.class.getName()); @@ -328,13 +329,14 @@ public void missingPredicates() { props.put("transforms.a.type", SimpleTransformation.class.getName()); props.put("transforms.a.magic.number", "42"); props.put("transforms.a.predicate", "my-pred"); + // technically not needed //props.put("predicates", "my-pred"); props.put("predicates.my-pred.type", TestPredicate.class.getName()); props.put("predicates.my-pred.int", "84"); new ConnectorConfig(MOCK_PLUGINS, props); } - @Test(expected = ConfigException.class) + @Test public void missingPredicateConfig() { Map props = new HashMap<>(); props.put("name", "test"); @@ -346,10 +348,11 @@ public void missingPredicateConfig() { props.put("predicates", "my-pred"); //props.put("predicates.my-pred.type", TestPredicate.class.getName()); //props.put("predicates.my-pred.int", "84"); - new ConnectorConfig(MOCK_PLUGINS, props); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("Not a Predicate")); } - @Test(expected = ConfigException.class) + @Test public void negatedButNoPredicate() { Map props = new HashMap<>(); props.put("name", "test"); @@ -358,7 +361,8 @@ public void negatedButNoPredicate() { props.put("transforms.a.type", SimpleTransformation.class.getName()); props.put("transforms.a.magic.number", "42"); props.put("transforms.a.negate", "true"); - new ConnectorConfig(MOCK_PLUGINS, props); + ConfigException e = assertThrows(ConfigException.class, () -> new ConnectorConfig(MOCK_PLUGINS, props)); + assertTrue(e.getMessage().contains("there is no config 'transforms.a.predicate' defining a predicate to be negated")); } public static class TestPredicate> implements Predicate { diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java index ee519aaf2d612..03f324e33e7f3 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKey.java @@ -31,7 +31,8 @@ public class HasHeaderKey> implements Predicate { private static final String NAME_CONFIG = "name"; - private static final ConfigDef CONFIG_DEF = new ConfigDef().define(NAME_CONFIG, ConfigDef.Type.STRING, ConfigDef.NO_DEFAULT_VALUE, + public static final ConfigDef CONFIG_DEF = new ConfigDef() + .define(NAME_CONFIG, ConfigDef.Type.STRING, ConfigDef.NO_DEFAULT_VALUE, new ConfigDef.NonEmptyString(), ConfigDef.Importance.MEDIUM, "The header name."); private String name; diff --git a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java index 439da3b1b52aa..fba8d514fb696 100644 --- a/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java +++ b/connect/transforms/src/main/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatches.java @@ -33,8 +33,10 @@ public class TopicNameMatches> implements Predicate { private static final String PATTERN_CONFIG = "pattern"; - private static final ConfigDef CONFIG_DEF = new ConfigDef().define(PATTERN_CONFIG, ConfigDef.Type.STRING, ".*", - new RegexValidator(), ConfigDef.Importance.MEDIUM, + public static final ConfigDef CONFIG_DEF = new ConfigDef() + .define(PATTERN_CONFIG, ConfigDef.Type.STRING, ConfigDef.NO_DEFAULT_VALUE, + ConfigDef.CompositeValidator.of(new ConfigDef.NonEmptyString(), new RegexValidator()), + ConfigDef.Importance.MEDIUM, "A Java regular expression for matching against the name of a record's topic."); private Pattern pattern; diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java index b2ad5450e6948..f6d1d330c5320 100644 --- a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/HasHeaderKeyTest.java @@ -18,22 +18,42 @@ import java.util.Arrays; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; import java.util.stream.Collectors; +import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.config.ConfigValue; import org.apache.kafka.connect.data.Schema; import org.apache.kafka.connect.header.Header; import org.apache.kafka.connect.source.SourceRecord; +import org.apache.kafka.connect.transforms.util.SimpleConfig; import org.junit.Test; import static java.util.Collections.singletonList; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; public class HasHeaderKeyTest { + @Test + public void testNameRequiredInConfig() { + Map props = new HashMap<>(); + ConfigException e = assertThrows(ConfigException.class, () -> config(props)); + assertTrue(e.getMessage().contains("Missing required configuration \"name\"")); + } + + @Test + public void testNameMayNotBeEmptyInConfig() { + Map props = new HashMap<>(); + props.put("name", ""); + ConfigException e = assertThrows(ConfigException.class, () -> config(props)); + assertTrue(e.getMessage().contains("String must be non-empty")); + } + @Test public void testConfig() { HasHeaderKey predicate = new HasHeaderKey<>(); @@ -58,6 +78,10 @@ public void testTest() { } + private SimpleConfig config(Map props) { + return new SimpleConfig(new HasHeaderKey().config(), props); + } + private SourceRecord recordWithHeaders(String... headers) { return new SourceRecord(null, null, null, null, null, null, null, null, null, Arrays.stream(headers).map(header -> new TestHeader(header)).collect(Collectors.toList())); diff --git a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java index 6b2fcb9d98afa..b0cc34bc36213 100644 --- a/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java +++ b/connect/transforms/src/test/java/org/apache/kafka/connect/transforms/predicates/TopicNameMatchesTest.java @@ -17,19 +17,46 @@ package org.apache.kafka.connect.transforms.predicates; import java.util.Collections; +import java.util.HashMap; import java.util.List; +import java.util.Map; +import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.config.ConfigValue; import org.apache.kafka.connect.source.SourceRecord; +import org.apache.kafka.connect.transforms.util.SimpleConfig; import org.junit.Test; -import static java.util.Collections.singletonList; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertThrows; import static org.junit.Assert.assertTrue; public class TopicNameMatchesTest { + @Test + public void testPatternRequiredInConfig() { + Map props = new HashMap<>(); + ConfigException e = assertThrows(ConfigException.class, () -> config(props)); + assertTrue(e.getMessage().contains("Missing required configuration \"pattern\"")); + } + + @Test + public void testPatternMayNotBeEmptyInConfig() { + Map props = new HashMap<>(); + props.put("pattern", ""); + ConfigException e = assertThrows(ConfigException.class, () -> config(props)); + assertTrue(e.getMessage().contains("String must be non-empty")); + } + + @Test + public void testPatternIsValidRegexInConfig() { + Map props = new HashMap<>(); + props.put("pattern", "["); + ConfigException e = assertThrows(ConfigException.class, () -> config(props)); + assertTrue(e.getMessage().contains("Invalid regex")); + } + @Test public void testConfig() { TopicNameMatches predicate = new TopicNameMatches<>(); @@ -56,6 +83,10 @@ public void testTest() { } + private SimpleConfig config(Map props) { + return new SimpleConfig(TopicNameMatches.CONFIG_DEF, props); + } + private SourceRecord recordWithTopicName(String topicName) { return new SourceRecord(null, null, topicName, null, null); } From cbc89818ce2a9ac9dae9ca02288f9d71b1a851da Mon Sep 17 00:00:00 2001 From: Tom Bentley Date: Thu, 28 May 2020 09:23:37 +0100 Subject: [PATCH 13/13] Trivial correction --- .../connect/integration/TransformationIntegrationTest.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java index ad9c02db17150..ca9d55dfd10b7 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/integration/TransformationIntegrationTest.java @@ -159,7 +159,7 @@ public void testFilterOnTopicNameWithSinkConnector() throws Exception { // consume all records from the source topic or fail, to ensure that they were correctly produced. assertEquals("Unexpected number of records consumed", numFooRecords, connect.kafka().consume(numFooRecords, RECORD_TRANSFER_DURATION_MS, fooTopic).count()); - assertEquals("Unexpected number of records consumed", numFooRecords, + assertEquals("Unexpected number of records consumed", numBarRecords, connect.kafka().consume(numBarRecords, RECORD_TRANSFER_DURATION_MS, barTopic).count()); // wait for the connector tasks to consume all records.