From 4576a59d68508a133d0858ae68f922816808e2a8 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Thu, 30 Jul 2026 23:35:45 +0800 Subject: [PATCH 1/5] Test: wait for Kafka partitions before publishing --- .../kafka/simulate/KafkaResource.java | 27 ++++++++++++++----- .../kafka/simulate/KafkaResourceTest.java | 14 +++++++++- 2 files changed, 34 insertions(+), 7 deletions(-) diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java index 8ef0c888b7e6..e052b4a1ebdd 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java @@ -27,10 +27,12 @@ import org.apache.kafka.clients.admin.CreatePartitionsResult; import org.apache.kafka.clients.admin.NewPartitions; import org.apache.kafka.clients.admin.NewTopic; +import org.apache.kafka.clients.admin.OffsetSpec; 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.clients.producer.ProducerRecord; +import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.TopicPartition; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; @@ -41,6 +43,7 @@ import java.util.List; import java.util.Map; import java.util.Set; +import java.util.concurrent.Future; import java.util.concurrent.ThreadLocalRandom; /** @@ -162,9 +165,9 @@ public void deleteTopic(String topicName) } /** - * Increases the number of partitions in the given Kakfa topic. The topic must - * already exist. This method waits until the increase in the partition count - * has started (but not necessarily finished). + * Increases the number of partitions in the given Kafka topic. The topic must + * already exist. This method waits until every partition is ready to handle + * requests. */ @Override public void increasePartitionsInTopic(String topic, int newPartitionCount) @@ -174,8 +177,16 @@ public void increasePartitionsInTopic(String topic, int newPartitionCount) Map.of(topic, NewPartitions.increaseTo(newPartitionCount)) ); - // Wait for the partitioning to start result.values().get(topic).get(); + + // createPartitions() may complete before the new partition leaders are + // ready to handle requests. Verify all partitions through their leaders + // before allowing callers to publish records. + final Map partitionOffsets = new HashMap<>(); + for (int partition = 0; partition < newPartitionCount; partition++) { + partitionOffsets.put(new TopicPartition(topic, partition), OffsetSpec.latest()); + } + admin.listOffsets(partitionOffsets).all().get(); } catch (Exception e) { throw new RuntimeException(e); @@ -218,8 +229,12 @@ public void produceRecordsWithoutTransaction(List props.remove(ProducerConfig.TRANSACTIONAL_ID_CONFIG); try (final KafkaProducer kafkaProducer = new KafkaProducer<>(props)) { - for (ProducerRecord record : records) { - kafkaProducer.send(record); + final List> sendResults = new ArrayList<>(records.size()); + for (final ProducerRecord record : records) { + sendResults.add(kafkaProducer.send(record)); + } + for (final Future sendResult : sendResults) { + sendResult.get(); } } catch (Exception e) { diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java index 9b8237417727..c69084b6de79 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java @@ -23,6 +23,7 @@ import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; +import java.util.Collections; import java.util.Map; import java.util.Set; @@ -58,8 +59,19 @@ public void testKafka() // Test topic creation final String topicName = "test-topic"; - resource.createTopicWithPartitions(topicName, 3); + resource.createTopicWithPartitions(topicName, 2); assertEquals(Set.of(topicName), resource.listTopics()); + + // Verify that records can be published immediately after adding partitions. + resource.increasePartitionsInTopic(topicName, 4); + resource.publishRecordsToTopicWithoutTransaction( + topicName, + Collections.nCopies(1_000, new byte[]{1}) + ); + final Map partitionOffsets = resource.getPartitionOffsets(topicName); + assertEquals(4, partitionOffsets.size()); + assertEquals(1_000, partitionOffsets.values().stream().mapToLong(Long::longValue).sum()); + resource.deleteTopic(topicName); resource.stop(); From f45aef2994da8729b8d8ddf7ba44f136be81b8e7 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Fri, 31 Jul 2026 00:39:05 +0800 Subject: [PATCH 2/5] Test: wait for Kafka topics before publishing --- .../kafka/simulate/KafkaResource.java | 23 +++++++++++-------- .../kafka/simulate/KafkaResourceTest.java | 12 +++++++++- 2 files changed, 25 insertions(+), 10 deletions(-) diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java index e052b4a1ebdd..465e4db35d43 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java @@ -137,6 +137,7 @@ public void createTopicWithPartitions(String topicName, int numPartitions) admin.createTopics( List.of(new NewTopic(topicName, numPartitions, (short) 1)) ).all().get(); + waitForPartitionsToBeReady(admin, topicName, numPartitions); } catch (Exception e) { throw new RuntimeException(e); @@ -178,15 +179,7 @@ public void increasePartitionsInTopic(String topic, int newPartitionCount) ); result.values().get(topic).get(); - - // createPartitions() may complete before the new partition leaders are - // ready to handle requests. Verify all partitions through their leaders - // before allowing callers to publish records. - final Map partitionOffsets = new HashMap<>(); - for (int partition = 0; partition < newPartitionCount; partition++) { - partitionOffsets.put(new TopicPartition(topic, partition), OffsetSpec.latest()); - } - admin.listOffsets(partitionOffsets).all().get(); + waitForPartitionsToBeReady(admin, topic, newPartitionCount); } catch (Exception e) { throw new RuntimeException(e); @@ -310,6 +303,18 @@ private KafkaProducer newProducer(Map extraPrope return new KafkaProducer<>(producerProperties); } + private void waitForPartitionsToBeReady(Admin admin, String topic, int partitionCount) throws Exception + { + // Topic and partition creation may complete before the partition leaders + // are ready to handle requests. Verify all partitions through their + // leaders before allowing callers to publish records. + final Map partitionOffsets = new HashMap<>(); + for (int partition = 0; partition < partitionCount; partition++) { + partitionOffsets.put(new TopicPartition(topic, partition), OffsetSpec.latest()); + } + admin.listOffsets(partitionOffsets).all().get(); + } + private Map commonClientProperties() { return Map.of("bootstrap.servers", getBootstrapServerUrl()); diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java index c69084b6de79..3fcd30503313 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java @@ -62,6 +62,16 @@ public void testKafka() resource.createTopicWithPartitions(topicName, 2); assertEquals(Set.of(topicName), resource.listTopics()); + // Verify that records can be published immediately after creating a topic. + resource.publishRecordsToTopicWithoutTransaction( + topicName, + Collections.nCopies(100, new byte[]{1}) + ); + assertEquals( + 100, + resource.getPartitionOffsets(topicName).values().stream().mapToLong(Long::longValue).sum() + ); + // Verify that records can be published immediately after adding partitions. resource.increasePartitionsInTopic(topicName, 4); resource.publishRecordsToTopicWithoutTransaction( @@ -70,7 +80,7 @@ public void testKafka() ); final Map partitionOffsets = resource.getPartitionOffsets(topicName); assertEquals(4, partitionOffsets.size()); - assertEquals(1_000, partitionOffsets.values().stream().mapToLong(Long::longValue).sum()); + assertEquals(1_100, partitionOffsets.values().stream().mapToLong(Long::longValue).sum()); resource.deleteTopic(topicName); From cb6d8b2e66ed2f7c50a134ee40f6dbfc780e1373 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Fri, 31 Jul 2026 06:45:41 +0800 Subject: [PATCH 3/5] Test: target newly added Kafka partitions --- .../kafka/simulate/KafkaResourceTest.java | 46 +++++++++++-------- 1 file changed, 28 insertions(+), 18 deletions(-) diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java index 3fcd30503313..7c0fb8ce6316 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java @@ -20,10 +20,11 @@ package org.apache.druid.indexing.kafka.simulate; import org.apache.druid.testing.embedded.EmbeddedDruidCluster; +import org.apache.kafka.clients.producer.ProducerRecord; import org.junit.jupiter.api.Test; import org.junit.jupiter.api.Timeout; -import java.util.Collections; +import java.util.List; import java.util.Map; import java.util.Set; @@ -59,30 +60,39 @@ public void testKafka() // Test topic creation final String topicName = "test-topic"; - resource.createTopicWithPartitions(topicName, 2); + resource.createTopicWithPartitions(topicName, 3); assertEquals(Set.of(topicName), resource.listTopics()); - // Verify that records can be published immediately after creating a topic. - resource.publishRecordsToTopicWithoutTransaction( - topicName, - Collections.nCopies(100, new byte[]{1}) + // Verify that every partition can accept records immediately after creating a topic. + resource.produceRecordsWithoutTransaction( + List.of( + new ProducerRecord<>(topicName, 0, null, new byte[]{1}), + new ProducerRecord<>(topicName, 1, null, new byte[]{1}), + new ProducerRecord<>(topicName, 2, null, new byte[]{1}) + ) ); assertEquals( - 100, - resource.getPartitionOffsets(topicName).values().stream().mapToLong(Long::longValue).sum() + Map.of("0", 1L, "1", 1L, "2", 1L), + resource.getPartitionOffsets(topicName) ); + resource.deleteTopic(topicName); - // Verify that records can be published immediately after adding partitions. - resource.increasePartitionsInTopic(topicName, 4); - resource.publishRecordsToTopicWithoutTransaction( - topicName, - Collections.nCopies(1_000, new byte[]{1}) - ); - final Map partitionOffsets = resource.getPartitionOffsets(topicName); - assertEquals(4, partitionOffsets.size()); - assertEquals(1_100, partitionOffsets.values().stream().mapToLong(Long::longValue).sum()); + final String expandedTopicName = "test-expanded-topic"; + resource.createTopicWithPartitions(expandedTopicName, 2); + resource.increasePartitionsInTopic(expandedTopicName, 4); - resource.deleteTopic(topicName); + // Verify that newly added partitions can accept records immediately. + resource.produceRecordsWithoutTransaction( + List.of( + new ProducerRecord<>(expandedTopicName, 2, null, new byte[]{1}), + new ProducerRecord<>(expandedTopicName, 3, null, new byte[]{1}) + ) + ); + assertEquals( + Map.of("0", 0L, "1", 0L, "2", 1L, "3", 1L), + resource.getPartitionOffsets(expandedTopicName) + ); + resource.deleteTopic(expandedTopicName); resource.stop(); assertFalse(resource.isRunning()); From 3b49ef3f3f2ef08f8c038a141016bb0c4619b491 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Fri, 31 Jul 2026 06:48:03 +0800 Subject: [PATCH 4/5] fix(test): strengthen Kafka partition readiness checks --- .../indexing/kafka/simulate/KafkaResource.java | 16 +++++++++++++++- .../kafka/simulate/KafkaResourceTest.java | 6 ++++-- 2 files changed, 19 insertions(+), 3 deletions(-) diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java index 465e4db35d43..f9a3cc19c7cb 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java @@ -21,6 +21,7 @@ import org.apache.druid.indexing.kafka.KafkaConsumerConfigs; import org.apache.druid.indexing.kafka.KafkaIndexTaskModule; +import org.apache.druid.java.util.common.RetryUtils; import org.apache.druid.testing.embedded.EmbeddedDruidCluster; import org.apache.druid.testing.embedded.StreamIngestResource; import org.apache.kafka.clients.admin.Admin; @@ -34,6 +35,7 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.TopicPartition; +import org.apache.kafka.common.errors.RetriableException; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.testcontainers.kafka.KafkaContainer; @@ -54,6 +56,8 @@ */ public class KafkaResource extends StreamIngestResource { + private static final int PARTITION_READINESS_MAX_TRIES = 5; + /** * Kafka Docker image used in embedded tests. The image name is * read from the system property {@code druid.testing.kafka.image} and @@ -312,7 +316,17 @@ private void waitForPartitionsToBeReady(Admin admin, String topic, int partition for (int partition = 0; partition < partitionCount; partition++) { partitionOffsets.put(new TopicPartition(topic, partition), OffsetSpec.latest()); } - admin.listOffsets(partitionOffsets).all().get(); + RetryUtils.retry( + () -> admin.listOffsets(partitionOffsets).all().get(), + KafkaResource::isRetriableKafkaException, + PARTITION_READINESS_MAX_TRIES + ); + } + + private static boolean isRetriableKafkaException(Throwable throwable) + { + return throwable instanceof RetriableException + || throwable.getCause() != null && isRetriableKafkaException(throwable.getCause()); } private Map commonClientProperties() diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java index 7c0fb8ce6316..a15492ca9929 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResourceTest.java @@ -81,15 +81,17 @@ public void testKafka() resource.createTopicWithPartitions(expandedTopicName, 2); resource.increasePartitionsInTopic(expandedTopicName, 4); - // Verify that newly added partitions can accept records immediately. + // Verify that every partition can accept records immediately after increasing the partition count. resource.produceRecordsWithoutTransaction( List.of( + new ProducerRecord<>(expandedTopicName, 0, null, new byte[]{1}), + new ProducerRecord<>(expandedTopicName, 1, null, new byte[]{1}), new ProducerRecord<>(expandedTopicName, 2, null, new byte[]{1}), new ProducerRecord<>(expandedTopicName, 3, null, new byte[]{1}) ) ); assertEquals( - Map.of("0", 0L, "1", 0L, "2", 1L, "3", 1L), + Map.of("0", 1L, "1", 1L, "2", 1L, "3", 1L), resource.getPartitionOffsets(expandedTopicName) ); resource.deleteTopic(expandedTopicName); From 3fc8dfab92fe2fdeb5a223890ab5e23dfba9ada7 Mon Sep 17 00:00:00 2001 From: Frank Chen Date: Fri, 31 Jul 2026 06:49:58 +0800 Subject: [PATCH 5/5] Test: rely on Kafka list offset retries --- .../indexing/kafka/simulate/KafkaResource.java | 16 +--------------- 1 file changed, 1 insertion(+), 15 deletions(-) diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java index f9a3cc19c7cb..465e4db35d43 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/KafkaResource.java @@ -21,7 +21,6 @@ import org.apache.druid.indexing.kafka.KafkaConsumerConfigs; import org.apache.druid.indexing.kafka.KafkaIndexTaskModule; -import org.apache.druid.java.util.common.RetryUtils; import org.apache.druid.testing.embedded.EmbeddedDruidCluster; import org.apache.druid.testing.embedded.StreamIngestResource; import org.apache.kafka.clients.admin.Admin; @@ -35,7 +34,6 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.TopicPartition; -import org.apache.kafka.common.errors.RetriableException; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; import org.testcontainers.kafka.KafkaContainer; @@ -56,8 +54,6 @@ */ public class KafkaResource extends StreamIngestResource { - private static final int PARTITION_READINESS_MAX_TRIES = 5; - /** * Kafka Docker image used in embedded tests. The image name is * read from the system property {@code druid.testing.kafka.image} and @@ -316,17 +312,7 @@ private void waitForPartitionsToBeReady(Admin admin, String topic, int partition for (int partition = 0; partition < partitionCount; partition++) { partitionOffsets.put(new TopicPartition(topic, partition), OffsetSpec.latest()); } - RetryUtils.retry( - () -> admin.listOffsets(partitionOffsets).all().get(), - KafkaResource::isRetriableKafkaException, - PARTITION_READINESS_MAX_TRIES - ); - } - - private static boolean isRetriableKafkaException(Throwable throwable) - { - return throwable instanceof RetriableException - || throwable.getCause() != null && isRetriableKafkaException(throwable.getCause()); + admin.listOffsets(partitionOffsets).all().get(); } private Map commonClientProperties()