From 910317a6da1db221666ec518cc0e86788447aaf4 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Mon, 2 Jul 2018 14:55:23 +0200 Subject: [PATCH 01/11] MINOR: use poll(Duration) in consumer example --- examples/src/main/java/kafka/examples/Consumer.java | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/examples/src/main/java/kafka/examples/Consumer.java b/examples/src/main/java/kafka/examples/Consumer.java index be062b309df03..26d6e23a3f8e8 100644 --- a/examples/src/main/java/kafka/examples/Consumer.java +++ b/examples/src/main/java/kafka/examples/Consumer.java @@ -22,6 +22,7 @@ import org.apache.kafka.clients.consumer.ConsumerRecords; import org.apache.kafka.clients.consumer.KafkaConsumer; +import java.time.Duration; import java.util.Collections; import java.util.Properties; @@ -47,7 +48,7 @@ public Consumer(String topic) { @Override public void doWork() { consumer.subscribe(Collections.singletonList(this.topic)); - ConsumerRecords records = consumer.poll(1000); + ConsumerRecords records = consumer.poll(Duration.ofSeconds(1)); for (ConsumerRecord record : records) { System.out.println("Received message: (" + record.key() + ", " + record.value() + ") at offset " + record.offset()); } From fe77b3cb47cdd5c863f38f4414980f60f633b9f4 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Thu, 5 Jul 2018 08:04:37 +0200 Subject: [PATCH 02/11] Replace other non-test usages as well --- core/src/main/scala/kafka/tools/ConsoleConsumer.scala | 5 +++-- core/src/main/scala/kafka/tools/ConsumerPerformance.scala | 4 ++-- core/src/main/scala/kafka/tools/EndToEndLatency.scala | 8 +++++--- core/src/main/scala/kafka/tools/MirrorMaker.scala | 3 ++- 4 files changed, 12 insertions(+), 8 deletions(-) diff --git a/core/src/main/scala/kafka/tools/ConsoleConsumer.scala b/core/src/main/scala/kafka/tools/ConsoleConsumer.scala index 7e2c5644a3b4a..97a5dd26f7719 100755 --- a/core/src/main/scala/kafka/tools/ConsoleConsumer.scala +++ b/core/src/main/scala/kafka/tools/ConsoleConsumer.scala @@ -19,6 +19,7 @@ package kafka.tools import java.io.PrintStream import java.nio.charset.StandardCharsets +import java.time.Duration import java.util.concurrent.CountDownLatch import java.util.regex.Pattern import java.util.{Collections, Locale, Properties, Random} @@ -388,7 +389,7 @@ object ConsoleConsumer extends Logging { private[tools] class ConsumerWrapper(topic: Option[String], partitionId: Option[Int], offset: Option[Long], whitelist: Option[String], consumer: Consumer[Array[Byte], Array[Byte]], val timeoutMs: Long = Long.MaxValue) { consumerInit() - var recordIter = consumer.poll(0).iterator + var recordIter = consumer.poll(Duration.ZERO).iterator def consumerInit() { (topic, partitionId, offset, whitelist) match { @@ -432,7 +433,7 @@ object ConsoleConsumer extends Logging { def receive(): ConsumerRecord[Array[Byte], Array[Byte]] = { if (!recordIter.hasNext) { - recordIter = consumer.poll(timeoutMs).iterator + recordIter = consumer.poll(Duration.ofMillis(timeoutMs)).iterator if (!recordIter.hasNext) throw new TimeoutException() } diff --git a/core/src/main/scala/kafka/tools/ConsumerPerformance.scala b/core/src/main/scala/kafka/tools/ConsumerPerformance.scala index 5af55a8d7f17b..2e7b8ddf094f1 100644 --- a/core/src/main/scala/kafka/tools/ConsumerPerformance.scala +++ b/core/src/main/scala/kafka/tools/ConsumerPerformance.scala @@ -28,8 +28,8 @@ import org.apache.kafka.common.utils.Utils import org.apache.kafka.common.{Metric, MetricName, TopicPartition} import kafka.utils.{CommandLineUtils, ToolsUtils} import java.util.{Collections, Properties, Random} - import java.text.SimpleDateFormat +import java.time.Duration import com.typesafe.scalalogging.LazyLogging @@ -127,7 +127,7 @@ object ConsumerPerformance extends LazyLogging { var currentTimeMillis = lastConsumedTime while (messagesRead < count && currentTimeMillis - lastConsumedTime <= timeout) { - val records = consumer.poll(100).asScala + val records = consumer.poll(Duration.ofMillis(100)).asScala currentTimeMillis = System.currentTimeMillis if (records.nonEmpty) lastConsumedTime = currentTimeMillis diff --git a/core/src/main/scala/kafka/tools/EndToEndLatency.scala b/core/src/main/scala/kafka/tools/EndToEndLatency.scala index 3beaf827f5959..ee92a8b6a5a88 100755 --- a/core/src/main/scala/kafka/tools/EndToEndLatency.scala +++ b/core/src/main/scala/kafka/tools/EndToEndLatency.scala @@ -18,11 +18,13 @@ package kafka.tools import java.nio.charset.StandardCharsets +import java.time.Duration import java.util.{Arrays, Collections, Properties} import kafka.utils.Exit import org.apache.kafka.clients.consumer.{ConsumerConfig, KafkaConsumer} import org.apache.kafka.clients.producer._ +import org.apache.kafka.common.TopicPartition import org.apache.kafka.common.utils.Utils import scala.collection.JavaConverters._ @@ -89,9 +91,9 @@ object EndToEndLatency { } //Ensure we are at latest offset. seekToEnd evaluates lazily, that is to say actually performs the seek only when - //a poll() or position() request is issued. Hence we need to poll after we seek to ensure we see our first write. + //a position() request is issued. Hence we need to position after we seek to ensure we see our first write. consumer.seekToEnd(Collections.emptyList()) - consumer.poll(0) + consumer.assignment().asScala.foreach(consumer.position) var totalTime = 0.0 val latencies = new Array[Long](numMessages) @@ -103,7 +105,7 @@ object EndToEndLatency { //Send message (of random bytes) synchronously then immediately poll for it producer.send(new ProducerRecord[Array[Byte], Array[Byte]](topic, message)).get() - val recordIter = consumer.poll(timeout).iterator + val recordIter = consumer.poll(Duration.ofMillis(timeout)).iterator val elapsed = System.nanoTime - begin diff --git a/core/src/main/scala/kafka/tools/MirrorMaker.scala b/core/src/main/scala/kafka/tools/MirrorMaker.scala index d7e09e4efdb09..d55d96bd65bc8 100755 --- a/core/src/main/scala/kafka/tools/MirrorMaker.scala +++ b/core/src/main/scala/kafka/tools/MirrorMaker.scala @@ -17,6 +17,7 @@ package kafka.tools +import java.time.Duration import java.util import java.util.concurrent.atomic.{AtomicBoolean, AtomicInteger} import java.util.concurrent.{CountDownLatch, TimeUnit} @@ -452,7 +453,7 @@ object MirrorMaker extends Logging with KafkaMetricsGroup { // uncommitted record since last poll. Using one second as poll's timeout ensures that // offsetCommitIntervalMs, of value greater than 1 second, does not see delays in offset // commit. - recordIter = consumer.poll(1000).iterator + recordIter = consumer.poll(Duration.ofSeconds(1)).iterator if (!recordIter.hasNext) throw new NoRecordsException } From ab31aea87b732308ffe0c54c47e3748968c0f44a Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Mon, 9 Jul 2018 15:54:23 +0200 Subject: [PATCH 03/11] Removing the rest of the deprecated non-test poll usages --- .../java/org/apache/kafka/connect/runtime/WorkerSinkTask.java | 3 ++- .../main/java/org/apache/kafka/connect/util/KafkaBasedLog.java | 3 ++- core/src/main/scala/kafka/tools/StreamsResetter.java | 2 +- .../org/apache/kafka/tools/TransactionalMessageCopier.java | 3 ++- .../main/java/org/apache/kafka/tools/VerifiableConsumer.java | 3 ++- .../org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java | 3 ++- .../org/apache/kafka/trogdor/workload/RoundTripWorker.java | 3 ++- 7 files changed, 13 insertions(+), 7 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 47f8529e2d149..692331ed13f71 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 @@ -53,6 +53,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.time.Duration; import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; @@ -441,7 +442,7 @@ public String toString() { } private ConsumerRecords pollConsumer(long timeoutMs) { - ConsumerRecords msgs = consumer.poll(timeoutMs); + ConsumerRecords msgs = consumer.poll(Duration.ofMillis(timeoutMs)); // Exceptions raised from the task during a rebalance should be rethrown to stop the worker if (rebalanceException != null) { diff --git a/connect/runtime/src/main/java/org/apache/kafka/connect/util/KafkaBasedLog.java b/connect/runtime/src/main/java/org/apache/kafka/connect/util/KafkaBasedLog.java index de1ceb3be1006..ea9b4c621f913 100644 --- a/connect/runtime/src/main/java/org/apache/kafka/connect/util/KafkaBasedLog.java +++ b/connect/runtime/src/main/java/org/apache/kafka/connect/util/KafkaBasedLog.java @@ -35,6 +35,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.time.Duration; import java.util.ArrayDeque; import java.util.ArrayList; import java.util.Iterator; @@ -253,7 +254,7 @@ private Consumer createConsumer() { private void poll(long timeoutMs) { try { - ConsumerRecords records = consumer.poll(timeoutMs); + ConsumerRecords records = consumer.poll(Duration.ofMillis(timeoutMs)); for (ConsumerRecord record : records) consumedCallback.onCompletion(null, record); } catch (WakeupException e) { diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 3c045c69eb23a..2b5fdd4d9e6d2 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -314,7 +314,7 @@ private int maybeResetInputAndSeekToEndIntermediateTopicOffsets(final Map consum try (final KafkaConsumer client = new KafkaConsumer<>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer())) { client.subscribe(topicsToSubscribe); - client.poll(1); + client.poll(java.time.Duration.ofMillis(1)); final Set partitions = client.assignment(); final Set inputTopicPartitions = new HashSet<>(); diff --git a/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java b/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java index 0d74645379ebe..27e7c7fda5d8e 100644 --- a/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java +++ b/tools/src/main/java/org/apache/kafka/tools/TransactionalMessageCopier.java @@ -35,6 +35,7 @@ import org.apache.kafka.common.errors.ProducerFencedException; import java.io.IOException; +import java.time.Duration; import java.util.HashMap; import java.util.Map; import java.util.Properties; @@ -287,7 +288,7 @@ public void run() { try { producer.beginTransaction(); while (messagesInCurrentTransaction < numMessagesForNextTransaction) { - ConsumerRecords records = consumer.poll(200L); + ConsumerRecords records = consumer.poll(Duration.ofMillis(200)); for (ConsumerRecord record : records) { producer.send(producerRecordFromConsumerRecord(outputTopic, record)); messagesInCurrentTransaction++; diff --git a/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java b/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java index cc09b2331678f..58f34718b8afd 100644 --- a/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java +++ b/tools/src/main/java/org/apache/kafka/tools/VerifiableConsumer.java @@ -47,6 +47,7 @@ import java.io.Closeable; import java.io.IOException; import java.io.PrintStream; +import java.time.Duration; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -220,7 +221,7 @@ public void run() { consumer.subscribe(Collections.singletonList(topic), this); while (!isFinished()) { - ConsumerRecords records = consumer.poll(Long.MAX_VALUE); + ConsumerRecords records = consumer.poll(Duration.ofMillis(Long.MAX_VALUE)); Map offsets = onRecordsReceived(records); if (!useAutoCommit) { diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java index 1a85296407099..c3a90e4da6a39 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/ConsumeBenchWorker.java @@ -38,6 +38,7 @@ import org.apache.kafka.trogdor.task.TaskWorker; +import java.time.Duration; import java.util.Collection; import java.util.HashSet; import java.util.Map; @@ -135,7 +136,7 @@ public Void call() throws Exception { long startBatchMs = startTimeMs; try { while (messagesConsumed < spec.maxMessages()) { - ConsumerRecords records = consumer.poll(50); + ConsumerRecords records = consumer.poll(Duration.ofMillis(50)); if (records.isEmpty()) { continue; } diff --git a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java index 570f6a11e34f1..669fafcc75ed5 100644 --- a/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java +++ b/tools/src/main/java/org/apache/kafka/trogdor/workload/RoundTripWorker.java @@ -50,6 +50,7 @@ import java.nio.ByteBuffer; import java.nio.ByteOrder; +import java.time.Duration; import java.util.ArrayList; import java.util.HashMap; import java.util.HashSet; @@ -337,7 +338,7 @@ public void run() { while (true) { try { pollInvoked++; - ConsumerRecords records = consumer.poll(50); + ConsumerRecords records = consumer.poll(Duration.ofMillis(50)); for (Iterator> iter = records.iterator(); iter.hasNext(); ) { ConsumerRecord record = iter.next(); int messageIndex = ByteBuffer.wrap(record.key()).order(ByteOrder.LITTLE_ENDIAN).getInt(); From a359c53d52b345f5b1f387462325e9495a71fda2 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Thu, 9 Aug 2018 16:37:58 +0200 Subject: [PATCH 04/11] Use empty iterator in ConsoleConsumer and manual assignment in EndToEndLatency --- .../scala/kafka/tools/ConsoleConsumer.scala | 2 +- .../scala/kafka/tools/EndToEndLatency.scala | 17 ++++++++++++----- 2 files changed, 13 insertions(+), 6 deletions(-) diff --git a/core/src/main/scala/kafka/tools/ConsoleConsumer.scala b/core/src/main/scala/kafka/tools/ConsoleConsumer.scala index 97a5dd26f7719..365652a75b5a3 100755 --- a/core/src/main/scala/kafka/tools/ConsoleConsumer.scala +++ b/core/src/main/scala/kafka/tools/ConsoleConsumer.scala @@ -389,7 +389,7 @@ object ConsoleConsumer extends Logging { private[tools] class ConsumerWrapper(topic: Option[String], partitionId: Option[Int], offset: Option[Long], whitelist: Option[String], consumer: Consumer[Array[Byte], Array[Byte]], val timeoutMs: Long = Long.MaxValue) { consumerInit() - var recordIter = consumer.poll(Duration.ZERO).iterator + var recordIter = Collections.emptyList[ConsumerRecord[Array[Byte], Array[Byte]]]().iterator() def consumerInit() { (topic, partitionId, offset, whitelist) match { diff --git a/core/src/main/scala/kafka/tools/EndToEndLatency.scala b/core/src/main/scala/kafka/tools/EndToEndLatency.scala index ee92a8b6a5a88..5a50438597b87 100755 --- a/core/src/main/scala/kafka/tools/EndToEndLatency.scala +++ b/core/src/main/scala/kafka/tools/EndToEndLatency.scala @@ -19,6 +19,7 @@ package kafka.tools import java.nio.charset.StandardCharsets import java.time.Duration +import java.util import java.util.{Arrays, Collections, Properties} import kafka.utils.Exit @@ -71,9 +72,7 @@ object EndToEndLatency { 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.put(ConsumerConfig.FETCH_MAX_WAIT_MS_CONFIG, "0") //ensure we have no temporal batching - val consumer = new KafkaConsumer[Array[Byte], Array[Byte]](consumerProps) - consumer.subscribe(Collections.singletonList(topic)) val producerProps = loadProps producerProps.put(ProducerConfig.BOOTSTRAP_SERVERS_CONFIG, brokerList) @@ -90,9 +89,17 @@ object EndToEndLatency { consumer.close() } - //Ensure we are at latest offset. seekToEnd evaluates lazily, that is to say actually performs the seek only when - //a position() request is issued. Hence we need to position after we seek to ensure we see our first write. - consumer.seekToEnd(Collections.emptyList()) + val topics = consumer.listTopics() + val tp = topics.get(topic) + if (tp == null) { + finalise() + throw new RuntimeException("The tested topic doesn't exist in the cluster") + } + val topicPartitions = tp.asScala + .map(pi => new TopicPartition(pi.topic(), pi.partition())) + .to[List].asJava + consumer.assign(topicPartitions) + consumer.seekToEnd(topicPartitions) consumer.assignment().asScala.foreach(consumer.position) var totalTime = 0.0 From 87a57b775412ea5ff72cb1b9ae0cd8d42af7dbf8 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Thu, 9 Aug 2018 18:36:33 +0200 Subject: [PATCH 05/11] Replace poll calls in WorkerSinkTaskTest --- .../kafka/connect/runtime/WorkerSinkTaskTest.java | 15 ++++++++------- 1 file changed, 8 insertions(+), 7 deletions(-) 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 4a7c760fc746a..33ab2ef06e083 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 @@ -58,6 +58,7 @@ import org.powermock.modules.junit4.PowerMockRunner; import org.powermock.reflect.Whitebox; +import java.time.Duration; import java.util.ArrayList; import java.util.Collection; import java.util.Collections; @@ -458,7 +459,7 @@ public void testWakeupInCommitSyncCausesRetry() throws Exception { sinkTask.open(partitions); EasyMock.expectLastCall(); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { @@ -893,7 +894,7 @@ public void run() { // Expect the next poll to discover and perform the rebalance, THEN complete the previous callback handler, // and then return one record for TP1 and one for TP3. final AtomicBoolean rebalanced = new AtomicBoolean(); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { @@ -1273,7 +1274,7 @@ private void expectRebalanceRevocationError(RuntimeException e) { sinkTask.preCommit(EasyMock.>anyObject()); EasyMock.expectLastCall().andReturn(Collections.emptyMap()); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { @@ -1298,7 +1299,7 @@ private void expectRebalanceAssignmentError(RuntimeException e) { sinkTask.open(partitions); EasyMock.expectLastCall().andThrow(e); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { @@ -1315,7 +1316,7 @@ private void expectPollInitialAssignment() { sinkTask.open(partitions); EasyMock.expectLastCall(); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer(new IAnswer>() { + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer(new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { rebalanceListener.getValue().onPartitionsAssigned(partitions); @@ -1332,7 +1333,7 @@ public ConsumerRecords answer() throws Throwable { private void expectConsumerWakeup() { consumer.wakeup(); EasyMock.expectLastCall(); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andThrow(new WakeupException()); + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andThrow(new WakeupException()); } private void expectConsumerPoll(final int numMessages) { @@ -1340,7 +1341,7 @@ private void expectConsumerPoll(final int numMessages) { } private void expectConsumerPoll(final int numMessages, final long timestamp, final TimestampType timestampType) { - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { From 68c6197cbad6e81d5bd61a1ab10f8a00df90b8f3 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Thu, 9 Aug 2018 18:38:20 +0200 Subject: [PATCH 06/11] Replace polls in ErrorHandlingTaskTest --- .../apache/kafka/connect/runtime/ErrorHandlingTaskTest.java | 5 +++-- 1 file changed, 3 insertions(+), 2 deletions(-) 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 1bf9c717068e3..6d92c34adef0d 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 @@ -65,6 +65,7 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.time.Duration; import java.util.Collections; import java.util.HashMap; import java.util.Map; @@ -180,8 +181,8 @@ public void testErrorHandlingInSinkTasks() throws Exception { // bad json ConsumerRecord record2 = new ConsumerRecord<>(TOPIC, PARTITION2, FIRST_OFFSET, null, "{\"a\" 10}".getBytes()); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andReturn(records(record1)); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andReturn(records(record2)); + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andReturn(records(record1)); + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andReturn(records(record2)); sinkTask.put(EasyMock.anyObject()); EasyMock.expectLastCall().times(2); From 6ad6f5ef11cc6a21d93d5e3e34ddebb622816625 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Thu, 9 Aug 2018 17:34:28 +0200 Subject: [PATCH 07/11] Manual partition assignment in StreamsResetter --- .../src/main/scala/kafka/tools/StreamsResetter.java | 13 ++++++++++--- 1 file changed, 10 insertions(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 2b5fdd4d9e6d2..96892f4a9aa28 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -33,6 +33,7 @@ import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndTimestamp; import org.apache.kafka.common.KafkaFuture; +import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.annotation.InterfaceStability; import org.apache.kafka.common.serialization.ByteArrayDeserializer; @@ -47,6 +48,7 @@ import java.text.SimpleDateFormat; import java.util.Arrays; import java.util.ArrayList; +import java.util.Collection; import java.util.Date; import java.util.HashMap; import java.util.HashSet; @@ -57,6 +59,7 @@ import java.util.Set; import java.util.concurrent.ExecutionException; import java.util.concurrent.TimeUnit; +import java.util.stream.Collectors; /** * {@link StreamsResetter} resets the processing state of a Kafka Streams application so that, for example, you can reprocess its input from scratch. @@ -313,10 +316,14 @@ private int maybeResetInputAndSeekToEndIntermediateTopicOffsets(final Map consum config.setProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); try (final KafkaConsumer client = new KafkaConsumer<>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer())) { - client.subscribe(topicsToSubscribe); - client.poll(java.time.Duration.ofMillis(1)); + Map> pi = client.listTopics(); + Collection partitions = pi.entrySet().stream() + .filter(entry -> topicsToSubscribe.contains(entry.getKey())) + .flatMap(entry -> entry.getValue().stream()) + .map(info -> new TopicPartition(info.topic(), info.partition())) + .collect(Collectors.toList()); + client.assign(partitions); - final Set partitions = client.assignment(); final Set inputTopicPartitions = new HashSet<>(); final Set intermediateTopicPartitions = new HashSet<>(); From 3482e732f938add48794951440e8a357ad4a4670 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Thu, 9 Aug 2018 18:54:31 +0200 Subject: [PATCH 08/11] Use partitionsFor in StreamsResetter --- core/src/main/scala/kafka/tools/StreamsResetter.java | 6 ++---- 1 file changed, 2 insertions(+), 4 deletions(-) diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 96892f4a9aa28..2a00126592f46 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -316,10 +316,8 @@ private int maybeResetInputAndSeekToEndIntermediateTopicOffsets(final Map consum config.setProperty(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, "false"); try (final KafkaConsumer client = new KafkaConsumer<>(config, new ByteArrayDeserializer(), new ByteArrayDeserializer())) { - Map> pi = client.listTopics(); - Collection partitions = pi.entrySet().stream() - .filter(entry -> topicsToSubscribe.contains(entry.getKey())) - .flatMap(entry -> entry.getValue().stream()) + Collection partitions = topicsToSubscribe.stream().map(client::partitionsFor) + .flatMap(Collection::stream) .map(info -> new TopicPartition(info.topic(), info.partition())) .collect(Collectors.toList()); client.assign(partitions); From 608624e88d28a1ee06fcc041b5a07cf3de3db8df Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Fri, 10 Aug 2018 08:13:27 +0200 Subject: [PATCH 09/11] Remove unused imports --- core/src/main/scala/kafka/tools/EndToEndLatency.scala | 3 +-- core/src/main/scala/kafka/tools/StreamsResetter.java | 1 - 2 files changed, 1 insertion(+), 3 deletions(-) diff --git a/core/src/main/scala/kafka/tools/EndToEndLatency.scala b/core/src/main/scala/kafka/tools/EndToEndLatency.scala index 5a50438597b87..ec56f2b6f8892 100755 --- a/core/src/main/scala/kafka/tools/EndToEndLatency.scala +++ b/core/src/main/scala/kafka/tools/EndToEndLatency.scala @@ -19,8 +19,7 @@ package kafka.tools import java.nio.charset.StandardCharsets import java.time.Duration -import java.util -import java.util.{Arrays, Collections, Properties} +import java.util.{Arrays, Properties} import kafka.utils.Exit import org.apache.kafka.clients.consumer.{ConsumerConfig, KafkaConsumer} diff --git a/core/src/main/scala/kafka/tools/StreamsResetter.java b/core/src/main/scala/kafka/tools/StreamsResetter.java index 2a00126592f46..09d3b9ea9a159 100644 --- a/core/src/main/scala/kafka/tools/StreamsResetter.java +++ b/core/src/main/scala/kafka/tools/StreamsResetter.java @@ -33,7 +33,6 @@ import org.apache.kafka.clients.consumer.KafkaConsumer; import org.apache.kafka.clients.consumer.OffsetAndTimestamp; import org.apache.kafka.common.KafkaFuture; -import org.apache.kafka.common.PartitionInfo; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.annotation.InterfaceStability; import org.apache.kafka.common.serialization.ByteArrayDeserializer; From 131b4227f47ff7f1225d571276ff4adea5548c5d Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Fri, 10 Aug 2018 10:11:03 +0200 Subject: [PATCH 10/11] Auto create topic in EndToEndLatency --- .../main/scala/kafka/tools/EndToEndLatency.scala | 15 ++++++--------- 1 file changed, 6 insertions(+), 9 deletions(-) diff --git a/core/src/main/scala/kafka/tools/EndToEndLatency.scala b/core/src/main/scala/kafka/tools/EndToEndLatency.scala index ec56f2b6f8892..4849b1ed8c627 100755 --- a/core/src/main/scala/kafka/tools/EndToEndLatency.scala +++ b/core/src/main/scala/kafka/tools/EndToEndLatency.scala @@ -82,21 +82,18 @@ object EndToEndLatency { producerProps.put(ProducerConfig.VALUE_SERIALIZER_CLASS_CONFIG, "org.apache.kafka.common.serialization.ByteArraySerializer") val producer = new KafkaProducer[Array[Byte], Array[Byte]](producerProps) + // sends a dummy message to create the topic if it doesn't exist + producer.send(new ProducerRecord[Array[Byte], Array[Byte]](topic, Array[Byte]())).get() + def finalise() { consumer.commitSync() producer.close() consumer.close() } - val topics = consumer.listTopics() - val tp = topics.get(topic) - if (tp == null) { - finalise() - throw new RuntimeException("The tested topic doesn't exist in the cluster") - } - val topicPartitions = tp.asScala - .map(pi => new TopicPartition(pi.topic(), pi.partition())) - .to[List].asJava + + val topicPartitions = consumer.partitionsFor(topic).asScala + .map(p => new TopicPartition(p.topic(), p.partition())).asJava consumer.assign(topicPartitions) consumer.seekToEnd(topicPartitions) consumer.assignment().asScala.foreach(consumer.position) From 581292f691b8aefea780ebf72207a56133c68c32 Mon Sep 17 00:00:00 2001 From: Viktor Somogyi Date: Fri, 10 Aug 2018 19:26:25 +0200 Subject: [PATCH 11/11] Fix failing tests in WorkerSinkTaskThreadedTest --- .../connect/runtime/WorkerSinkTaskThreadedTest.java | 9 +++++---- 1 file changed, 5 insertions(+), 4 deletions(-) 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 73689d3571023..d0089e92b6bf5 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 @@ -55,6 +55,7 @@ import org.powermock.modules.junit4.PowerMockRunner; import org.powermock.reflect.Whitebox; +import java.time.Duration; import java.util.Arrays; import java.util.Collection; import java.util.Collections; @@ -525,7 +526,7 @@ private void expectPollInitialAssignment() throws Exception { sinkTask.open(partitions); EasyMock.expectLastCall(); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer(new IAnswer>() { + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer(new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { rebalanceListener.getValue().onPartitionsAssigned(partitions); @@ -557,7 +558,7 @@ private void expectStopTask() throws Exception { private Capture> expectPolls(final long pollDelayMs) throws Exception { // Stub out all the consumer stream/iterator responses, which we just want to verify occur, // but don't care about the exact details here. - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andStubAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andStubAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { @@ -595,7 +596,7 @@ private IExpectationSetters expectOnePoll() { // Currently the SinkTask's put() method will not be invoked unless we provide some data, so instead of // returning empty data, we return one record. The expectation is that the data will be ignored by the // response behavior specified using the return value of this method. - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable { @@ -625,7 +626,7 @@ private IExpectationSetters expectRebalanceDuringPoll() throws Exception final Map offsets = new HashMap<>(); offsets.put(TOPIC_PARTITION, startOffset); - EasyMock.expect(consumer.poll(EasyMock.anyLong())).andAnswer( + EasyMock.expect(consumer.poll(Duration.ofMillis(EasyMock.anyLong()))).andAnswer( new IAnswer>() { @Override public ConsumerRecords answer() throws Throwable {