> 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 extends Connector> 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 extends Connector> 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 extends Connector> 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 extends Connector> 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));
}