Skip to content

test: stabilize embedded Kafka supervisor tests - #20479

Merged
FrankChen021 merged 2 commits into
apache:masterfrom
gitedmond:fix-embedded-kafka-supervisor-flakes
Oct 6, 2026
Merged

FrankChen021 merged 2 commits into
apache:masterfrom
gitedmond:fix-embedded-kafka-supervisor-flakes

Conversation

@gitedmond

@gitedmond gitedmond commented Oct 4, 2026 •

Copy link
Copy Markdown
Contributor

Fixes #20478.

Description

Stabilize EmbeddedKafkaSupervisorTest when Kafka tasks roll over or segment handoff is slower than ingestion.

  • List active tasks only and assert that the list is nonempty and every task is RUNNING, allowing the publishing task and its replacement to overlap.
  • Track the submitted supervisor and suspend it in @AfterEach, cancel remaining active tasks, and wait for them to finish so a failed assertion cannot leave ingestion running into the next test.
  • Emit Indexer metrics every 100 ms rather than waiting for the task's final emission.
  • Keep the empty-dimension test's 100 records within one day, use the default segment row limit, and allow ingestion to finish before time-based rollover. Removing the one-row limit alone would still create 100 daily segments.
  • Wait for handoff before checking empty-dimension rows and schema. Preserve the first test's ten-segment handoff and lock-release assertions.

Validation

  • On commit ddc23b5ac124df5ac3e5d6ae877e9f37411fedf7, the unmodified EmbeddedKafkaSupervisorTest passed three serial full-class runs using Testcontainers 2.0.5 and the default apache/kafka:4.3.0 image on Docker 28.4.0, Java 25.0.2, and Maven 3.9.11. Surefire reruns were disabled: 9 tests passed with no failures, errors, or skips. The earlier standalone Kafka adapter was not used. Surefire class times were 40.540 s, 42.052 s, and 51.333 s.
  • Repeated the controlled rollover diagnostic with Docker on the original parent 1265b47c9a51252295a4aed42d95c459b00bf2a3 and the fix. A temporary barrier waited for a completed task and an active replacement. The original exact-one assertion failed (expected 1, got 2); the fixed test passed with the same barrier. This exercises a completed-plus-active rollover state, rather than claiming to reproduce the precise two-active publishing overlap from CI.
  • Repeated the cleanup diagnostic with Docker: a temporary assertion failure after task startup was the only failure, with no teardown error. The following test checked that all supervisors were suspended and no tasks remained active globally, and both subsequent tests passed. Reruns were disabled for both diagnostics.
  • Prior validation: Java 25 compilation and Checkstyle passed on the final source. Three full-class runs also passed with reruns disabled using real standalone Kafka 4.3.0 through a temporary resource adapter: 9 tests, no failures.

The adapter, rollover barrier, ordering annotations, and injected failure are absent from this PR; the checkout was clean for the three Docker full-class runs and after diagnostics. The Kafka image digest was sha256:0be4c9eb3565733612d2836d65636fd611a219ebcf5b4162e5ac259ea0ecb907.

These are local Docker smoke checks, not a CI pass or proof that a rare flake is impossible. The empty-dimension timeout was confirmed in the referenced CI log but did not occur in the local baseline run. As of 2026-10-04, CI has completed: PR Commit Checks, Static Checks CI, and CodeQL passed; Unit & Integration tests CI failed. Its only failed job, the Java 25 K/U/Z/Y/X job, failed while resolving the Surefire provider in druid-processing because Maven Central returned HTTP 429, before reaching the Kafka test class.

Key changed classes
  • EmbeddedKafkaSupervisorTest

This PR has:

  • been self-reviewed.
  • modified existing tests.
  • added comments explaining timing and cleanup behavior.

Allow task overlap during rollover, exclude completed tasks from active-task assertions, and suspend supervisors and cancel remaining tasks after each test. Emit ingestion metrics promptly and keep the empty-dimension fixture within one daily segment, waiting for handoff before querying.

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

changes are good, left some comments which require some comments in the test cases. although PR description already explains but extra comments in the source code make the code better to understand without finding the PR and reading the description of it

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟢 Approval recommended

No actionable issues found in this review. The active-task query excludes completed tasks while retaining overlapping rollover tasks, and the tracked-supervisor teardown suspends ingestion, cancels remaining active tasks, and waits for the shared test cluster to drain before the next test.

Reviewed 1 of 1 changed files.

Validation: git diff --check 1265b47c9a51252295a4aed42d95c459b00bf2a3 ddc23b5ac124df5ac3e5d6ae877e9f37411fedf7 -- extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/EmbeddedKafkaSupervisorTest.java passed. No build or test commands were run; this was a static review.


This is an automated review by Codex GPT-5.6-Luna(max)

@gitedmond
gitedmond requested a review from FrankChen021 October 5, 2026 18:59
@FrankChen021
FrankChen021 merged commit 205b50f into apache:master Oct 6, 2026
27 checks passed
@github-actions github-actions Bot added this to the 39.0.0 milestone Oct 6, 2026

@FrankChen021 FrankChen021 left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

🟡 Changes recommended

The teardown now attempts to isolate each test by suspending the submitted supervisor, cancelling its active tasks, and waiting for the shared embedded cluster to drain. The failure path still discards the only cleanup handle when any of those operations fails, so the P2 teardown-isolation issue should be addressed before relying on this stabilization.

Reviewed 1 of 1 changed files.

Validation: git diff --check 1265b47c9a51252295a4aed42d95c459b00bf2a3 a06216c184436ada3a4d1ae92382aec14f848641 -- extensions-core/kafka-indexing-service/src/test/java/org/apache/druid/indexing/kafka/simulate/EmbeddedKafkaSupervisorTest.java passed. No build or test commands were run; this was a static review.

Severity Findings
P0 0
P1 0
P2 1
P3 0
Total 1

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.

cluster.callApi().waitForResult(this::getActiveTasks, List::isEmpty).go();
}
finally {
supervisorSpec = null;

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

[P2] Teardown loses the cleanup handle after a failure

Finding: If suspending the supervisor, cancelling a task, or the default 10-second drain wait throws, this finally block still clears the only tracked supervisor spec. Because the test class shares one cluster across methods, any task left running after that failure can continue into the next test and this teardown cannot retry the suspension or cancellation, causing cross-test interference.

Suggestion: Only clear the tracked spec after cleanup succeeds, and preserve or retry it when cancellation or draining fails so a later cleanup can still isolate the following test.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Projects

None yet

Development

Successfully merging this pull request may close these issues.

Flaky test: EmbeddedKafkaSupervisorTest (test_runKafkaSupervisor, test_runSupervisor_withEmptyDimension)

3 participants