li3zhi4 opened a new pull request, #11835:
URL: https://github.com/apache/seatunnel/pull/11835
## Why
`KafkaIT` e2e tests fail intermittently on CI with a **topic-readiness
race**: a job submitted right after topic creation can hit
`UnknownTopicOrPartitionException` inside
`KafkaSourceSplitEnumerator.getTopicInfo` (line 383), because the broker the
job's own Kafka client connects to has not finished assigning/propagating the
partition leader yet.
Evidence (PR #11633 head `0dbea3718`, run #38, three consecutive reruns with
the same failure signature — annotation line numbers drift by only tens of
lines):
| Run | Job | `SeaTunnel job executed failed` lines | Spark failure lines |
|---|---|---|---|
| 1 | 94330487649 | 423547 / 419198 / 414894 / 414603 / 412919 | 276989 /
276300 |
| 2 | 94415607650 | 423600 / 419254 / 414856 / 414613 / 412976 | 277444 /
276753 |
| 3 | 94668401633 | 423348 / 423204 / 418821 / 414496 / 414249 | 276967 /
276226 |
The same race previously hit `KafkaJsonDefaultValueIT`; the warm-up fix in
PR #11633 (produce+consume a throwaway record to force metadata propagation
end-to-end, then `deleteRecords` so it does not leak into `start_mode =
earliest` reads) has kept that test green on every CI attempt and locally.
`KafkaIT` is not protected, so unrelated PRs keep flaking on the
`kafka-connector-it` job — this PR applies the same pattern to the shared Kafka
e2e base.
## What changed
New shared base class `AbstractKafkaIT` (in `connector-kafka-e2e` test
sources), with the warm-up logic extracted from `KafkaJsonDefaultValueIT` and
generalized to multiple topics:
1. `waitForKafkaTopicsReady(Collection<String> topics)` — waits until every
partition of every topic reports a non-null leader in the admin metadata view
(stronger than the existing count-only check in `KafkaIT.startUp`).
2. `warmUpKafkaTopics(Collection<String> topics)` — for each topic, produces
and consumes a throwaway record (explicit partition 0) to force leader
propagation end-to-end through the broker, then
`deleteRecords(beforeOffset(1L))` removes the records so they do not leak into
earliest-start reads.
3. `kafkaBootstrapServers()` — abstract accessor implemented by each
concrete test class.
`KafkaIT` now `extends AbstractKafkaIT` and calls both checks in `startUp()`
right after `createTopics` and before writing any test data. Test-only change;
no production code touched.
## File changes
| File | Change |
|---|---|
|
`seatunnel-e2e/seatunnel-connector-v2-e2e/connector-kafka-e2e/src/test/java/org/apache/seatunnel/e2e/connector/kafka/AbstractKafkaIT.java`
| **new** — shared base with `waitForKafkaTopicsReady` + `warmUpKafkaTopics` +
`kafkaBootstrapServers` |
|
`seatunnel-e2e/seatunnel-connector-v2-e2e/connector-kafka-e2e/src/test/java/org/apache/seatunnel/e2e/connector/kafka/KafkaIT.java`
| extend `AbstractKafkaIT`; hook both checks into `startUp()`; drop
now-redundant imports |
`KafkaJsonDefaultValueIT` is intentionally left unchanged: it is part of the
unmerged PR #11633 and does not exist on `dev` yet; it can reuse the shared
base when #11633 lands and the branch is rebased.
## Validation
- `connector-kafka-e2e` `test-compile` ✅ (includes spotless check)
- `git diff --check` ✅
- Local official E2E
(`-Dit.test='KafkaIT#testKafkaToKafkaExactlyOnceOnStreaming+testSourceKafkaJsonToConsole'`,
`-Dtestcontainer.version=1.21.4`, Docker 29.1.3):
- `testKafkaToKafkaExactlyOnceOnStreaming` (the CI-failing target) **all
pass** ✅
- Zeta engine `testSourceKafkaJsonToConsole` passes ✅
- Flink/Spark `testSourceKafkaJsonToConsole` fail with
`InvalidClassException: JsonToRowConverters$19` (stream `-4116979196724038267`
vs local `8925809659107082274`) — a **local stale-artifact issue** (format-json
jar version mismatch on Flink/Spark containers): the diff touches only e2e test
code, not the `format-json` production code; the exactly-once test on the same
engines passes, which rules out this change as the cause.
--
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.
To unsubscribe, e-mail: [email protected]
For queries about this service, please contact Infrastructure at:
[email protected]