From 79e14c3c8e2747fcf001401e5e964aa3e6246302 Mon Sep 17 00:00:00 2001 From: Magesh Nandakumar Date: Mon, 22 Oct 2018 15:38:20 -0700 Subject: [PATCH 1/6] Worker Task refactoring Worker Task refactoring mark the new methods private --- .../apache/kafka/connect/runtime/Worker.java | 60 ++++++++++++++----- .../kafka/connect/runtime/WorkerSinkTask.java | 30 +--------- 2 files changed, 48 insertions(+), 42 deletions(-) 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 6e021b90a07dc..94db81a776211 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 @@ -16,6 +16,8 @@ */ package org.apache.kafka.connect.runtime; +import org.apache.kafka.clients.consumer.ConsumerConfig; +import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.producer.KafkaProducer; import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.MetricName; @@ -50,6 +52,7 @@ import org.apache.kafka.connect.storage.OffsetStorageReaderImpl; import org.apache.kafka.connect.storage.OffsetStorageWriter; import org.apache.kafka.connect.util.ConnectorTaskId; +import org.apache.kafka.connect.util.SinkUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -89,7 +92,6 @@ public class Worker { private final Converter internalKeyConverter; private final Converter internalValueConverter; private final OffsetBackingStore offsetBackingStore; - private final Map producerProps; private final ConcurrentMap connectors = new ConcurrentHashMap<>(); private final ConcurrentMap tasks = new ConcurrentHashMap<>(); @@ -129,19 +131,6 @@ public Worker( this.workerConfigTransformer = initConfigTransformer(); - 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 are designed to ensure there is no data loss. They *may* be overridden via configs passed to the - // worker, but this may compromise the delivery guarantees of Kafka Connect. - producerProps.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); - producerProps.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); - 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)); - // User-specified overrides - producerProps.putAll(config.originalsWithPrefix("producer.")); } private WorkerConfigTransformer initConfigTransformer() { @@ -498,6 +487,7 @@ private WorkerTask buildWorkerTask(ClusterConfigState configState, internalKeyConverter, internalValueConverter); OffsetStorageWriter offsetWriter = new OffsetStorageWriter(offsetBackingStore, id.connector(), internalKeyConverter, internalValueConverter); + Map producerProps = producerConfigs(); KafkaProducer producer = new KafkaProducer<>(producerProps); // Note we pass the configState as it performs dynamic transformations under the covers @@ -508,15 +498,54 @@ private WorkerTask buildWorkerTask(ClusterConfigState configState, TransformationChain transformationChain = new TransformationChain<>(connConfig.transformations(), retryWithToleranceOperator); SinkConnectorConfig sinkConfig = new SinkConnectorConfig(plugins, connConfig.originalsStrings()); retryWithToleranceOperator.reporters(sinkTaskReporters(id, sinkConfig, errorHandlingMetrics)); + + Map props = consumerConfig(connConfig, id); + KafkaConsumer consumer = new KafkaConsumer<>(props); + return new WorkerSinkTask(id, (SinkTask) task, statusListener, initialState, config, configState, metrics, keyConverter, valueConverter, headerConverter, transformationChain, loader, time, - retryWithToleranceOperator); + retryWithToleranceOperator, consumer); } else { log.error("Tasks must be a subclass of either SourceTask or SinkTask", task); throw new ConnectException("Tasks must be a subclass of either SourceTask or SinkTask"); } } + private Map producerConfigs() { + 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 are designed to ensure there is no data loss. They *may* be overridden via configs passed to the + // worker, but this may compromise the delivery guarantees of Kafka Connect. + producerProps.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + producerProps.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); + 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)); + // User-specified overrides + producerProps.putAll(config.originalsWithPrefix("producer.")); + return producerProps; + } + + + private Map consumerConfig(ConnectorConfig connConfig, ConnectorTaskId id) { + // Include any unknown worker configs so consumer configs can be set globally on the worker + // and through to the task + Map props = new HashMap<>(); + + props.put(ConsumerConfig.GROUP_ID_CONFIG, SinkUtils.consumerGroupId(id.connector())); + props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, + Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); + props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); + props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + + props.putAll(config.originalsWithPrefix("consumer.")); + return props; + } + ErrorHandlingMetrics errorHandlingMetrics(ConnectorTaskId id) { return new ErrorHandlingMetrics(id, metrics); } @@ -530,6 +559,7 @@ private List sinkTaskReporters(ConnectorTaskId id, SinkConnectorC // check if topic for dead letter queue exists String topic = connConfig.dlqTopicName(); if (topic != null && !topic.isEmpty()) { + Map producerProps = producerConfigs(); DeadLetterQueueReporter reporter = DeadLetterQueueReporter.createAndSetup(config, id, connConfig, producerProps, errorHandlingMetrics); reporters.add(reporter); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java index 39e0c6d53f645..bf3b911de87b6 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java @@ -107,7 +107,8 @@ public WorkerSinkTask(ConnectorTaskId id, TransformationChain transformationChain, ClassLoader loader, Time time, - RetryWithToleranceOperator retryWithToleranceOperator) { + RetryWithToleranceOperator retryWithToleranceOperator, + KafkaConsumer consumer) { super(id, statusListener, initialState, loader, connectMetrics, retryWithToleranceOperator); this.workerConfig = workerConfig; @@ -131,13 +132,13 @@ public WorkerSinkTask(ConnectorTaskId id, this.commitFailures = 0; this.sinkTaskMetricsGroup = new SinkTaskMetricsGroup(id, connectMetrics); this.sinkTaskMetricsGroup.recordOffsetSequenceNumber(commitSeqno); + this.consumer = consumer; } @Override public void initialize(TaskConfig taskConfig) { try { this.taskConfig = taskConfig.originalsStrings(); - this.consumer = createConsumer(); this.context = new WorkerSinkTaskContext(consumer, this, configState); } catch (Throwable t) { log.error("{} Task failed initialization and will not be started.", this, t); @@ -455,31 +456,6 @@ private ConsumerRecords pollConsumer(long timeoutMs) { return msgs; } - private KafkaConsumer createConsumer() { - // Include any unknown worker configs so consumer configs can be set globally on the worker - // and through to the task - Map props = new HashMap<>(); - - props.put(ConsumerConfig.GROUP_ID_CONFIG, SinkUtils.consumerGroupId(id.connector())); - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, - Utils.join(workerConfig.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); - props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); - props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - - props.putAll(workerConfig.originalsWithPrefix("consumer.")); - - KafkaConsumer newConsumer; - try { - newConsumer = new KafkaConsumer<>(props); - } catch (Throwable t) { - throw new ConnectException("Failed to create consumer", t); - } - - return newConsumer; - } - private void convertMessages(ConsumerRecords msgs) { origOffsets.clear(); for (ConsumerRecord msg : msgs) { From de7adece4bbccaf208829a832d0a63ea13264126 Mon Sep 17 00:00:00 2001 From: Magesh Nandakumar Date: Thu, 25 Oct 2018 11:17:48 -0700 Subject: [PATCH 2/6] Remove unused imports --- .../java/org/apache/kafka/connect/runtime/WorkerSinkTask.java | 3 --- 1 file changed, 3 deletions(-) diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java index bf3b911de87b6..14bd387504077 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java @@ -16,7 +16,6 @@ */ package org.apache.kafka.connect.runtime; -import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.consumer.ConsumerRebalanceListener; import org.apache.kafka.clients.consumer.ConsumerRecord; import org.apache.kafka.clients.consumer.ConsumerRecords; @@ -33,7 +32,6 @@ import org.apache.kafka.common.metrics.stats.Total; import org.apache.kafka.common.metrics.stats.Value; import org.apache.kafka.common.utils.Time; -import org.apache.kafka.common.utils.Utils; import org.apache.kafka.connect.data.SchemaAndValue; import org.apache.kafka.connect.errors.ConnectException; import org.apache.kafka.connect.errors.RetriableException; @@ -49,7 +47,6 @@ import org.apache.kafka.connect.storage.HeaderConverter; import org.apache.kafka.connect.util.ConnectUtils; import org.apache.kafka.connect.util.ConnectorTaskId; -import org.apache.kafka.connect.util.SinkUtils; import org.slf4j.Logger; import org.slf4j.LoggerFactory; From 57dc817fce2e96fb8ab15d88cdac4c6dabd18994 Mon Sep 17 00:00:00 2001 From: Magesh Nandakumar Date: Thu, 1 Nov 2018 15:26:49 -0700 Subject: [PATCH 3/6] Fix review comments --- .../apache/kafka/connect/runtime/Worker.java | 34 +++---- .../kafka/connect/runtime/WorkerSinkTask.java | 4 +- .../runtime/ErrorHandlingTaskTest.java | 10 +- .../connect/runtime/WorkerSinkTaskTest.java | 13 +-- .../runtime/WorkerSinkTaskThreadedTest.java | 6 +- .../kafka/connect/runtime/WorkerTest.java | 93 ++++++++++++++++++- 6 files changed, 118 insertions(+), 42 deletions(-) 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 94db81a776211..111dc477efedf 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 @@ -487,7 +487,7 @@ private WorkerTask buildWorkerTask(ClusterConfigState configState, internalKeyConverter, internalValueConverter); OffsetStorageWriter offsetWriter = new OffsetStorageWriter(offsetBackingStore, id.connector(), internalKeyConverter, internalValueConverter); - Map producerProps = producerConfigs(); + Map producerProps = producerConfigs(connConfig, id); KafkaProducer producer = new KafkaProducer<>(producerProps); // Note we pass the configState as it performs dynamic transformations under the covers @@ -499,19 +499,19 @@ private WorkerTask buildWorkerTask(ClusterConfigState configState, SinkConnectorConfig sinkConfig = new SinkConnectorConfig(plugins, connConfig.originalsStrings()); retryWithToleranceOperator.reporters(sinkTaskReporters(id, sinkConfig, errorHandlingMetrics)); - Map props = consumerConfig(connConfig, id); - KafkaConsumer consumer = new KafkaConsumer<>(props); + Map consumerProps = consumerConfigs(connConfig, id); + KafkaConsumer consumer = new KafkaConsumer<>(consumerProps); return new WorkerSinkTask(id, (SinkTask) task, statusListener, initialState, config, configState, metrics, keyConverter, - valueConverter, headerConverter, transformationChain, loader, time, - retryWithToleranceOperator, consumer); + valueConverter, headerConverter, transformationChain, consumer, loader, time, + retryWithToleranceOperator); } else { log.error("Tasks must be a subclass of either SourceTask or SinkTask", task); throw new ConnectException("Tasks must be a subclass of either SourceTask or SinkTask"); } } - private Map producerConfigs() { + Map producerConfigs(ConnectorConfig connConfig, ConnectorTaskId id) { 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"); @@ -529,21 +529,21 @@ private Map producerConfigs() { } - private Map consumerConfig(ConnectorConfig connConfig, ConnectorTaskId id) { + Map consumerConfigs(ConnectorConfig connConfig, ConnectorTaskId id) { // Include any unknown worker configs so consumer configs can be set globally on the worker // and through to the task - Map props = new HashMap<>(); + Map consumerProps = new HashMap<>(); - props.put(ConsumerConfig.GROUP_ID_CONFIG, SinkUtils.consumerGroupId(id.connector())); - props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, + consumerProps.put(ConsumerConfig.GROUP_ID_CONFIG, SinkUtils.consumerGroupId(id.connector())); + consumerProps.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, Utils.join(config.getList(WorkerConfig.BOOTSTRAP_SERVERS_CONFIG), ",")); - props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); - props.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + 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"); - props.putAll(config.originalsWithPrefix("consumer.")); - return props; + consumerProps.putAll(config.originalsWithPrefix("consumer.")); + return consumerProps; } ErrorHandlingMetrics errorHandlingMetrics(ConnectorTaskId id) { @@ -559,7 +559,7 @@ private List sinkTaskReporters(ConnectorTaskId id, SinkConnectorC // check if topic for dead letter queue exists String topic = connConfig.dlqTopicName(); if (topic != null && !topic.isEmpty()) { - Map producerProps = producerConfigs(); + Map producerProps = producerConfigs(connConfig, id); DeadLetterQueueReporter reporter = DeadLetterQueueReporter.createAndSetup(config, id, connConfig, producerProps, errorHandlingMetrics); reporters.add(reporter); } diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java index 14bd387504077..a112bfa41aefa 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/runtime/WorkerSinkTask.java @@ -102,10 +102,10 @@ public WorkerSinkTask(ConnectorTaskId id, Converter valueConverter, HeaderConverter headerConverter, TransformationChain transformationChain, + KafkaConsumer consumer, ClassLoader loader, Time time, - RetryWithToleranceOperator retryWithToleranceOperator, - KafkaConsumer consumer) { + RetryWithToleranceOperator retryWithToleranceOperator) { super(id, statusListener, initialState, loader, connectMetrics, retryWithToleranceOperator); this.workerConfig = workerConfig; diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java index 6d92c34adef0d..5d223f4713ec0 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/ErrorHandlingTaskTest.java @@ -352,7 +352,6 @@ private void assertErrorHandlingMetricValue(String name, double expected) { } private void expectInitializeTask() throws Exception { - PowerMock.expectPrivate(workerSinkTask, "createConsumer").andReturn(consumer); consumer.subscribe(EasyMock.eq(singletonList(TOPIC)), EasyMock.capture(rebalanceListener)); PowerMock.expectLastCall(); @@ -371,11 +370,10 @@ private void createSinkTask(TargetState initialState, RetryWithToleranceOperator TransformationChain sinkTransforms = new TransformationChain<>(singletonList(new FaultyPassthrough()), retryWithToleranceOperator); - workerSinkTask = PowerMock.createPartialMock( - WorkerSinkTask.class, new String[]{"createConsumer"}, - taskId, sinkTask, statusListener, initialState, workerConfig, - ClusterConfigState.EMPTY, metrics, converter, converter, - headerConverter, sinkTransforms, pluginLoader, time, retryWithToleranceOperator); + workerSinkTask = new WorkerSinkTask( + taskId, sinkTask, statusListener, initialState, workerConfig, + ClusterConfigState.EMPTY, metrics, converter, converter, + headerConverter, sinkTransforms, consumer, pluginLoader, time, retryWithToleranceOperator); } private void createSourceTask(TargetState initialState, RetryWithToleranceOperator retryWithToleranceOperator) { diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java index 7223c3bce228d..3e047ff9ed5ce 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskTest.java @@ -162,12 +162,11 @@ public void setUp() { } private void createTask(TargetState initialState) { - workerTask = PowerMock.createPartialMock( - WorkerSinkTask.class, new String[]{"createConsumer"}, - taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, - keyConverter, valueConverter, headerConverter, - transformationChain, pluginLoader, time, - RetryWithToleranceOperatorTest.NOOP_OPERATOR); + workerTask = new WorkerSinkTask( + taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, + keyConverter, valueConverter, headerConverter, + transformationChain, consumer, pluginLoader, time, + RetryWithToleranceOperatorTest.NOOP_OPERATOR); } @After @@ -1167,7 +1166,6 @@ public void testTopicsRegex() throws Exception { createTask(TargetState.PAUSED); - PowerMock.expectPrivate(workerTask, "createConsumer").andReturn(consumer); consumer.subscribe(EasyMock.capture(topicsRegex), EasyMock.capture(rebalanceListener)); PowerMock.expectLastCall(); @@ -1255,7 +1253,6 @@ public void testMetricsGroup() { } private void expectInitializeTask() throws Exception { - PowerMock.expectPrivate(workerTask, "createConsumer").andReturn(consumer); consumer.subscribe(EasyMock.eq(asList(TOPIC)), EasyMock.capture(rebalanceListener)); PowerMock.expectLastCall(); diff --git a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java index d49c1cd99ff93..6e2b01ce7bc19 100644 --- a/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java +++ b/connect/runtime/src/test/java/org/apache/kafka/connect/runtime/WorkerSinkTaskThreadedTest.java @@ -137,12 +137,11 @@ public void setup() { workerProps.put("offset.storage.file.filename", "/tmp/connect.offsets"); pluginLoader = PowerMock.createMock(PluginClassLoader.class); workerConfig = new StandaloneConfig(workerProps); - workerTask = PowerMock.createPartialMock( - WorkerSinkTask.class, new String[]{"createConsumer"}, + workerTask = new WorkerSinkTask( taskId, sinkTask, statusListener, initialState, workerConfig, ClusterConfigState.EMPTY, metrics, keyConverter, valueConverter, headerConverter, new TransformationChain<>(Collections.emptyList(), RetryWithToleranceOperatorTest.NOOP_OPERATOR), - pluginLoader, time, RetryWithToleranceOperatorTest.NOOP_OPERATOR); + consumer, pluginLoader, time, RetryWithToleranceOperatorTest.NOOP_OPERATOR); recordsReturned = 0; } @@ -509,7 +508,6 @@ public Object answer() throws Throwable { } private void expectInitializeTask() throws Exception { - PowerMock.expectPrivate(workerTask, "createConsumer").andReturn(consumer); consumer.subscribe(EasyMock.eq(Arrays.asList(TOPIC)), EasyMock.capture(rebalanceListener)); PowerMock.expectLastCall(); 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 f3cacc469daf7..88dd7566180ac 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 @@ -17,7 +17,9 @@ package org.apache.kafka.connect.runtime; import org.apache.kafka.clients.CommonClientConfigs; +import org.apache.kafka.clients.consumer.ConsumerConfig; import org.apache.kafka.clients.producer.KafkaProducer; +import org.apache.kafka.clients.producer.ProducerConfig; import org.apache.kafka.common.Configurable; import org.apache.kafka.common.config.AbstractConfig; import org.apache.kafka.common.config.ConfigDef; @@ -90,6 +92,33 @@ public class WorkerTest extends ThreadedTest { private WorkerConfig config; private Worker worker; + private final Map defaultProducerConfigs = defaultProducerConfigs(); + + private Map defaultProducerConfigs() { + Map prodConfigs = new HashMap<>(); + prodConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); + prodConfigs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + prodConfigs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + prodConfigs.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + prodConfigs.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); + prodConfigs.put(ProducerConfig.ACKS_CONFIG, "all"); + prodConfigs.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "1"); + prodConfigs.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + return prodConfigs; + } + + private final Map defaultConsumerConfigs = defaultConsumerConfigs(); + + private Map defaultConsumerConfigs() { + Map consumerConfigs = new HashMap<>(); + consumerConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); + consumerConfigs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); + consumerConfigs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + consumerConfigs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + consumerConfigs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + return consumerConfigs; + } + @Mock private Plugins plugins; @Mock @@ -804,6 +833,63 @@ public void testConverterOverrides() throws Exception { PowerMock.verifyAll(); } + @Test + public void testProducerConfigsWithoutOverrides() { + expectConverters(); + expectStartStorage(); + PowerMock.replayAll(); + worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore); + assertEquals(defaultProducerConfigs, worker.producerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + } + + @Test + public void testProducerConfigsWithOverrides() { + Map workerProps = config.originals(); + workerProps.put("producer.acks", "-1"); + workerProps.put("producer.linger.ms", "1000"); + WorkerConfig configWithOverrides = new StandaloneConfig(workerProps); + + expectConverters(JsonConverter.class, false, configWithOverrides); + expectStartStorage(); + PowerMock.replayAll(); + + worker = new Worker(WORKER_ID, new MockTime(), plugins, configWithOverrides, offsetBackingStore); + + Map expectedConfigs = new HashMap<>(defaultProducerConfigs); + expectedConfigs.put("acks", "-1"); + expectedConfigs.put("linger.ms", "1000"); + assertEquals(expectedConfigs, worker.producerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + } + + @Test + public void testConsumerConfigsWithoutOverrides() { + expectConverters(); + expectStartStorage(); + PowerMock.replayAll(); + worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore); + Map expectedConfigs = new HashMap<>(defaultConsumerConfigs); + expectedConfigs.put("group.id", "connect-test"); + assertEquals(expectedConfigs, worker.consumerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + } + + @Test + public void testConsumerConfigsWithOverrides() { + Map workerProps = config.originals(); + workerProps.put("consumer.auto.offset.reset", "latest"); + workerProps.put("consumer.max.poll.records", "1000"); + WorkerConfig configWithOverrides = new StandaloneConfig(workerProps); + + expectConverters(JsonConverter.class, false, configWithOverrides); + expectStartStorage(); + PowerMock.replayAll(); + worker = new Worker(WORKER_ID, new MockTime(), plugins, configWithOverrides, offsetBackingStore); + Map expectedConfigs = new HashMap<>(defaultConsumerConfigs); + expectedConfigs.put("group.id", "connect-test"); + expectedConfigs.put("auto.offset.reset", "latest"); + expectedConfigs.put("max.poll.records", "1000"); + assertEquals(expectedConfigs, worker.consumerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + } + private void assertStatistics(Worker worker, int connectors, int tasks) { MetricGroup workerMetrics = worker.workerMetricsGroup().metricGroup(); assertEquals(connectors, MockConnectMetrics.currentMetricValueAsDouble(worker.metrics(), workerMetrics, "connector-count"), 0.0001d); @@ -852,15 +938,12 @@ private void expectStopStorage() { } private void expectConverters() { - expectConverters(JsonConverter.class, false); + expectConverters(JsonConverter.class, false, config); } - private void expectConverters(Boolean expectDefaultConverters) { - expectConverters(JsonConverter.class, expectDefaultConverters); - } @SuppressWarnings("deprecation") - private void expectConverters(Class converterClass, Boolean expectDefaultConverters) { + private void expectConverters(Class converterClass, Boolean expectDefaultConverters, WorkerConfig config) { // As default converters are instantiated when a task starts, they are expected only if the `startTask` method is called if (expectDefaultConverters) { From 0467de3e88ebc7548ffa09683fa5b6b28ec3b6f8 Mon Sep 17 00:00:00 2001 From: Magesh Nandakumar Date: Tue, 27 Nov 2018 23:14:37 -0800 Subject: [PATCH 4/6] Fix review comments --- .../apache/kafka/connect/runtime/Worker.java | 10 +-- .../kafka/connect/runtime/WorkerTest.java | 86 ++++++++----------- 2 files changed, 39 insertions(+), 57 deletions(-) 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 111dc477efedf..49ba4fdc38b93 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 @@ -487,7 +487,7 @@ private WorkerTask buildWorkerTask(ClusterConfigState configState, internalKeyConverter, internalValueConverter); OffsetStorageWriter offsetWriter = new OffsetStorageWriter(offsetBackingStore, id.connector(), internalKeyConverter, internalValueConverter); - Map producerProps = producerConfigs(connConfig, id); + Map producerProps = producerConfigs(config); KafkaProducer producer = new KafkaProducer<>(producerProps); // Note we pass the configState as it performs dynamic transformations under the covers @@ -499,7 +499,7 @@ private WorkerTask buildWorkerTask(ClusterConfigState configState, SinkConnectorConfig sinkConfig = new SinkConnectorConfig(plugins, connConfig.originalsStrings()); retryWithToleranceOperator.reporters(sinkTaskReporters(id, sinkConfig, errorHandlingMetrics)); - Map consumerProps = consumerConfigs(connConfig, id); + Map consumerProps = consumerConfigs(id, config); KafkaConsumer consumer = new KafkaConsumer<>(consumerProps); return new WorkerSinkTask(id, (SinkTask) task, statusListener, initialState, config, configState, metrics, keyConverter, @@ -511,7 +511,7 @@ private WorkerTask buildWorkerTask(ClusterConfigState configState, } } - Map producerConfigs(ConnectorConfig connConfig, ConnectorTaskId id) { + static Map producerConfigs(WorkerConfig config) { 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"); @@ -529,7 +529,7 @@ Map producerConfigs(ConnectorConfig connConfig, ConnectorTaskId } - Map consumerConfigs(ConnectorConfig connConfig, ConnectorTaskId id) { + static Map consumerConfigs(ConnectorTaskId id, WorkerConfig config) { // Include any unknown worker configs so consumer configs can be set globally on the worker // and through to the task Map consumerProps = new HashMap<>(); @@ -559,7 +559,7 @@ private List sinkTaskReporters(ConnectorTaskId id, SinkConnectorC // check if topic for dead letter queue exists String topic = connConfig.dlqTopicName(); if (topic != null && !topic.isEmpty()) { - Map producerProps = producerConfigs(connConfig, id); + Map producerProps = producerConfigs(config); DeadLetterQueueReporter reporter = DeadLetterQueueReporter.createAndSetup(config, id, connConfig, producerProps, errorHandlingMetrics); reporters.add(reporter); } 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 88dd7566180ac..daa6665faefe7 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 @@ -92,31 +92,28 @@ public class WorkerTest extends ThreadedTest { private WorkerConfig config; private Worker worker; - private final Map defaultProducerConfigs = defaultProducerConfigs(); - - private Map defaultProducerConfigs() { - Map prodConfigs = new HashMap<>(); - prodConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); - prodConfigs.put(ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); - prodConfigs.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); - prodConfigs.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); - prodConfigs.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); - prodConfigs.put(ProducerConfig.ACKS_CONFIG, "all"); - prodConfigs.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "1"); - prodConfigs.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); - return prodConfigs; - } - - private final Map defaultConsumerConfigs = defaultConsumerConfigs(); - - private Map defaultConsumerConfigs() { - Map consumerConfigs = new HashMap<>(); - consumerConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); - consumerConfigs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); - consumerConfigs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - consumerConfigs.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - consumerConfigs.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - return consumerConfigs; + private static final Map DEFAULT_PRODUCER_CONFIGS = new HashMap<>(); + private static final Map DEFAULT_CONSUMER_CONFIGS = new HashMap<>(); + + static { + DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); + DEFAULT_PRODUCER_CONFIGS.put( + ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + DEFAULT_PRODUCER_CONFIGS.put( + ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); + DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.ACKS_CONFIG, "all"); + DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "1"); + DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + + DEFAULT_CONSUMER_CONFIGS.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); + DEFAULT_CONSUMER_CONFIGS.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); + DEFAULT_CONSUMER_CONFIGS.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + DEFAULT_CONSUMER_CONFIGS + .put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + DEFAULT_CONSUMER_CONFIGS + .put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); } @Mock @@ -835,11 +832,7 @@ public void testConverterOverrides() throws Exception { @Test public void testProducerConfigsWithoutOverrides() { - expectConverters(); - expectStartStorage(); - PowerMock.replayAll(); - worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore); - assertEquals(defaultProducerConfigs, worker.producerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + assertEquals(DEFAULT_PRODUCER_CONFIGS, worker.producerConfigs(config)); } @Test @@ -849,27 +842,17 @@ public void testProducerConfigsWithOverrides() { workerProps.put("producer.linger.ms", "1000"); WorkerConfig configWithOverrides = new StandaloneConfig(workerProps); - expectConverters(JsonConverter.class, false, configWithOverrides); - expectStartStorage(); - PowerMock.replayAll(); - - worker = new Worker(WORKER_ID, new MockTime(), plugins, configWithOverrides, offsetBackingStore); - - Map expectedConfigs = new HashMap<>(defaultProducerConfigs); + Map expectedConfigs = new HashMap<>(DEFAULT_PRODUCER_CONFIGS); expectedConfigs.put("acks", "-1"); expectedConfigs.put("linger.ms", "1000"); - assertEquals(expectedConfigs, worker.producerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + assertEquals(expectedConfigs, worker.producerConfigs(configWithOverrides)); } @Test public void testConsumerConfigsWithoutOverrides() { - expectConverters(); - expectStartStorage(); - PowerMock.replayAll(); - worker = new Worker(WORKER_ID, new MockTime(), plugins, config, offsetBackingStore); - Map expectedConfigs = new HashMap<>(defaultConsumerConfigs); + Map expectedConfigs = new HashMap<>(DEFAULT_CONSUMER_CONFIGS); expectedConfigs.put("group.id", "connect-test"); - assertEquals(expectedConfigs, worker.consumerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + assertEquals(expectedConfigs, worker.consumerConfigs(new ConnectorTaskId("test", 1), config)); } @Test @@ -879,15 +862,11 @@ public void testConsumerConfigsWithOverrides() { workerProps.put("consumer.max.poll.records", "1000"); WorkerConfig configWithOverrides = new StandaloneConfig(workerProps); - expectConverters(JsonConverter.class, false, configWithOverrides); - expectStartStorage(); - PowerMock.replayAll(); - worker = new Worker(WORKER_ID, new MockTime(), plugins, configWithOverrides, offsetBackingStore); - Map expectedConfigs = new HashMap<>(defaultConsumerConfigs); + Map expectedConfigs = new HashMap<>(DEFAULT_CONSUMER_CONFIGS); expectedConfigs.put("group.id", "connect-test"); expectedConfigs.put("auto.offset.reset", "latest"); expectedConfigs.put("max.poll.records", "1000"); - assertEquals(expectedConfigs, worker.consumerConfigs(EasyMock.mock(ConnectorConfig.class), new ConnectorTaskId("test", 1))); + assertEquals(expectedConfigs, worker.consumerConfigs(new ConnectorTaskId("test", 1), configWithOverrides)); } private void assertStatistics(Worker worker, int connectors, int tasks) { @@ -938,12 +917,15 @@ private void expectStopStorage() { } private void expectConverters() { - expectConverters(JsonConverter.class, false, config); + expectConverters(JsonConverter.class, false); } + private void expectConverters(Boolean expectDefaultConverters) { + expectConverters(JsonConverter.class, expectDefaultConverters); + } @SuppressWarnings("deprecation") - private void expectConverters(Class converterClass, Boolean expectDefaultConverters, WorkerConfig config) { + private void expectConverters(Class converterClass, Boolean expectDefaultConverters) { // As default converters are instantiated when a task starts, they are expected only if the `startTask` method is called if (expectDefaultConverters) { From 992226994cd285ca69a925d9466bfa6476e6cca0 Mon Sep 17 00:00:00 2001 From: Magesh Nandakumar Date: Wed, 28 Nov 2018 09:29:22 -0800 Subject: [PATCH 5/6] Fix build failures --- .../kafka/connect/runtime/WorkerTest.java | 77 +++++++++---------- 1 file changed, 37 insertions(+), 40 deletions(-) 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 daa6665faefe7..6482f07d95716 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 @@ -89,32 +89,12 @@ public class WorkerTest extends ThreadedTest { private static final ConnectorTaskId TASK_ID = new ConnectorTaskId("job", 0); private static final String WORKER_ID = "localhost:8083"; + private Map workerProps = new HashMap<>(); private WorkerConfig config; private Worker worker; - private static final Map DEFAULT_PRODUCER_CONFIGS = new HashMap<>(); - private static final Map DEFAULT_CONSUMER_CONFIGS = new HashMap<>(); - - static { - DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); - DEFAULT_PRODUCER_CONFIGS.put( - ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); - DEFAULT_PRODUCER_CONFIGS.put( - ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); - DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); - DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); - DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.ACKS_CONFIG, "all"); - DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "1"); - DEFAULT_PRODUCER_CONFIGS.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); - - DEFAULT_CONSUMER_CONFIGS.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); - DEFAULT_CONSUMER_CONFIGS.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); - DEFAULT_CONSUMER_CONFIGS.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); - DEFAULT_CONSUMER_CONFIGS - .put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - DEFAULT_CONSUMER_CONFIGS - .put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); - } + private Map defaultProducerConfigs = new HashMap<>(); + private Map defaultConsumerConfigs = new HashMap<>(); @Mock private Plugins plugins; @@ -142,8 +122,6 @@ public class WorkerTest extends ThreadedTest { @Before 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("internal.key.converter", "org.apache.kafka.connect.json.JsonConverter"); @@ -154,6 +132,25 @@ public void setup() { workerProps.put(CommonClientConfigs.METRIC_REPORTER_CLASSES_CONFIG, MockMetricsReporter.class.getName()); config = new StandaloneConfig(workerProps); + defaultProducerConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); + defaultProducerConfigs.put( + ProducerConfig.KEY_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + defaultProducerConfigs.put( + ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer"); + defaultProducerConfigs.put(ProducerConfig.REQUEST_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + defaultProducerConfigs.put(ProducerConfig.MAX_BLOCK_MS_CONFIG, Long.toString(Long.MAX_VALUE)); + defaultProducerConfigs.put(ProducerConfig.ACKS_CONFIG, "all"); + defaultProducerConfigs.put(ProducerConfig.MAX_IN_FLIGHT_REQUESTS_PER_CONNECTION, "1"); + defaultProducerConfigs.put(ProducerConfig.DELIVERY_TIMEOUT_MS_CONFIG, Integer.toString(Integer.MAX_VALUE)); + + defaultConsumerConfigs.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092"); + defaultConsumerConfigs.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); + defaultConsumerConfigs.put(ConsumerConfig.AUTO_OFFSET_RESET_CONFIG, "earliest"); + defaultConsumerConfigs + .put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + defaultConsumerConfigs + .put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArrayDeserializer"); + PowerMock.mockStatic(Plugins.class); } @@ -832,41 +829,41 @@ public void testConverterOverrides() throws Exception { @Test public void testProducerConfigsWithoutOverrides() { - assertEquals(DEFAULT_PRODUCER_CONFIGS, worker.producerConfigs(config)); + assertEquals(defaultProducerConfigs, Worker.producerConfigs(config)); } @Test public void testProducerConfigsWithOverrides() { - Map workerProps = config.originals(); - workerProps.put("producer.acks", "-1"); - workerProps.put("producer.linger.ms", "1000"); - WorkerConfig configWithOverrides = new StandaloneConfig(workerProps); + Map props = new HashMap<>(workerProps); + props.put("consumer.auto.offset.reset", "latest"); + props.put("consumer.max.poll.records", "1000"); + WorkerConfig configWithOverrides = new StandaloneConfig(props); - Map expectedConfigs = new HashMap<>(DEFAULT_PRODUCER_CONFIGS); + Map expectedConfigs = new HashMap<>(defaultProducerConfigs); expectedConfigs.put("acks", "-1"); expectedConfigs.put("linger.ms", "1000"); - assertEquals(expectedConfigs, worker.producerConfigs(configWithOverrides)); + assertEquals(expectedConfigs, Worker.producerConfigs(configWithOverrides)); } @Test public void testConsumerConfigsWithoutOverrides() { - Map expectedConfigs = new HashMap<>(DEFAULT_CONSUMER_CONFIGS); + Map expectedConfigs = new HashMap<>(defaultConsumerConfigs); expectedConfigs.put("group.id", "connect-test"); - assertEquals(expectedConfigs, worker.consumerConfigs(new ConnectorTaskId("test", 1), config)); + assertEquals(expectedConfigs, Worker.consumerConfigs(new ConnectorTaskId("test", 1), config)); } @Test public void testConsumerConfigsWithOverrides() { - Map workerProps = config.originals(); - workerProps.put("consumer.auto.offset.reset", "latest"); - workerProps.put("consumer.max.poll.records", "1000"); - WorkerConfig configWithOverrides = new StandaloneConfig(workerProps); + Map props = new HashMap<>(workerProps); + props.put("consumer.auto.offset.reset", "latest"); + props.put("consumer.max.poll.records", "1000"); + WorkerConfig configWithOverrides = new StandaloneConfig(props); - Map expectedConfigs = new HashMap<>(DEFAULT_CONSUMER_CONFIGS); + Map expectedConfigs = new HashMap<>(defaultConsumerConfigs); expectedConfigs.put("group.id", "connect-test"); expectedConfigs.put("auto.offset.reset", "latest"); expectedConfigs.put("max.poll.records", "1000"); - assertEquals(expectedConfigs, worker.consumerConfigs(new ConnectorTaskId("test", 1), configWithOverrides)); + assertEquals(expectedConfigs, Worker.consumerConfigs(new ConnectorTaskId("test", 1), configWithOverrides)); } private void assertStatistics(Worker worker, int connectors, int tasks) { From 9124ecf0762c566bbdc128247665c26358af5296 Mon Sep 17 00:00:00 2001 From: Magesh Nandakumar Date: Wed, 28 Nov 2018 12:44:38 -0800 Subject: [PATCH 6/6] Fix test failures --- .../java/org/apache/kafka/connect/runtime/WorkerTest.java | 4 ++-- 1 file changed, 2 insertions(+), 2 deletions(-) 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 6482f07d95716..8f15c87c5016c 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 @@ -835,8 +835,8 @@ public void testProducerConfigsWithoutOverrides() { @Test public void testProducerConfigsWithOverrides() { Map props = new HashMap<>(workerProps); - props.put("consumer.auto.offset.reset", "latest"); - props.put("consumer.max.poll.records", "1000"); + props.put("producer.acks", "-1"); + props.put("producer.linger.ms", "1000"); WorkerConfig configWithOverrides = new StandaloneConfig(props); Map expectedConfigs = new HashMap<>(defaultProducerConfigs);