From efd73838e667796aa6602165e3584268a43a76a2 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Sat, 2 Apr 2022 01:09:31 -0400 Subject: [PATCH 1/3] KAFKA-7509: Clean up incorrect warnings logged by Connect --- .../kafka/clients/CommonClientConfigs.java | 10 ++ .../clients/admin/AdminClientConfig.java | 3 + .../kafka/clients/admin/KafkaAdminClient.java | 4 +- .../clients/consumer/ConsumerConfig.java | 9 +- .../kafka/clients/consumer/KafkaConsumer.java | 5 +- .../kafka/clients/producer/KafkaProducer.java | 4 +- .../clients/producer/ProducerConfig.java | 9 +- .../kafka/common/config/AbstractConfig.java | 9 +- .../kafka/connect/cli/ConnectDistributed.java | 3 +- .../kafka/connect/cli/ConnectStandalone.java | 2 +- .../kafka/connect/runtime/AbstractHerder.java | 6 +- .../kafka/connect/runtime/ConnectMetrics.java | 4 +- .../runtime/SourceTaskOffsetCommitter.java | 14 +-- .../apache/kafka/connect/runtime/Worker.java | 101 ++++++++++-------- .../kafka/connect/runtime/WorkerConfig.java | 32 +++++- .../distributed/DistributedConfig.java | 8 ++ .../distributed/DistributedHerder.java | 4 +- .../distributed/WorkerGroupMember.java | 5 +- .../runtime/standalone/StandaloneHerder.java | 11 +- .../storage/KafkaConfigBackingStore.java | 7 +- .../storage/KafkaOffsetBackingStore.java | 7 +- .../storage/KafkaStatusBackingStore.java | 7 +- .../kafka/connect/util/ConnectUtils.java | 4 +- .../connect/runtime/AbstractHerderTest.java | 49 +++++---- .../SourceTaskOffsetCommitterTest.java | 18 +--- .../kafka/connect/runtime/WorkerTest.java | 4 +- .../distributed/WorkerGroupMemberTest.java | 5 +- .../standalone/StandaloneHerderTest.java | 3 +- .../kafka/connect/util/ConnectUtilsTest.java | 9 +- 29 files changed, 222 insertions(+), 134 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java b/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java index 5371a73ece192..b470a92fc31ec 100644 --- a/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java +++ b/clients/src/main/java/org/apache/kafka/clients/CommonClientConfigs.java @@ -24,6 +24,7 @@ import java.util.HashMap; import java.util.Map; +import java.util.stream.Stream; /** * Configurations shared by Kafka client applications: producer, consumer, connect, etc. @@ -182,6 +183,15 @@ public class CommonClientConfigs { public static final String DEFAULT_API_TIMEOUT_MS_DOC = "Specifies the timeout (in milliseconds) for client APIs. " + "This configuration is used as the default timeout for all client operations that do not specify a timeout parameter."; + public static final String CONNECT_KAFKA_CLUSTER_ID = "connect.kafka.cluster.id"; + public static final String CONNECT_GROUP_ID = "connect.group.id"; + + public static void ignoreAutoPopulatedMetricsContextProperties(AbstractConfig config) { + Stream.of(CONNECT_KAFKA_CLUSTER_ID, CONNECT_GROUP_ID) + .map(property -> METRICS_CONTEXT_PREFIX + property) + .forEach(config::ignore); + } + /** * Postprocess the configuration so that exponential backoff is disabled when reconnect backoff * is explicitly configured but the maximum reconnect backoff is not explicitly configured. diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java b/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java index 16feef66d4351..5c53093bf0427 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java @@ -218,6 +218,8 @@ public class AdminClientConfig extends AbstractConfig { .withClientSaslSupport(); } + final boolean isSubConfig; + @Override protected Map postProcessParsedConfig(final Map parsedValues) { return CommonClientConfigs.postProcessReconnectBackoffConfigs(this, parsedValues); @@ -229,6 +231,7 @@ public AdminClientConfig(Map props) { protected AdminClientConfig(Map props, boolean doLog) { super(CONFIG, props, doLog); + this.isSubConfig = props instanceof RecordingMap; } public static Set configNames() { diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java index cf99556b0f2cf..b873f594f5bc5 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/KafkaAdminClient.java @@ -592,7 +592,9 @@ private KafkaAdminClient(AdminClientConfig config, new TimeoutProcessorFactory() : timeoutProcessorFactory; this.maxRetries = config.getInt(AdminClientConfig.RETRIES_CONFIG); this.retryBackoffMs = config.getLong(AdminClientConfig.RETRY_BACKOFF_MS_CONFIG); - config.logUnused(); + CommonClientConfigs.ignoreAutoPopulatedMetricsContextProperties(config); + if (!config.isSubConfig) + config.logUnused(); AppInfoParser.registerAppInfo(JMX_PREFIX, clientId, metrics, time.milliseconds()); log.debug("Kafka admin client initialized"); thread.start(); diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java index 48f1ccbf19513..5649bd49275fe 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java @@ -577,6 +577,8 @@ public class ConsumerConfig extends AbstractConfig { .withClientSaslSupport(); } + final boolean isSubConfig; + @Override protected Map postProcessParsedConfig(final Map parsedValues) { Map refinedConfigs = CommonClientConfigs.postProcessReconnectBackoffConfigs(this, parsedValues); @@ -601,7 +603,9 @@ private void maybeOverrideClientId(Map configs) { protected static Map appendDeserializerToConfig(Map configs, Deserializer keyDeserializer, Deserializer valueDeserializer) { - Map newConfigs = new HashMap<>(configs); + Map newConfigs = configs instanceof RecordingMap ? + ((RecordingMap) configs).copy() : + new HashMap<>(configs); if (keyDeserializer != null) newConfigs.put(KEY_DESERIALIZER_CLASS_CONFIG, keyDeserializer.getClass()); if (valueDeserializer != null) @@ -624,14 +628,17 @@ boolean maybeOverrideEnableAutoCommit() { public ConsumerConfig(Properties props) { super(CONFIG, props); + this.isSubConfig = false; } public ConsumerConfig(Map props) { super(CONFIG, props); + this.isSubConfig = props instanceof RecordingMap; } protected ConsumerConfig(Map props, boolean doLog) { super(CONFIG, props, doLog); + this.isSubConfig = props instanceof RecordingMap; } public static Set configNames() { diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java index d0ef1a0cbebc2..ef030c2c6792e 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/KafkaConsumer.java @@ -812,7 +812,10 @@ public KafkaConsumer(Map configs, this.kafkaConsumerMetrics = new KafkaConsumerMetrics(metrics, metricGrpPrefix); - config.logUnused(); + CommonClientConfigs.ignoreAutoPopulatedMetricsContextProperties(config); + + if (!config.isSubConfig) + config.logUnused(); AppInfoParser.registerAppInfo(JMX_PREFIX, clientId, metrics, time.milliseconds()); log.debug("Kafka consumer initialized"); } catch (Throwable t) { diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java index 96d9125831dad..cbf82ff616149 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/KafkaProducer.java @@ -429,7 +429,9 @@ public KafkaProducer(Properties properties, Serializer keySerializer, Seriali String ioThreadName = NETWORK_THREAD_PREFIX + " | " + clientId; this.ioThread = new KafkaThread(ioThreadName, this.sender, true); this.ioThread.start(); - config.logUnused(); + CommonClientConfigs.ignoreAutoPopulatedMetricsContextProperties(config); + if (!config.isSubConfig) + config.logUnused(); AppInfoParser.registerAppInfo(JMX_PREFIX, clientId, metrics, time.milliseconds()); log.debug("Kafka producer started"); } catch (Throwable t) { diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java index afc1e55cdfdad..03087fb77d805 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java @@ -453,6 +453,8 @@ public class ProducerConfig extends AbstractConfig { TRANSACTIONAL_ID_DOC); } + final boolean isSubConfig; + @Override protected Map postProcessParsedConfig(final Map parsedValues) { Map refinedConfigs = CommonClientConfigs.postProcessReconnectBackoffConfigs(this, parsedValues); @@ -537,7 +539,9 @@ private static String parseAcks(String acksString) { static Map appendSerializerToConfig(Map configs, Serializer keySerializer, Serializer valueSerializer) { - Map newConfigs = new HashMap<>(configs); + Map newConfigs = configs instanceof RecordingMap ? + ((RecordingMap) configs).copy() : + new HashMap<>(configs); if (keySerializer != null) newConfigs.put(KEY_SERIALIZER_CLASS_CONFIG, keySerializer.getClass()); if (valueSerializer != null) @@ -547,14 +551,17 @@ static Map appendSerializerToConfig(Map configs, public ProducerConfig(Properties props) { super(CONFIG, props); + this.isSubConfig = false; } public ProducerConfig(Map props) { super(CONFIG, props); + this.isSubConfig = props instanceof RecordingMap; } ProducerConfig(Map props, boolean doLog) { super(CONFIG, props, doLog); + this.isSubConfig = props instanceof RecordingMap; } public static Set configNames() { diff --git a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java index e3fda4d9f5406..36a6f8e2e3b97 100644 --- a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java +++ b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java @@ -382,7 +382,7 @@ private void logAll() { public void logUnused() { Set unusedkeys = unused(); if (!unusedkeys.isEmpty()) { - log.warn("These configurations '{}' were supplied but are not used yet.", unusedkeys); + log.warn("These configurations were supplied but are not used yet: {}", unusedkeys); } } @@ -609,7 +609,7 @@ public int hashCode() { * Marks keys retrieved via `get` as used. This is needed because `Configurable.configure` takes a `Map` instead * of an `AbstractConfig` and we can't change that without breaking public API like `Partitioner`. */ - private class RecordingMap extends HashMap { + protected class RecordingMap extends HashMap { private final String prefix; private final boolean withIgnoreFallback; @@ -649,6 +649,11 @@ public V get(Object key) { } return super.get(key); } + + public RecordingMap copy() { + return new RecordingMap<>(this); + } + } /** diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java index 8d93e795911b9..c2061da9751ae 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectDistributed.java @@ -43,7 +43,6 @@ import java.net.URI; import java.util.Arrays; import java.util.Collections; -import java.util.HashMap; import java.util.Map; /** @@ -104,7 +103,7 @@ public Connect startConnect(Map workerProps) { String workerId = advertisedUrl.getHost() + ":" + advertisedUrl.getPort(); // Create the admin client to be shared by all backing stores. - Map adminProps = new HashMap<>(config.originals()); + Map adminProps = config.originals(); ConnectUtils.addMetricsContextProperties(adminProps, config, kafkaClusterId); SharedTopicAdmin sharedAdmin = new SharedTopicAdmin(adminProps); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java index 19cc115d9ddc4..b292748c5cc18 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/cli/ConnectStandalone.java @@ -94,7 +94,7 @@ public static void main(String[] args) { Worker worker = new Worker(workerId, time, plugins, config, new FileOffsetBackingStore(), connectorClientConfigOverridePolicy); - Herder herder = new StandaloneHerder(worker, kafkaClusterId, connectorClientConfigOverridePolicy); + Herder herder = new StandaloneHerder(worker, kafkaClusterId, connectorClientConfigOverridePolicy, config); final Connect connect = new Connect(herder, rest); log.info("Kafka Connect standalone worker initialization took {}ms", time.hiResClockMs() - initStart); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java index 2fe75a955b06c..7483877eade15 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/AbstractHerder.java @@ -104,6 +104,7 @@ public abstract class AbstractHerder implements Herder, TaskStatus.Listener, Con protected final ConfigBackingStore configBackingStore; private final ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy; protected volatile boolean running = false; + private final WorkerConfig config; private final ExecutorService connectorExecutor; private final ConcurrentMap tempConnectors = new ConcurrentHashMap<>(); @@ -113,7 +114,8 @@ public AbstractHerder(Worker worker, String kafkaClusterId, StatusBackingStore statusBackingStore, ConfigBackingStore configBackingStore, - ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) { + ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, + WorkerConfig config) { this.worker = worker; this.worker.herder = this; this.workerId = workerId; @@ -121,6 +123,7 @@ public AbstractHerder(Worker worker, this.statusBackingStore = statusBackingStore; this.configBackingStore = configBackingStore; this.connectorClientConfigOverridePolicy = connectorClientConfigOverridePolicy; + this.config = config; this.connectorExecutor = Executors.newCachedThreadPool(); } @@ -135,6 +138,7 @@ protected void startServices() { this.worker.start(); this.statusBackingStore.start(); this.configBackingStore.start(); + config.logUnused(); } protected void stopServices() { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java index 7dad6aec0af1b..e3f90caf5195d 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/ConnectMetrics.java @@ -87,10 +87,10 @@ public ConnectMetrics(String workerId, WorkerConfig config, Time time, String cl Map contextLabels = new HashMap<>(); contextLabels.putAll(config.originalsWithPrefix(CommonClientConfigs.METRICS_CONTEXT_PREFIX)); - contextLabels.put(WorkerConfig.CONNECT_KAFKA_CLUSTER_ID, clusterId); + contextLabels.put(CommonClientConfigs.CONNECT_KAFKA_CLUSTER_ID, clusterId); Object groupId = config.originals().get(DistributedConfig.GROUP_ID_CONFIG); if (groupId != null) { - contextLabels.put(WorkerConfig.CONNECT_GROUP_ID, groupId); + contextLabels.put(CommonClientConfigs.CONNECT_GROUP_ID, groupId); } MetricsContext metricsContext = new KafkaMetricsContext(JMX_PREFIX, contextLabels); this.metrics = new Metrics(metricConfig, reporters, time, metricsContext); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitter.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitter.java index c3416be45c724..2c9e4cb94498b 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitter.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitter.java @@ -47,22 +47,25 @@ class SourceTaskOffsetCommitter { private static final Logger log = LoggerFactory.getLogger(SourceTaskOffsetCommitter.class); - private final WorkerConfig config; + private final long commitIntervalMs; private final ScheduledExecutorService commitExecutorService; private final ConcurrentMap> committers; // visible for testing - SourceTaskOffsetCommitter(WorkerConfig config, + SourceTaskOffsetCommitter(long commitIntervalMs, ScheduledExecutorService commitExecutorService, ConcurrentMap> committers) { - this.config = config; + this.commitIntervalMs = commitIntervalMs; this.commitExecutorService = commitExecutorService; this.committers = committers; } public SourceTaskOffsetCommitter(WorkerConfig config) { - this(config, Executors.newSingleThreadScheduledExecutor(ThreadUtils.createThreadFactory( - SourceTaskOffsetCommitter.class.getSimpleName() + "-%d", false)), + this( + config.getLong(WorkerConfig.OFFSET_COMMIT_INTERVAL_MS_CONFIG), + Executors.newSingleThreadScheduledExecutor(ThreadUtils.createThreadFactory( + SourceTaskOffsetCommitter.class.getSimpleName() + "-%d", false) + ), new ConcurrentHashMap<>()); } @@ -78,7 +81,6 @@ public void close(long timeoutMs) { } public void schedule(final ConnectorTaskId id, final WorkerSourceTask workerTask) { - long commitIntervalMs = config.getLong(WorkerConfig.OFFSET_COMMIT_INTERVAL_MS_CONFIG); ScheduledFuture commitFuture = commitExecutorService.scheduleWithFixedDelay(() -> { try (LoggingContext loggingContext = LoggingContext.forOffsets(id)) { commit(workerTask); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java index 4adf6ff5e0416..a0d121be76dc4 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/Worker.java @@ -78,6 +78,11 @@ import java.util.concurrent.Executors; import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; +import java.util.stream.Stream; + +import static org.apache.kafka.connect.runtime.WorkerConfig.CONNECTOR_ADMIN_PREFIX; +import static org.apache.kafka.connect.runtime.WorkerConfig.CONNECTOR_CONSUMER_PREFIX; +import static org.apache.kafka.connect.runtime.WorkerConfig.CONNECTOR_PRODUCER_PREFIX; /** *

@@ -641,26 +646,14 @@ static Map producerConfigs(ConnectorTaskId id, Class connectorClass, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, String clusterId) { + // Start with overrides specified by the user in the worker config + // We intentionally create a copy instead of using config::originalsWithPrefix directly + // since warnings for unused configs aren't issued by producers when they're configured + // with the result of AbstractConfig::originals or AbstractConfig::originalsWithPrefix Map producerProps = new HashMap<>(); - producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); - producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); - producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); - // These settings will execute infinite retries on retriable exceptions. They *may* be overridden via configs passed to the worker, - // but this may compromise the delivery guarantees of Kafka Connect. - producerProps.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); - // By default, Connect disables idempotent behavior for all producers, even though idempotence became - // default for Kafka producers. This is to ensure Connect continues to work with many Kafka broker versions, including older brokers that do not support - // idempotent producers or require explicit steps to enable them (e.g. adding the IDEMPOTENT_WRITE ACL to brokers older than 2.8). - // These settings might change when https://cwiki.apache.org/confluence/display/KAFKA/KIP-318%3A+Make+Kafka+Connect+Source+idempotent - // gets approved and scheduled for release. - producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "false"); - producerProps.put(ProducerConfig.ACKS_CONFIG, "all"); - producerProps.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "1"); - producerProps.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); - producerProps.put(ProducerConfig.CLIENT_ID_CONFIG, defaultClientId); - // User-specified overrides - producerProps.putAll(config.originalsWithPrefix("producer.")); - //add client metrics.context properties + producerProps.putAll(config.originalsWithPrefix(CONNECTOR_PRODUCER_PREFIX)); + + // Metrics context properties ConnectUtils.addMetricsContextProperties(producerProps, config, clusterId); // Connector-specified overrides @@ -670,6 +663,24 @@ static Map producerConfigs(ConnectorTaskId id, connectorClientConfigOverridePolicy); producerProps.putAll(producerOverrides); + // Default properties to use if no values for them have been specified anywhere else + producerProps.putIfAbsent(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); + producerProps.putIfAbsent(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + producerProps.putIfAbsent(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + // These settings will execute infinite retries on retriable exceptions. They *may* be overridden via configs passed to the worker, + // but this may compromise the delivery guarantees of Kafka Connect. + producerProps.putIfAbsent(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); + // By default, Connect disables idempotent behavior for all producers, even though idempotence became + // default for Kafka producers. This is to ensure Connect continues to work with many Kafka broker versions, including older brokers that do not support + // idempotent producers or require explicit steps to enable them (e.g. adding the IDEMPOTENT_WRITE ACL to brokers older than 2.8). + // These settings might change when https://cwiki.apache.org/confluence/display/KAFKA/KIP-318%3A+Make+Kafka+Connect+Source+idempotent + // gets approved and scheduled for release. + producerProps.putIfAbsent(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "false"); + producerProps.putIfAbsent(ProducerConfig.ACKS_CONFIG, "all"); + producerProps.putIfAbsent(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "1"); + producerProps.putIfAbsent(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + producerProps.putIfAbsent(ProducerConfig.CLIENT_ID_CONFIG, defaultClientId); + return producerProps; } @@ -679,22 +690,16 @@ static Map consumerConfigs(ConnectorTaskId id, Class connectorClass, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, String clusterId) { - // Include any unknown worker configs so consumer configs can be set globally on the worker - // and through to the task + // Start with overrides specified by the user in the worker config + // We intentionally create a copy instead of using config::originalsWithPrefix directly + // since warnings for unused configs aren't issued by consumers when they're configured + // with the result of AbstractConfig::originals or AbstractConfig::originalsWithPrefix Map consumerProps = new HashMap<>(); + consumerProps.putAll(config.originalsWithPrefix(CONNECTOR_CONSUMER_PREFIX)); - consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, SinkUtils.consumerGroupId(id.connector())); - consumerProps.put(ConsumerConfig.CLIENT_ID_CONFIG, "connector-consumer-" + id); - consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, - Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); - consumerProps.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); - consumerProps.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - - consumerProps.putAll(config.originalsWithPrefix("consumer.")); - //add client metrics.context properties + // Metrics context properties ConnectUtils.addMetricsContextProperties(consumerProps, config, clusterId); + // Connector-specified overrides Map consumerOverrides = connectorClientConfigOverrides(id, connConfig, connectorClass, ConnectorConfig.CONNECTOR_CLIENT_CONSUMER_OVERRIDES_PREFIX, @@ -702,6 +707,16 @@ static Map consumerConfigs(ConnectorTaskId id, connectorClientConfigOverridePolicy); consumerProps.putAll(consumerOverrides); + // Default properties to use if no values for them have been specified anywhere else + consumerProps.putIfAbsent(ConsumerConfig.GROUP_ID_CONFIG, SinkUtils.consumerGroupId(id.connector())); + consumerProps.putIfAbsent(ConsumerConfig.CLIENT_ID_CONFIG, "connector-consumer-" + id); + consumerProps.putIfAbsent(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, + Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); + consumerProps.putIfAbsent(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); + consumerProps.putIfAbsent(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + consumerProps.putIfAbsent(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + consumerProps.putIfAbsent(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + return consumerProps; } @@ -712,21 +727,18 @@ static Map adminConfigs(ConnectorTaskId id, Class connectorClass, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, String clusterId) { - Map adminProps = new HashMap<>(); // Use the top-level worker configs to retain backwards compatibility with older releases which // did not require a prefix for connector admin client configs in the worker configuration file + // Use config::originals directly here in order to prevent the admin client from issuing warnings + // for unused configs since it's basically guaranteed that some properties in the worker config + // are not used by admin clients + Map adminProps = config.originals(); + // Ignore configs that begin with "admin." since those will be added next (with the prefix stripped) // and those that begin with "producer." and "consumer.", since we know they aren't intended for // the admin client - Map nonPrefixedWorkerConfigs = config.originals().entrySet().stream() - .filter(e -> !e.getKey().startsWith("admin.") - && !e.getKey().startsWith("producer.") - && !e.getKey().startsWith("consumer.")) - .collect(Collectors.toMap(Map.Entry::getKey, Map.Entry::getValue)); - adminProps.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, - Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); - adminProps.put(AdminClientConfig.CLIENT_ID_CONFIG, defaultClientId); - adminProps.putAll(nonPrefixedWorkerConfigs); + Stream.of(CONNECTOR_ADMIN_PREFIX, CONNECTOR_CONSUMER_PREFIX, CONNECTOR_PRODUCER_PREFIX) + .forEach(prefix -> adminProps.keySet().removeIf(property -> property.startsWith(prefix))); // Admin client-specific overrides in the worker config adminProps.putAll(config.originalsWithPrefix("admin.")); @@ -738,7 +750,10 @@ static Map adminConfigs(ConnectorTaskId id, connectorClientConfigOverridePolicy); adminProps.putAll(adminOverrides); - //add client metrics.context properties + // Default client ID to use if none has been specified anywhere else + adminProps.putIfAbsent(AdminClientConfig.CLIENT_ID_CONFIG, defaultClientId); + + // Metrics context properties ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); return adminProps; diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java index 3224a230f90ef..92732c62e5131 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerConfig.java @@ -207,9 +207,6 @@ public class WorkerConfig extends AbstractConfig { + "user requests to reset the set of active topics per connector."; protected static final boolean TOPIC_TRACKING_ALLOW_RESET_DEFAULT = true; - public static final String CONNECT_KAFKA_CLUSTER_ID = "connect.kafka.cluster.id"; - public static final String CONNECT_GROUP_ID = "connect.group.id"; - public static final String TOPIC_CREATION_ENABLE_CONFIG = "topic.creation.enable"; protected static final String TOPIC_CREATION_ENABLE_DOC = "Whether to allow " + "automatic creation of topics used by source connectors, when source connectors " @@ -222,6 +219,10 @@ public class WorkerConfig extends AbstractConfig { protected static final String RESPONSE_HTTP_HEADERS_DOC = "Rules for REST API HTTP response headers"; protected static final String RESPONSE_HTTP_HEADERS_DEFAULT = ""; + public static final String CONNECTOR_ADMIN_PREFIX = "admin."; + public static final String CONNECTOR_CONSUMER_PREFIX = "consumer."; + public static final String CONNECTOR_PRODUCER_PREFIX = "producer."; + /** * Get a basic ConfigDef for a WorkerConfig. This includes all the common settings. Subclasses can use this to * bootstrap their own ConfigDef. @@ -327,6 +328,8 @@ private void logInternalConverterRemovalWarnings(Map props) { private void logPluginPathConfigProviderWarning(Map rawOriginals) { String rawPluginPath = rawOriginals.get(PLUGIN_PATH_CONFIG); + if (rawPluginPath == null) + return; // Can't use AbstractConfig::originalsStrings here since some values may be null, which // causes that method to fail String transformedPluginPath = Objects.toString(originals().get(PLUGIN_PATH_CONFIG)); @@ -366,6 +369,29 @@ public WorkerConfig(ConfigDef definition, Map props) { super(definition, props); logInternalConverterRemovalWarnings(props); logPluginPathConfigProviderWarning(props); + ignoreSubConfigs(); + } + + private void ignoreSubConfigs() { + subConfigPrefixes().forEach(this::ignoreAllWithPrefixes); + Arrays.asList( + KEY_CONVERTER_CLASS_CONFIG, VALUE_CONVERTER_CLASS_CONFIG, HEADER_CONVERTER_CLASS_CONFIG + ).forEach(this::ignore); + } + + protected List subConfigPrefixes() { + return new ArrayList<>(Arrays.asList( + KEY_CONVERTER_CLASS_CONFIG + ".", + VALUE_CONVERTER_CLASS_CONFIG + ".", + HEADER_CONVERTER_CLASS_CONFIG + ".", + CONNECTOR_ADMIN_PREFIX, + CONNECTOR_CONSUMER_PREFIX, + CONNECTOR_PRODUCER_PREFIX + )); + } + + protected void ignoreAllWithPrefixes(String prefix) { + originalsWithPrefix(prefix, false).keySet().forEach(this::ignore); } // Visible for testing diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedConfig.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedConfig.java index 0823fbcc30ad3..73bd6bbdd1ad1 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedConfig.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedConfig.java @@ -29,6 +29,7 @@ import javax.crypto.Mac; import java.security.InvalidParameterException; import java.security.NoSuchAlgorithmException; +import java.util.Arrays; import java.util.Collections; import java.util.List; import java.util.Map; @@ -401,6 +402,13 @@ public Integer getRebalanceTimeout() { return getInt(DistributedConfig.REBALANCE_TIMEOUT_MS_CONFIG); } + @Override + protected List subConfigPrefixes() { + List result = super.subConfigPrefixes(); + result.addAll(Arrays.asList(CONFIG_STORAGE_PREFIX, OFFSET_STORAGE_PREFIX, STATUS_STORAGE_PREFIX)); + return result; + } + public DistributedConfig(Map props) { super(CONFIG, props); getInternalRequestKeyGenerator(); // Check here for a valid key size + key algorithm to fail fast if either are invalid diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java index 65a8e7e15b818..0a9871794eb87 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/DistributedHerder.java @@ -241,7 +241,7 @@ public DistributedHerder(DistributedConfig config, Time time, ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, AutoCloseable... uponShutdown) { - super(worker, workerId, kafkaClusterId, statusBackingStore, configBackingStore, connectorClientConfigOverridePolicy); + super(worker, workerId, kafkaClusterId, statusBackingStore, configBackingStore, connectorClientConfigOverridePolicy, config); this.time = time; this.herderMetrics = new HerderMetrics(metrics); @@ -1250,7 +1250,7 @@ private boolean handleRebalanceCompleted() { } } else { if (configState.offset() < assignment.offset()) { - log.warn("Catching up to assignment's config offset."); + log.info("Catching up to assignment's config offset."); needsReadToEnd = true; } } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java index 4c1d6a5b9e1cc..78941e6bf7912 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMember.java @@ -37,7 +37,6 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.Time; import org.apache.kafka.common.utils.Utils; -import org.apache.kafka.connect.runtime.WorkerConfig; import org.apache.kafka.connect.storage.ConfigBackingStore; import org.apache.kafka.connect.util.ConnectUtils; import org.apache.kafka.connect.util.ConnectorTaskId; @@ -98,8 +97,8 @@ public WorkerGroupMember(DistributedConfig config, Map contextLabels = new HashMap<>(); contextLabels.putAll(config.originalsWithPrefix(CommonClientConfigs.METRICS_CONTEXT_PREFIX)); - contextLabels.put(WorkerConfig.CONNECT_KAFKA_CLUSTER_ID, ConnectUtils.lookupKafkaClusterId(config)); - contextLabels.put(WorkerConfig.CONNECT_GROUP_ID, config.getString(DistributedConfig.GROUP_ID_CONFIG)); + contextLabels.put(CommonClientConfigs.CONNECT_KAFKA_CLUSTER_ID, ConnectUtils.lookupKafkaClusterId(config)); + contextLabels.put(CommonClientConfigs.CONNECT_GROUP_ID, config.getString(DistributedConfig.GROUP_ID_CONFIG)); MetricsContext metricsContext = new KafkaMetricsContext(JMX_PREFIX, contextLabels); this.metrics = new Metrics(metricConfig, reporters, time, metricsContext); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java index dac389ba0e346..a4b44326db352 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerder.java @@ -71,13 +71,15 @@ public class StandaloneHerder extends AbstractHerder { private ClusterConfigState configState; public StandaloneHerder(Worker worker, String kafkaClusterId, - ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) { + ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, + StandaloneConfig config) { this(worker, worker.workerId(), kafkaClusterId, new MemoryStatusBackingStore(), new MemoryConfigBackingStore(worker.configTransformer()), - connectorClientConfigOverridePolicy); + connectorClientConfigOverridePolicy, + config); } // visible for testing @@ -86,8 +88,9 @@ public StandaloneHerder(Worker worker, String kafkaClusterId, String kafkaClusterId, StatusBackingStore statusBackingStore, MemoryConfigBackingStore configBackingStore, - ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy) { - super(worker, workerId, kafkaClusterId, statusBackingStore, configBackingStore, connectorClientConfigOverridePolicy); + ConnectorClientConfigOverridePolicy connectorClientConfigOverridePolicy, + StandaloneConfig config) { + super(worker, workerId, kafkaClusterId, statusBackingStore, configBackingStore, connectorClientConfigOverridePolicy, config); this.configState = ClusterConfigState.EMPTY; this.requestExecutorService = Executors.newSingleThreadScheduledExecutor(); configBackingStore.setUpdateListener(new ConfigUpdateListener()); diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java index 94b98cb9eb38d..a0c5ed237efa5 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaConfigBackingStore.java @@ -506,8 +506,7 @@ public void putRestartRequest(RestartRequest restartRequest) { // package private for testing KafkaBasedLog setupAndCreateKafkaBasedLog(String topic, final WorkerConfig config) { String clusterId = ConnectUtils.lookupKafkaClusterId(config); - Map originals = config.originals(); - Map producerProps = new HashMap<>(originals); + Map producerProps = config.originals(); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName()); producerProps.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.MAX_VALUE); @@ -519,12 +518,12 @@ KafkaBasedLog setupAndCreateKafkaBasedLog(String topic, final Wo producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "false"); ConnectUtils.addMetricsContextProperties(producerProps, config, clusterId); - Map consumerProps = new HashMap<>(originals); + Map consumerProps = config.originals(); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName()); ConnectUtils.addMetricsContextProperties(consumerProps, config, clusterId); - Map adminProps = new HashMap<>(originals); + Map adminProps = config.originals(); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); Supplier adminSupplier; if (topicAdminSupplier != null) { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java index f3cbb686fff50..798d65a675504 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaOffsetBackingStore.java @@ -86,8 +86,7 @@ public void configure(final WorkerConfig config) { String clusterId = ConnectUtils.lookupKafkaClusterId(config); data = new HashMap<>(); - Map originals = config.originals(); - Map producerProps = new HashMap<>(originals); + Map producerProps = config.originals(); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName()); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName()); producerProps.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.MAX_VALUE); @@ -99,12 +98,12 @@ public void configure(final WorkerConfig config) { producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "false"); ConnectUtils.addMetricsContextProperties(producerProps, config, clusterId); - Map consumerProps = new HashMap<>(originals); + Map consumerProps = config.originals(); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName()); ConnectUtils.addMetricsContextProperties(consumerProps, config, clusterId); - Map adminProps = new HashMap<>(originals); + Map adminProps = config.originals(); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); Supplier adminSupplier; if (topicAdminSupplier != null) { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java index 3ba6996da8ab7..266a7cdebd94f 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/storage/KafkaStatusBackingStore.java @@ -165,8 +165,7 @@ public void configure(final WorkerConfig config) { throw new ConfigException("Must specify topic for connector status."); String clusterId = ConnectUtils.lookupKafkaClusterId(config); - Map originals = config.originals(); - Map producerProps = new HashMap<>(originals); + Map producerProps = config.originals(); producerProps.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, StringSerializer.class.getName()); producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, ByteArraySerializer.class.getName()); producerProps.put(ProducerConfig.RETRIES_CONFIG, 0); // we handle retries in this class @@ -178,12 +177,12 @@ public void configure(final WorkerConfig config) { producerProps.put(ProducerConfig.ENABLE_IDEMPOTENCE_CONFIG, "false"); // disable idempotence since retries is force to 0 ConnectUtils.addMetricsContextProperties(producerProps, config, clusterId); - Map consumerProps = new HashMap<>(originals); + Map consumerProps = config.originals(); consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName()); consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName()); ConnectUtils.addMetricsContextProperties(consumerProps, config, clusterId); - Map adminProps = new HashMap<>(originals); + Map adminProps = config.originals(); ConnectUtils.addMetricsContextProperties(adminProps, config, clusterId); Supplier adminSupplier; if (topicAdminSupplier != null) { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java index 7adbd8f92dfd8..8f2911801787c 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/ConnectUtils.java @@ -145,10 +145,10 @@ public static void addMetricsContextProperties(Map prop, WorkerC //add all properties predefined with "metrics.context." prop.putAll(config.originalsWithPrefix(CommonClientConfigs.METRICS_CONTEXT_PREFIX, false)); //add connect properties - prop.put(CommonClientConfigs.METRICS_CONTEXT_PREFIX + WorkerConfig.CONNECT_KAFKA_CLUSTER_ID, clusterId); + prop.put(CommonClientConfigs.METRICS_CONTEXT_PREFIX + CommonClientConfigs.CONNECT_KAFKA_CLUSTER_ID, clusterId); Object groupId = config.originals().get(DistributedConfig.GROUP_ID_CONFIG); if (groupId != null) { - prop.put(CommonClientConfigs.METRICS_CONTEXT_PREFIX + WorkerConfig.CONNECT_GROUP_ID, groupId); + prop.put(CommonClientConfigs.METRICS_CONTEXT_PREFIX + CommonClientConfigs.CONNECT_GROUP_ID, groupId); } } 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 5b9e199e5a1ee..83234f5d045e6 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 @@ -50,6 +50,7 @@ import org.apache.kafka.connect.util.ConnectorTaskId; import org.easymock.Capture; import org.easymock.EasyMock; +import org.easymock.Mock; import org.junit.Test; import org.junit.runner.RunWith; import org.powermock.api.easymock.PowerMock; @@ -146,6 +147,7 @@ public class AbstractHerderTest { @MockStrict private ClassLoader classLoader; @MockStrict private ConfigBackingStore configStore; @MockStrict private StatusBackingStore statusStore; + @Mock private WorkerConfig config; @Test public void testConnectors() { @@ -156,9 +158,10 @@ public void testConnectors() { String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class + ConnectorClientConfigOverridePolicy.class, + WorkerConfig.class ) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -180,9 +183,10 @@ public void testConnectorStatus() { String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class + ConnectorClientConfigOverridePolicy.class, + WorkerConfig.class ) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -205,8 +209,8 @@ public void connectorStatus() { AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -246,8 +250,8 @@ public void taskStatus() { AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -278,8 +282,8 @@ public void testBuildRestartPlanForUnknownConnector() { RestartRequest restartRequest = new RestartRequest(connectorName, false, true); AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -345,8 +349,8 @@ public void testBuildRestartPlanForConnectorAndTasks() { AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -387,8 +391,8 @@ public void testBuildRestartPlanForNoRestart() { AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -913,9 +917,10 @@ public void testConnectorPluginConfig() throws Exception { String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class + ConnectorClientConfigOverridePolicy.class, + WorkerConfig.class ) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); @@ -963,8 +968,8 @@ public void testGetConnectorConfigDefWithBadName() throws Exception { String connName = "AnotherPlugin"; AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); EasyMock.expect(worker.getPlugins()).andStubReturn(plugins); @@ -978,8 +983,8 @@ public void testGetConnectorConfigDefWithInvalidPluginType() throws Exception { String connName = "AnotherPlugin"; AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, noneConnectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); EasyMock.expect(worker.getPlugins()).andStubReturn(plugins); @@ -1049,8 +1054,8 @@ private AbstractHerder createConfigValidationHerder(Class c AbstractHerder herder = partialMockBuilder(AbstractHerder.class) .withConstructor(Worker.class, String.class, String.class, StatusBackingStore.class, ConfigBackingStore.class, - ConnectorClientConfigOverridePolicy.class) - .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, connectorClientConfigOverridePolicy) + ConnectorClientConfigOverridePolicy.class, WorkerConfig.class) + .withArgs(worker, workerId, kafkaClusterId, statusStore, configStore, connectorClientConfigOverridePolicy, config) .addMockedMethod("generation") .createMock(); EasyMock.expect(herder.generation()).andStubReturn(generation); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitterTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitterTest.java index 278a73d16d4df..8b5fa9e00d137 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitterTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/SourceTaskOffsetCommitterTest.java @@ -17,7 +17,6 @@ package org.apache.kafka.connect.runtime; import org.apache.kafka.connect.errors.ConnectException; -import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; import org.apache.kafka.connect.util.ConnectorTaskId; import org.apache.kafka.connect.util.ThreadedTest; import org.easymock.Capture; @@ -30,8 +29,6 @@ import org.powermock.reflect.Whitebox; import org.slf4j.Logger; -import java.util.HashMap; -import java.util.Map; import java.util.concurrent.CancellationException; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledExecutorService; @@ -59,19 +56,12 @@ public class SourceTaskOffsetCommitterTest extends ThreadedTest { private SourceTaskOffsetCommitter committer; - private static final long DEFAULT_OFFSET_COMMIT_INTERVAL_MS = 1000; + private static final long OFFSET_COMMIT_INTERVAL_MS = 1000; @Override public void setup() { super.setup(); - Map workerProps = new HashMap<>(); - workerProps.put("key.converter", "org.apache.kafka.connect.json.JsonConverter"); - workerProps.put("value.converter", "org.apache.kafka.connect.json.JsonConverter"); - workerProps.put("offset.storage.file.filename", "/tmp/connect.offsets"); - workerProps.put("offset.flush.interval.ms", - Long.toString(DEFAULT_OFFSET_COMMIT_INTERVAL_MS)); - WorkerConfig config = new StandaloneConfig(workerProps); - committer = new SourceTaskOffsetCommitter(config, executor, committers); + committer = new SourceTaskOffsetCommitter(OFFSET_COMMIT_INTERVAL_MS, executor, committers); Whitebox.setInternalState(SourceTaskOffsetCommitter.class, "log", mockLog); } @@ -81,8 +71,8 @@ public void testSchedule() { Capture taskWrapper = EasyMock.newCapture(); EasyMock.expect(executor.scheduleWithFixedDelay( - EasyMock.capture(taskWrapper), eq(DEFAULT_OFFSET_COMMIT_INTERVAL_MS), - eq(DEFAULT_OFFSET_COMMIT_INTERVAL_MS), eq(TimeUnit.MILLISECONDS)) + EasyMock.capture(taskWrapper), eq(OFFSET_COMMIT_INTERVAL_MS), + eq(OFFSET_COMMIT_INTERVAL_MS), eq(TimeUnit.MILLISECONDS)) ).andReturn((ScheduledFuture) commitFuture); PowerMock.replayAll(); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java index dcd9286480e87..7827cb7ca88c0 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerTest.java @@ -88,6 +88,7 @@ import static org.apache.kafka.connect.runtime.TopicCreationConfig.DEFAULT_TOPIC_CREATION_PREFIX; import static org.apache.kafka.connect.runtime.TopicCreationConfig.PARTITIONS_CONFIG; import static org.apache.kafka.connect.runtime.TopicCreationConfig.REPLICATION_FACTOR_CONFIG; +import static org.apache.kafka.connect.runtime.WorkerConfig.BOOTSTRAP_SERVERS_CONFIG; import static org.apache.kafka.connect.runtime.WorkerConfig.TOPIC_CREATION_ENABLE_CONFIG; import static org.hamcrest.CoreMatchers.instanceOf; import static org.hamcrest.MatcherAssert.assertThat; @@ -196,6 +197,7 @@ public void setup() { .strictness(Strictness.STRICT_STUBS) .startMocking(); + workerProps.put(BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); workerProps.put("key.converter", "org.apache.kafka.connect.json.JsonConverter"); workerProps.put("value.converter", "org.apache.kafka.connect.json.JsonConverter"); workerProps.put("offset.storage.file.filename", "/tmp/connect.offsets"); @@ -1163,7 +1165,7 @@ public void testWorkerMetrics() throws Exception { if (reporter instanceof MockMetricsReporter) { MockMetricsReporter mockMetricsReporter = (MockMetricsReporter) reporter; //verify connect cluster is set in MetricsContext - assertEquals(CLUSTER_ID, mockMetricsReporter.getMetricsContext().contextLabels().get(WorkerConfig.CONNECT_KAFKA_CLUSTER_ID)); + assertEquals(CLUSTER_ID, mockMetricsReporter.getMetricsContext().contextLabels().get(CommonClientConfigs.CONNECT_KAFKA_CLUSTER_ID)); } } //verify metric is created with correct jmx prefix diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMemberTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMemberTest.java index 05cd01734fef8..feefb7cd0bf01 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMemberTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/distributed/WorkerGroupMemberTest.java @@ -23,7 +23,6 @@ import org.apache.kafka.common.utils.LogContext; import org.apache.kafka.common.utils.Time; import org.apache.kafka.connect.runtime.MockConnectMetrics; -import org.apache.kafka.connect.runtime.WorkerConfig; import org.apache.kafka.connect.storage.ConfigBackingStore; import org.apache.kafka.connect.storage.StatusBackingStore; import org.apache.kafka.connect.util.ConnectUtils; @@ -82,8 +81,8 @@ public void testMetrics() throws Exception { if (reporter instanceof MockConnectMetrics.MockMetricsReporter) { entered = true; MockConnectMetrics.MockMetricsReporter mockMetricsReporter = (MockConnectMetrics.MockMetricsReporter) reporter; - assertEquals("cluster-1", mockMetricsReporter.getMetricsContext().contextLabels().get(WorkerConfig.CONNECT_KAFKA_CLUSTER_ID)); - assertEquals("group-1", mockMetricsReporter.getMetricsContext().contextLabels().get(WorkerConfig.CONNECT_GROUP_ID)); + assertEquals("cluster-1", mockMetricsReporter.getMetricsContext().contextLabels().get(CommonClientConfigs.CONNECT_KAFKA_CLUSTER_ID)); + assertEquals("group-1", mockMetricsReporter.getMetricsContext().contextLabels().get(CommonClientConfigs.CONNECT_GROUP_ID)); } } assertTrue("Failed to verify MetricsReporter", entered); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java index f5ee4ccd310d7..38e291c7ea6c7 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/standalone/StandaloneHerderTest.java @@ -118,6 +118,7 @@ private enum SourceSink { private DelegatingClassLoader delegatingLoader; protected FutureCallback> createCallback; @Mock protected StatusBackingStore statusBackingStore; + @Mock private StandaloneConfig config; private final ConnectorClientConfigOverridePolicy noneConnectorClientConfigOverridePolicy = new NoneConnectorClientConfigOverridePolicy(); @@ -128,7 +129,7 @@ public void setup() { worker = PowerMock.createMock(Worker.class); String[] methodNames = new String[]{"connectorTypeForClass"/*, "validateConnectorConfig"*/, "buildRestartPlan", "recordRestarting"}; herder = PowerMock.createPartialMock(StandaloneHerder.class, methodNames, - worker, WORKER_ID, KAFKA_CLUSTER_ID, statusBackingStore, new MemoryConfigBackingStore(transformer), noneConnectorClientConfigOverridePolicy); + worker, WORKER_ID, KAFKA_CLUSTER_ID, statusBackingStore, new MemoryConfigBackingStore(transformer), noneConnectorClientConfigOverridePolicy, config); createCallback = new FutureCallback<>(); plugins = PowerMock.createMock(Plugins.class); pluginLoader = PowerMock.createMock(PluginClassLoader.class); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/util/ConnectUtilsTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/util/ConnectUtilsTest.java index d1330277e8e45..041aaddf400fb 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/util/ConnectUtilsTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/util/ConnectUtilsTest.java @@ -20,7 +20,6 @@ import org.apache.kafka.clients.admin.MockAdminClient; import org.apache.kafka.common.Node; import org.apache.kafka.connect.errors.ConnectException; -import org.apache.kafka.connect.runtime.WorkerConfig; import org.apache.kafka.connect.runtime.distributed.DistributedConfig; import org.apache.kafka.connect.runtime.standalone.StandaloneConfig; import org.junit.Test; @@ -84,8 +83,8 @@ public void testAddMetricsContextPropertiesDistributed() { Map prop = new HashMap<>(); ConnectUtils.addMetricsContextProperties(prop, config, "cluster-1"); - assertEquals("connect-cluster", prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + WorkerConfig.CONNECT_GROUP_ID)); - assertEquals("cluster-1", prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + WorkerConfig.CONNECT_KAFKA_CLUSTER_ID)); + assertEquals("connect-cluster", prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + CommonClientConfigs.CONNECT_GROUP_ID)); + assertEquals("cluster-1", prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + CommonClientConfigs.CONNECT_KAFKA_CLUSTER_ID)); } @Test @@ -99,8 +98,8 @@ public void testAddMetricsContextPropertiesStandalone() { Map prop = new HashMap<>(); ConnectUtils.addMetricsContextProperties(prop, config, "cluster-1"); - assertNull(prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + WorkerConfig.CONNECT_GROUP_ID)); - assertEquals("cluster-1", prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + WorkerConfig.CONNECT_KAFKA_CLUSTER_ID)); + assertNull(prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + CommonClientConfigs.CONNECT_GROUP_ID)); + assertEquals("cluster-1", prop.get(CommonClientConfigs.METRICS_CONTEXT_PREFIX + CommonClientConfigs.CONNECT_KAFKA_CLUSTER_ID)); } From 8295a6c42513c6475ca0e427acb29f5b974e9659 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Fri, 8 Apr 2022 13:55:52 -0400 Subject: [PATCH 2/3] KAFKA-7509: Extract RecordingMap into separate file in internal package --- .../clients/admin/AdminClientConfig.java | 1 + .../clients/consumer/ConsumerConfig.java | 1 + .../clients/producer/ProducerConfig.java | 1 + .../kafka/common/config/AbstractConfig.java | 70 +++------------- .../common/config/internals/RecordingMap.java | 84 +++++++++++++++++++ 5 files changed, 97 insertions(+), 60 deletions(-) create mode 100644 clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java diff --git a/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java b/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java index 5c53093bf0427..8f70ebc73e12f 100644 --- a/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/admin/AdminClientConfig.java @@ -24,6 +24,7 @@ import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; import org.apache.kafka.common.config.SecurityConfig; +import org.apache.kafka.common.config.internals.RecordingMap; import org.apache.kafka.common.metrics.Sensor; import java.util.Map; diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java index 5649bd49275fe..d99d0b9d69d82 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java @@ -24,6 +24,7 @@ import org.apache.kafka.common.config.ConfigDef.Importance; import org.apache.kafka.common.config.ConfigDef.Type; import org.apache.kafka.common.config.SecurityConfig; +import org.apache.kafka.common.config.internals.RecordingMap; import org.apache.kafka.common.errors.InvalidConfigurationException; import org.apache.kafka.common.metrics.Sensor; import org.apache.kafka.common.requests.JoinGroupRequest; diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java index 03087fb77d805..5189de3b0f442 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java @@ -25,6 +25,7 @@ import org.apache.kafka.common.config.ConfigDef.Type; import org.apache.kafka.common.config.ConfigException; import org.apache.kafka.common.config.SecurityConfig; +import org.apache.kafka.common.config.internals.RecordingMap; import org.apache.kafka.common.metrics.Sensor; import org.apache.kafka.common.serialization.Serializer; import org.slf4j.Logger; diff --git a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java index 36a6f8e2e3b97..9c39ffb1cd2fb 100644 --- a/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java +++ b/clients/src/main/java/org/apache/kafka/common/config/AbstractConfig.java @@ -18,6 +18,7 @@ import org.apache.kafka.common.Configurable; import org.apache.kafka.common.KafkaException; +import org.apache.kafka.common.config.internals.RecordingMap; import org.apache.kafka.common.config.types.Password; import org.apache.kafka.common.utils.Utils; import org.slf4j.Logger; @@ -227,13 +228,13 @@ public Set unused() { } public Map originals() { - Map copy = new RecordingMap<>(); + Map copy = new RecordingMap<>(this); copy.putAll(originals); return copy; } public Map originals(Map configOverrides) { - Map copy = new RecordingMap<>(); + Map copy = new RecordingMap<>(this); copy.putAll(originals); copy.putAll(configOverrides); return copy; @@ -245,7 +246,7 @@ public Map originals(Map configOverrides) { * @throws ClassCastException if any of the values are not strings */ public Map originalsStrings() { - Map copy = new RecordingMap<>(); + Map copy = new RecordingMap<>(this); for (Map.Entry entry : originals.entrySet()) { if (!(entry.getValue() instanceof String)) throw new ClassCastException("Non-string value found in original settings for key " + entry.getKey() + @@ -273,7 +274,7 @@ public Map originalsWithPrefix(String prefix) { * @return a Map containing the settings with the prefix */ public Map originalsWithPrefix(String prefix, boolean strip) { - Map result = new RecordingMap<>(prefix, false); + Map result = new RecordingMap<>(this, prefix, false); for (Map.Entry entry : originals.entrySet()) { if (entry.getKey().startsWith(prefix) && entry.getKey().length() > prefix.length()) { if (strip) @@ -302,7 +303,7 @@ public Map originalsWithPrefix(String prefix, boolean strip) { *

*/ public Map valuesWithPrefixOverride(String prefix) { - Map result = new RecordingMap<>(values(), prefix, true); + Map result = new RecordingMap<>(this, values(), prefix, true); for (Map.Entry entry : originals.entrySet()) { if (entry.getKey().startsWith(prefix) && entry.getKey().length() > prefix.length()) { String keyWithNoPrefix = entry.getKey().substring(prefix.length()); @@ -331,9 +332,9 @@ public Map valuesWithPrefixAllOrNothing(String prefix) { Map withPrefix = originalsWithPrefix(prefix, true); if (withPrefix.isEmpty()) { - return new RecordingMap<>(values(), "", true); + return new RecordingMap<>(this, values(), "", true); } else { - Map result = new RecordingMap<>(prefix, true); + Map result = new RecordingMap<>(this, prefix, true); for (Map.Entry entry : withPrefix.entrySet()) { ConfigDef.ConfigKey configKey = definition.configKeys().get(entry.getKey()); @@ -346,11 +347,11 @@ public Map valuesWithPrefixAllOrNothing(String prefix) { } public Map values() { - return new RecordingMap<>(values); + return new RecordingMap<>(this, values); } public Map nonInternalValues() { - Map nonInternalConfigs = new RecordingMap<>(); + Map nonInternalConfigs = new RecordingMap<>(this); values.forEach((key, value) -> { ConfigDef.ConfigKey configKey = definition.configKeys().get(key); if (configKey == null || !configKey.internalConfig) { @@ -605,57 +606,6 @@ public int hashCode() { return originals.hashCode(); } - /** - * Marks keys retrieved via `get` as used. This is needed because `Configurable.configure` takes a `Map` instead - * of an `AbstractConfig` and we can't change that without breaking public API like `Partitioner`. - */ - protected class RecordingMap extends HashMap { - - private final String prefix; - private final boolean withIgnoreFallback; - - RecordingMap() { - this("", false); - } - - RecordingMap(String prefix, boolean withIgnoreFallback) { - this.prefix = prefix; - this.withIgnoreFallback = withIgnoreFallback; - } - - RecordingMap(Map m) { - this(m, "", false); - } - - RecordingMap(Map m, String prefix, boolean withIgnoreFallback) { - super(m); - this.prefix = prefix; - this.withIgnoreFallback = withIgnoreFallback; - } - - @Override - public V get(Object key) { - if (key instanceof String) { - String stringKey = (String) key; - String keyWithPrefix; - if (prefix.isEmpty()) { - keyWithPrefix = stringKey; - } else { - keyWithPrefix = prefix + stringKey; - } - ignore(keyWithPrefix); - if (withIgnoreFallback) - ignore(stringKey); - } - return super.get(key); - } - - public RecordingMap copy() { - return new RecordingMap<>(this); - } - - } - /** * ResolvingMap keeps a track of the original map instance and the resolved configs. * The originals are tracked in a separate nested map and may be a `RecordingMap`; thus diff --git a/clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java b/clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java new file mode 100644 index 0000000000000..549a824adaa8f --- /dev/null +++ b/clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java @@ -0,0 +1,84 @@ +/* + * 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.common.config.internals; + +import org.apache.kafka.common.config.AbstractConfig; + +import java.util.HashMap; +import java.util.Map; +import java.util.Objects; + +/** + * Marks keys retrieved via `get` as used. This is needed because `Configurable.configure` takes a `Map` instead + * of an `AbstractConfig` and we can't change that without breaking public API like `Partitioner`. + */ +public class RecordingMap extends HashMap { + + private static final long serialVersionUID = 8756361579758562482L; + + private final String prefix; + private final boolean withIgnoreFallback; + private final transient AbstractConfig config; + + public RecordingMap(AbstractConfig config) { + this(config, "", false); + } + + public RecordingMap(AbstractConfig config, String prefix, boolean withIgnoreFallback) { + this.config = config; + this.prefix = Objects.requireNonNull(prefix); + this.withIgnoreFallback = withIgnoreFallback; + } + + public RecordingMap(AbstractConfig config, Map m) { + this(config, m, "", false); + } + + public RecordingMap(AbstractConfig config, Map m, String prefix, boolean withIgnoreFallback) { + super(m); + this.config = config; + this.prefix = Objects.requireNonNull(prefix); + this.withIgnoreFallback = withIgnoreFallback; + } + + @Override + public V get(Object key) { + if (key instanceof String) { + String stringKey = (String) key; + String keyWithPrefix; + if (prefix.isEmpty()) { + keyWithPrefix = stringKey; + } else { + keyWithPrefix = prefix + stringKey; + } + record(keyWithPrefix); + if (withIgnoreFallback) + record(stringKey); + } + return super.get(key); + } + + public RecordingMap copy() { + return new RecordingMap<>(config, this); + } + + private void record(String key) { + if (config != null) + config.ignore(key); + } + +} From d35e8668d3f46280043673fa321e2afd837addc4 Mon Sep 17 00:00:00 2001 From: Chris Egerton Date: Fri, 8 Apr 2022 14:10:23 -0400 Subject: [PATCH 3/3] KAFKA-7509: Minor improvement to RecordingMap copy API --- .../kafka/clients/consumer/ConsumerConfig.java | 5 +---- .../kafka/clients/producer/ProducerConfig.java | 5 +---- .../kafka/common/config/internals/RecordingMap.java | 13 +++++++++---- 3 files changed, 11 insertions(+), 12 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java index d99d0b9d69d82..116fc3635e673 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/ConsumerConfig.java @@ -32,7 +32,6 @@ import java.util.Arrays; import java.util.Collections; -import java.util.HashMap; import java.util.List; import java.util.Locale; import java.util.Map; @@ -604,9 +603,7 @@ private void maybeOverrideClientId(Map configs) { protected static Map appendDeserializerToConfig(Map configs, Deserializer keyDeserializer, Deserializer valueDeserializer) { - Map newConfigs = configs instanceof RecordingMap ? - ((RecordingMap) configs).copy() : - new HashMap<>(configs); + Map newConfigs = RecordingMap.copyAndPreserve(configs); if (keyDeserializer != null) newConfigs.put(KEY_DESERIALIZER_CLASS_CONFIG, keyDeserializer.getClass()); if (valueDeserializer != null) diff --git a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java index 5189de3b0f442..9ac4fa186aaf6 100644 --- a/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java +++ b/clients/src/main/java/org/apache/kafka/clients/producer/ProducerConfig.java @@ -32,7 +32,6 @@ import org.slf4j.LoggerFactory; import java.util.Collections; -import java.util.HashMap; import java.util.Map; import java.util.Properties; import java.util.Set; @@ -540,9 +539,7 @@ private static String parseAcks(String acksString) { static Map appendSerializerToConfig(Map configs, Serializer keySerializer, Serializer valueSerializer) { - Map newConfigs = configs instanceof RecordingMap ? - ((RecordingMap) configs).copy() : - new HashMap<>(configs); + Map newConfigs = RecordingMap.copyAndPreserve(configs); if (keySerializer != null) newConfigs.put(KEY_SERIALIZER_CLASS_CONFIG, keySerializer.getClass()); if (valueSerializer != null) diff --git a/clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java b/clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java index 549a824adaa8f..b2e6ce6c255b9 100644 --- a/clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java +++ b/clients/src/main/java/org/apache/kafka/common/config/internals/RecordingMap.java @@ -55,6 +55,15 @@ public RecordingMap(AbstractConfig config, Map m, String pr this.withIgnoreFallback = withIgnoreFallback; } + public static Map copyAndPreserve(Map map) { + if (map instanceof RecordingMap) { + RecordingMap recordingMap = (RecordingMap) map; + return new RecordingMap<>(recordingMap.config, recordingMap, recordingMap.prefix, recordingMap.withIgnoreFallback); + } else { + return new HashMap<>(map); + } + } + @Override public V get(Object key) { if (key instanceof String) { @@ -72,10 +81,6 @@ public V get(Object key) { return super.get(key); } - public RecordingMap copy() { - return new RecordingMap<>(config, this); - } - private void record(String key) { if (config != null) config.ignore(key);