Repository navigation
fix: stabilize task queue and ingestion CI tests - #20291
Conversation
There was a problem hiding this comment.
🟡 Changes recommended
The new datasource-removal hook can lead to duplicate DATASOURCE_REMOVED emissions due to stale rebuild state unless the rebuild set is cleared when the datasource is removed (see stored review comment).
Once you've addressed the issues Copilot identified, you can request another Copilot review.
Pull request overview
This PR stabilizes several flaky CI tests around task queue state accounting, Kafka bounded ingestion offsets, and Broker datasource-removal metric emission by making the tests and cache behavior less timing- and distribution-sensitive.
Changes:
- Add a datasource-removal hook to the segment metadata cache and emit
DATASOURCE_REMOVEDfrom the Broker cache when the last segment is removed (without requiring a refresh cycle). - Make
TaskQueueScaleTestavoid summing non-atomic counters and wait for cleanup completion instead of using a fixed sleep. - Ensure bounded Kafka supervisor tests publish records to both partitions so end offsets are deterministic, and add a regression test for datasource-removal metric emission.
File summaries
| File | Description |
|---|---|
| sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java | Adds regression coverage ensuring last-segment removal emits the datasource-removed metric without a background refresh. |
| sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java | Emits DATASOURCE_REMOVED via the new datasource-removal hook when the cache drops the last segment of a datasource. |
| server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java | Introduces a removeDataSourceAction hook invoked when the last segment of a datasource is removed. |
| indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueScaleTest.java | Removes a flaky assertion based on independent snapshots; waits for cleanup to complete before asserting queue/storage emptiness. |
| embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaBoundedSupervisorTest.java | Publishes records explicitly to each partition to stabilize per-partition end-offset assertions. |
Review details
- Files reviewed: 5/5 changed files
- Comments generated: 1
- Review effort level: Lite
💡 Add a code-review agent skill or configure MCP servers for context-aware, tailored reviews. Learn more in the docs.
FrankChen021
left a comment
There was a problem hiding this comment.
Review deferred: GitHub currently reports this PR as conflicting with master.
Current head: ee5d92b55d98678655c1bb41819702c50bb60f8a (codex/stabilize-ci-tests). Please resolve the conflicts with master and push the updated head so the changes can be reviewed later.
Reviewed 0 of 0 changed files; review was deferred before preparation.
This is an automated review by Codex GPT-5.6-Luna(max)
…tests # Conflicts: # indexing-service/src/test/java/org/apache/druid/indexing/overlord/supervisor/SupervisorManagerTest.java
FrankChen021
left a comment
There was a problem hiding this comment.
| Severity | Findings |
|---|---|
| P0 | 0 |
| P1 | 0 |
| P2 | 1 |
| P3 | 0 |
| Total | 1 |
This is an automated review by Codex GPT-5.6-Luna(max)
FrankChen021
left a comment
There was a problem hiding this comment.
| Severity | Findings |
|---|---|
| P0 | 0 |
| P1 | 0 |
| P2 | 1 |
| P3 | 0 |
| Total | 1 |
Reviewed 6 of 6 changed files.
This is an automated review by Codex GPT-5.6-Luna(max)
# Conflicts: # embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaBoundedSupervisorTest.java
The last-segment callback and an in-flight refresh can both observe the same datasource removal. Gate the metric on the tables.remove() result in both paths so exactly one segment/schemaCache/dataSource/removed event is emitted, regardless of which path removes the table first.
FrankChen021
left a comment
There was a problem hiding this comment.
🟡 Changes recommended
Reviewed all 6 changed files at the current head. Coverage: embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/IngestionSmokeTest.java; embedded-tests/src/test/java/org/apache/druid/testing/embedded/indexing/KafkaBoundedSupervisorTest.java; indexing-service/src/test/java/org/apache/druid/indexing/overlord/TaskQueueScaleTest.java; server/src/main/java/org/apache/druid/segment/metadata/AbstractSegmentMetadataCache.java; sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java; sql/src/test/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCacheTest.java.
The cleanup, bounded Kafka test stabilization, task-queue synchronization, and stale rebuild-state changes otherwise look consistent. One remaining P2 metric/concurrency issue is called out inline: the empty-row-signature path can still duplicate DATASOURCE_REMOVED when a last-segment callback wins an in-flight refresh.
Static review only. Validation actually run: git diff --check against the PR merge base (passed). No builds, tests, dependency installs, formatters, or broad validations were run.
| Severity | Findings |
|---|---|
| P0 | 0 |
| P1 | 0 |
| P2 | 1 |
| P3 | 0 |
| Total | 1 |
Findings that could not be attached inline:
-
sql/src/main/java/org/apache/druid/sql/calcite/schema/BrokerSegmentMetadataCache.java:268 - [P2] Guard the empty-signature removal metric
Finding: The new at-most-once gate covers only the rowSignature == null path. buildDataSourceRowSignature returns an empty signature when the map still contains an unrefreshed segment whose row signature is null; if that last segment is removed while this refresh is in flight, removeSegment removes the table and the new callback emits DATASOURCE_REMOVED, then this branch calls tables.remove (which returns null) but emits the metric unconditionally. Normal segment churn can therefore still produce two removal events for one datasource removal.
Suggestion: Apply the same successful-removal guard or shared helper to the empty-signature branch, and add a regression test covering an in-flight refresh with an uninitialized segment.
This is an automated review by Codex GPT-5.6-Luna(max)
After addressing the findings or replying to the comments, you can request another review from me to trigger a new automated review.
buildDataSourceRowSignature returns an empty (non-null) signature when a datasource only has unrefreshed segments. That branch of refresh() still emitted dataSource/removed unconditionally, so a last-segment callback racing an in-flight refresh could report the removal twice. Apply the same tables.remove() guard, and add a latch-controlled test for the interleaving plus one for a never-initialised datasource.
|
Re the remaining P2 (
Tests added:
Both fail on the previous head ( |
FrankChen021
left a comment
There was a problem hiding this comment.
🟢 Approval recommended
No actionable issues found in this review. The incremental fix correctly closes the previously reported duplicate dataSource/removed path: both removal branches and the last-segment callback gate emission on successful table removal. I reviewed the complete current diff and surrounding task cleanup, bounded Kafka offset setup, task-queue synchronization, and metadata-cache changes; no remaining merge-blocking correctness issue was found.
Reviewed 6 of 6 changed files.
Static review only. Validation actually run: git diff --check 2fed7ca11d9852412ce76798a3fa150ddae9afdb a59f5f5d4d0856b568dbfe3d31ebd3135b47107d (passed). No builds, tests, dependency installs, formatters, or broad validations were run.
This is an automated review by Codex GPT-5.6-Luna(max)
Description
Fix three CI failures observed in Dependabot PRs:
Validation
IngestionSmokeTestmethods passed locally (retries disabled), plus Checkstyle. Terminate the owned supervisor and drain/cancel datasource-scoped tasks between methods so leftover Kafka tasks cannot occupy the two worker slots; keep assertion timeouts unchanged.Release note
Emit the Broker datasource-removal metric when the last segment is removed directly from the cache.
Review