From 5cdf99464bd38abd791ff81c0d6551a33420b41c Mon Sep 17 00:00:00 2001 From: Philip Nee Date: Mon, 24 Jul 2023 14:58:42 -0700 Subject: [PATCH 1/5] Test assign in integration test --- .../scala/integration/kafka/api/BaseAsyncConsumerTest.scala | 6 ++++++ 1 file changed, 6 insertions(+) diff --git a/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala b/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala index a0252abf9e696..716912ae31d12 100644 --- a/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala +++ b/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala @@ -17,8 +17,11 @@ package kafka.api import kafka.utils.TestUtils.waitUntilTrue +import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test +import scala.jdk.CollectionConverters.SeqHasAsJava + class BaseAsyncConsumerTest extends AbstractConsumerTest { @Test @@ -42,6 +45,9 @@ class BaseAsyncConsumerTest extends AbstractConsumerTest { val numRecords = 10000 val startingTimestamp = System.currentTimeMillis() sendRecords(producer, numRecords, tp, startingTimestamp = startingTimestamp) + consumer.assign(List(tp).asJava) consumer.commitSync(); + + assertTrue(consumer.assignment.contains(tp)) } } From f76b94176cdcb870f28a9fa4782a2942cbb672e0 Mon Sep 17 00:00:00 2001 From: Philip Nee Date: Mon, 24 Jul 2023 15:10:56 -0700 Subject: [PATCH 2/5] Update BaseAsyncConsumerTest.scala --- .../scala/integration/kafka/api/BaseAsyncConsumerTest.scala | 2 ++ 1 file changed, 2 insertions(+) diff --git a/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala b/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala index 716912ae31d12..4d03e41ed690d 100644 --- a/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala +++ b/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala @@ -32,10 +32,12 @@ class BaseAsyncConsumerTest extends AbstractConsumerTest { val startingTimestamp = System.currentTimeMillis() val cb = new CountConsumerCommitCallback sendRecords(producer, numRecords, tp, startingTimestamp = startingTimestamp) + consumer.assign(List(tp).asJava) consumer.commitAsync(cb) waitUntilTrue(() => { cb.successCount == 1 }, "wait until commit is completed successfully", 5000) + assertTrue(consumer.assignment.contains(tp)) } @Test From 614b9a389fa823897799798b1f3643c9f3368d9b Mon Sep 17 00:00:00 2001 From: Philip Nee Date: Tue, 25 Jul 2023 13:05:49 -0700 Subject: [PATCH 3/5] A but in committed API --- .../clients/consumer/internals/NetworkClientDelegate.java | 7 ++++++- .../clients/consumer/internals/PrototypeAsyncConsumer.java | 2 +- 2 files changed, 7 insertions(+), 2 deletions(-) diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/NetworkClientDelegate.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/NetworkClientDelegate.java index 9fab7f8ef7520..5b1952162f9f4 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/NetworkClientDelegate.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/NetworkClientDelegate.java @@ -236,7 +236,12 @@ AbstractRequest.Builder requestBuilder() { @Override public String toString() { - return "UnsentRequest(builder=" + requestBuilder + ")"; + return "UnsentRequest{" + + "requestBuilder=" + requestBuilder + + ", handler=" + handler + + ", node=" + node + + ", timer=" + timer + + '}'; } } diff --git a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/PrototypeAsyncConsumer.java b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/PrototypeAsyncConsumer.java index be67251bcb93a..178fd74fd17b6 100644 --- a/clients/src/main/java/org/apache/kafka/clients/consumer/internals/PrototypeAsyncConsumer.java +++ b/clients/src/main/java/org/apache/kafka/clients/consumer/internals/PrototypeAsyncConsumer.java @@ -333,7 +333,7 @@ public Map committed(final Set Date: Thu, 27 Jul 2023 09:35:34 -0700 Subject: [PATCH 4/5] Use CollectionConverters._ to resolve a build failures. --- .../scala/integration/kafka/api/BaseAsyncConsumerTest.scala | 3 +-- 1 file changed, 1 insertion(+), 2 deletions(-) diff --git a/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala b/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala index 4d03e41ed690d..76e5648a77d2f 100644 --- a/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala +++ b/core/src/test/scala/integration/kafka/api/BaseAsyncConsumerTest.scala @@ -20,8 +20,7 @@ import kafka.utils.TestUtils.waitUntilTrue import org.junit.jupiter.api.Assertions.assertTrue import org.junit.jupiter.api.Test -import scala.jdk.CollectionConverters.SeqHasAsJava - +import scala.jdk.CollectionConverters._ class BaseAsyncConsumerTest extends AbstractConsumerTest { @Test From 76a430f9dec2bbe6815b5e1f9b5b0f2e7d0a2cf8 Mon Sep 17 00:00:00 2001 From: Philip Nee Date: Thu, 27 Jul 2023 13:36:31 -0700 Subject: [PATCH 5/5] Update BaseAsyncConsumerTest.scala Update BaseAsyncConsumerTest.scala