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..c66d20bc7b3b 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,6 +27,7 @@ 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; @@ -134,6 +135,15 @@ public void createTopicWithPartitions(String topicName, int numPartitions) admin.createTopics( List.of(new NewTopic(topicName, numPartitions, (short) 1)) ).all().get(); + + // createTopics() may complete before the partition leaders are ready to + // handle requests. Verify every partition through its leader before + // allowing callers to start a supervisor or publish records. + final Map partitionOffsets = new HashMap<>(); + for (int partition = 0; partition < numPartitions; partition++) { + partitionOffsets.put(new TopicPartition(topicName, partition), OffsetSpec.latest()); + } + admin.listOffsets(partitionOffsets).all().get(); } catch (Exception e) { throw new RuntimeException(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..ea8e6edc87cf 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; @@ -60,6 +61,16 @@ public void testKafka() final String topicName = "test-topic"; resource.createTopicWithPartitions(topicName, 3); assertEquals(Set.of(topicName), resource.listTopics()); + + // Verify that callers can publish immediately after topic creation. + resource.publishRecordsToTopicWithoutTransaction( + topicName, + Collections.nCopies(1_000, new byte[]{1}) + ); + final Map partitionOffsets = resource.getPartitionOffsets(topicName); + assertEquals(3, partitionOffsets.size()); + assertEquals(1_000, partitionOffsets.values().stream().mapToLong(Long::longValue).sum()); + resource.deleteTopic(topicName); resource.stop();