diff --git a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/EmbeddedKafkaSupervisorTest.java b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/EmbeddedKafkaSupervisorTest.java index 6afb70a79a40..ea0a3ca49ff8 100644 --- a/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/EmbeddedKafkaSupervisorTest.java +++ b/extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/EmbeddedKafkaSupervisorTest.java @@ -215,8 +215,9 @@ public void test_runKafkaSupervisorWithHeaderFiltering() Assertions.assertTrue(supervisorStatus.isHealthy()); Assertions.assertEquals("RUNNING", supervisorStatus.getState()); - // Suspend the supervisor and wait for segment handoff - cluster.callApi().postSupervisor(kafkaSupervisorSpec.createSuspendedSpec()); + // Wait for segment handoff without suspending the supervisor; its short task duration publishes the segments. + // Suspending while ingestion is in progress can replace the supervisor while the task's checkpoint request is in + // flight, which fails the task before it publishes anything. indexer.latchableEmitter().waitForEventAggregate( event -> event.hasMetricName("ingest/handoff/count") .hasDimension(DruidMetrics.DATASOURCE, dataSource), @@ -345,24 +346,8 @@ private KafkaSupervisorSpec createKafkaSupervisorWithHeaderFilter(String supervi InDimFilter filter = new InDimFilter("environment", ImmutableSet.of("production")); KafkaHeaderBasedFilterConfig headerFilterConfig = new KafkaHeaderBasedFilterConfig(filter, "UTF-8", 1000); - return new KafkaSupervisorSpecBuilder() - .withDataSchema( - schema -> schema - .withTimestamp(new TimestampSpec("timestamp", null, null)) - .withDimensions(DimensionsSpec.EMPTY) - ) - .withTuningConfig( - tuningConfig -> tuningConfig - .withMaxRowsPerSegment(1) - .withReleaseLocksOnHandoff(true) - ) - .withIoConfig( - ioConfig -> ioConfig - .withInputFormat(new CsvInputFormat(List.of("timestamp", "item"), null, null, false, 0, false)) - .withConsumerProperties(kafkaServer.consumerProperties()) - .withUseEarliestSequenceNumber(true) - .withHeaderBasedFilterConfig(headerFilterConfig) - ) + return newKafkaSupervisor() + .withIoConfig(ioConfig -> ioConfig.withHeaderBasedFilterConfig(headerFilterConfig)) .withId(supervisorId) .build(dataSource, topic); }