Skip to content
Closed
Show file tree
Hide file tree
Changes from 1 commit
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -49,6 +49,7 @@
import org.apache.kafka.connect.storage.OffsetStorageReader;
import org.apache.kafka.connect.storage.OffsetStorageReaderImpl;
import org.apache.kafka.connect.storage.OffsetStorageWriter;
import org.apache.kafka.connect.util.ConnectUtils;
import org.apache.kafka.connect.util.ConnectorTaskId;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -141,6 +142,9 @@ public Worker(
producerProps.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE));
// User-specified overrides
producerProps.putAll(config.originalsWithPrefix("producer."));
// Prevent logging unused config warnings
ConnectUtils.retainConfigs(producerProps, ProducerConfig.configNames());

}

private WorkerConfigTransformer initConfigTransformer() {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -472,6 +472,8 @@ private KafkaConsumer<byte[], byte[]> createConsumer() {

KafkaConsumer<byte[], byte[]> newConsumer;
try {
// Prevent logging unused config warnings
props = ConnectUtils.retainConfigs(props, ConsumerConfig.configNames());
newConsumer = new KafkaConsumer<>(props);
} catch (Throwable t) {
throw new ConnectException("Failed to create consumer", t);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -17,16 +17,19 @@
package org.apache.kafka.connect.runtime.errors;

import org.apache.kafka.clients.admin.AdminClient;
import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.NewTopic;
import org.apache.kafka.clients.consumer.ConsumerRecord;
import org.apache.kafka.clients.producer.KafkaProducer;
import org.apache.kafka.clients.producer.ProducerConfig;
import org.apache.kafka.clients.producer.ProducerRecord;
import org.apache.kafka.common.errors.TopicExistsException;
import org.apache.kafka.common.header.Headers;
import org.apache.kafka.common.record.RecordBatch;
import org.apache.kafka.connect.errors.ConnectException;
import org.apache.kafka.connect.runtime.SinkConnectorConfig;
import org.apache.kafka.connect.runtime.WorkerConfig;
import org.apache.kafka.connect.util.ConnectUtils;
import org.apache.kafka.connect.util.ConnectorTaskId;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
Expand Down Expand Up @@ -77,7 +80,8 @@ public static DeadLetterQueueReporter createAndSetup(WorkerConfig workerConfig,
ErrorHandlingMetrics errorHandlingMetrics) {
String topic = sinkConfig.dlqTopicName();

try (AdminClient admin = AdminClient.create(workerConfig.originals())) {
Map<String, Object> adminConfig = ConnectUtils.retainConfigs(workerConfig.originals(), AdminClientConfig.configNames());
try (AdminClient admin = AdminClient.create(adminConfig)) {
if (!admin.listTopics().names().get().contains(topic)) {
log.error("Topic {} doesn't exist. Will attempt to create topic.", topic);
NewTopic schemaTopicRequest = new NewTopic(topic, DLQ_NUM_DESIRED_PARTITIONS, sinkConfig.dlqTopicReplicationFactor());
Expand All @@ -91,6 +95,7 @@ public static DeadLetterQueueReporter createAndSetup(WorkerConfig workerConfig,
}
}

ConnectUtils.retainConfigs(producerProps, ProducerConfig.configNames());
KafkaProducer<byte[], byte[]> dlqProducer = new KafkaProducer<>(producerProps);
return new DeadLetterQueueReporter(dlqProducer, sinkConfig, id, errorHandlingMetrics);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
import org.apache.kafka.connect.runtime.distributed.ClusterConfigState;
import org.apache.kafka.connect.runtime.distributed.DistributedConfig;
import org.apache.kafka.connect.util.Callback;
import org.apache.kafka.connect.util.ConnectUtils;
import org.apache.kafka.connect.util.ConnectorTaskId;
import org.apache.kafka.connect.util.KafkaBasedLog;
import org.apache.kafka.connect.util.TopicAdmin;
Expand Down Expand Up @@ -419,6 +420,9 @@ KafkaBasedLog<String, byte[]> setupAndCreateKafkaBasedLog(String topic, final Wo
Map<String, Object> consumerProps = new HashMap<>(originals);
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName());
// Prevent logging unused config warnings
ConnectUtils.retainConfigs(consumerProps, ConsumerConfig.configNames());
ConnectUtils.retainConfigs(producerProps, ProducerConfig.configNames());

Map<String, Object> adminProps = new HashMap<>(originals);
NewTopic topicDescription = TopicAdmin.defineTopic(topic).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -28,6 +28,7 @@
import org.apache.kafka.connect.runtime.WorkerConfig;
import org.apache.kafka.connect.runtime.distributed.DistributedConfig;
import org.apache.kafka.connect.util.Callback;
import org.apache.kafka.connect.util.ConnectUtils;
import org.apache.kafka.connect.util.ConvertingFutureCallback;
import org.apache.kafka.connect.util.KafkaBasedLog;
import org.apache.kafka.connect.util.TopicAdmin;
Expand Down Expand Up @@ -72,10 +73,14 @@ public void configure(final WorkerConfig config) {
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);
// Prevent logging unused config warnings
ConnectUtils.retainConfigs(producerProps, ProducerConfig.configNames());

Map<String, Object> consumerProps = new HashMap<>(originals);
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName());
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName());
// Prevent logging unused config warnings
ConnectUtils.retainConfigs(consumerProps, ConsumerConfig.configNames());

Map<String, Object> adminProps = new HashMap<>(originals);
NewTopic topicDescription = TopicAdmin.defineTopic(topic).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -39,6 +39,7 @@
import org.apache.kafka.connect.runtime.WorkerConfig;
import org.apache.kafka.connect.runtime.distributed.DistributedConfig;
import org.apache.kafka.connect.util.Callback;
import org.apache.kafka.connect.util.ConnectUtils;
import org.apache.kafka.connect.util.ConnectorTaskId;
import org.apache.kafka.connect.util.KafkaBasedLog;
import org.apache.kafka.connect.util.Table;
Expand Down Expand Up @@ -129,10 +130,14 @@ public void configure(final WorkerConfig config) {
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
// Prevent logging unused config warnings
ConnectUtils.retainConfigs(producerProps, ProducerConfig.configNames());

Map<String, Object> consumerProps = new HashMap<>(originals);
consumerProps.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class.getName());
consumerProps.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, ByteArrayDeserializer.class.getName());
// Prevent logging unused config warnings
ConnectUtils.retainConfigs(consumerProps, ConsumerConfig.configNames());

Map<String, Object> adminProps = new HashMap<>(originals);
NewTopic topicDescription = TopicAdmin.defineTopic(topic).
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,10 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;

import java.util.Iterator;
import java.util.Map;
import java.util.Map.Entry;
import java.util.Set;
import java.util.concurrent.ExecutionException;

public final class ConnectUtils {
Expand Down Expand Up @@ -65,4 +69,23 @@ static String lookupKafkaClusterId(AdminClient adminClient) {
+ "Check worker's broker connection and security properties.", e);
}
}

/**
* Modify the supplied map of configurations to remove all configuration name-value pairs that are not included in the
* specified set of names.
*
* @param configs the map of configurations to be modified; may not be null
* @param configNames the names of the configuration properties that are to be retained; may not be null
* @return the supplied {@code configs} parameter, returned for convenience
*/
public static Map<String, Object> retainConfigs(Map<String, Object> configs, Set<String> configNames) {
Iterator<Entry<String, Object>> entryIter = configs.entrySet().iterator();
while (entryIter.hasNext()) {
Map.Entry<String, Object> entry = entryIter.next();
if (!configNames.contains(entry.getKey())) {
entryIter.remove();
Comment thread
rhauch marked this conversation as resolved.
}
}
return configs;
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -165,7 +165,8 @@ public static NewTopicBuilder defineTopic(String topicName) {
* @param adminConfig the configuration for the {@link AdminClient}
*/
public TopicAdmin(Map<String, Object> adminConfig) {
this(adminConfig, AdminClient.create(adminConfig));
// Prevent logging unused config warnings
this(adminConfig, AdminClient.create(ConnectUtils.retainConfigs(adminConfig, AdminClientConfig.configNames())));
}

// visible for testing
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -16,16 +16,21 @@
*/
package org.apache.kafka.connect.util;

import org.apache.kafka.clients.admin.AdminClientConfig;
import org.apache.kafka.clients.admin.MockAdminClient;
import org.apache.kafka.common.Node;
import org.apache.kafka.connect.errors.ConnectException;
import org.junit.Test;

import java.util.Arrays;
import java.util.HashMap;
import java.util.List;
import java.util.Map;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
import static org.junit.Assert.assertNull;
import static org.junit.Assert.assertTrue;

public class ConnectUtilsTest {

Expand Down Expand Up @@ -60,4 +65,27 @@ public void testLookupKafkaClusterIdTimeout() {
ConnectUtils.lookupKafkaClusterId(adminClient);
}

@Test
public void removeNonAdminClientConfigurations() {
Map<String, Object> configs = new HashMap<>();
configs.put(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG, "bootstrap1");
configs.put(AdminClientConfig.CLIENT_ID_CONFIG, "clientId");
configs.put(AdminClientConfig.RETRIES_CONFIG, "1");
configs.put(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG, "100");
configs.put("some.other.property", "value");
configs.put("other.property", "value2");
Map<String, Object> filtered = ConnectUtils.retainConfigs(new HashMap<>(configs), AdminClientConfig.configNames());
assertEquals(configs.get(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG),
filtered.get(AdminClientConfig.BOOTSTRAP_SERVERS_CONFIG));
assertEquals(configs.get(AdminClientConfig.CLIENT_ID_CONFIG),
filtered.get(AdminClientConfig.CLIENT_ID_CONFIG));
assertEquals(configs.get(AdminClientConfig.RETRIES_CONFIG),
filtered.get(AdminClientConfig.RETRIES_CONFIG));
assertEquals(configs.get(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG),
filtered.get(AdminClientConfig.REQUEST_TIMEOUT_MS_CONFIG));
assertFalse(filtered.containsKey("some.other.property"));
assertFalse(filtered.containsKey("other.property"));
assertEquals(configs.size() - 2, filtered.size());
assertTrue(configs.keySet().containsAll(filtered.keySet()));
}
}