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]

Reply via email to