Skip to content
Merged
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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),
Expand Down Expand Up @@ -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);
}
Expand Down
Loading