li3zhi4 commented on PR #11835: URL: https://github.com/apache/seatunnel/pull/11835#issuecomment-5323538821
@DanielLeens — the first CI run surfaced a regression from the initial warm-up placement, now fixed on head `22405c9e0`. **Root cause:** `startUp()` ran `warmUpKafkaTopics` on the static topics **before** `generateTestData` wrote them. The warm-up produces a throwaway record at offset 0 and then `deleteRecords(beforeOffset(1))`, which advances the partition's log start offset by one — so the 100 records written afterwards land at offsets 1..100 instead of 0..99. Tests that read with **absolute offsets** (`specific_offsets` / restore) then observe every id one less than expected (e.g. `testKafkaSpecificOffsetsToConsole` read id 49 where the assertion expects `MIN = 50`), which broke 14 spark-runs of the `kafka-connector-it (8)` job. **Fix:** keep only the read-only `waitForKafkaTopicsReady` leader check in `startUp()` and drop the warm-up call there. The static topics are already written by `generateTestData` right after, and the produce itself forces end-to-end leader propagation before any job is submitted — the warm-up added nothing for those topics while the `deleteRecords` corrupted absolute-offset reads. The shared `warmUpKafkaTopics` helper stays in `AbstractKafkaIT` for test classes whose topics are empty at job-submission time (the scenario it was designed for). Verified locally on the new head: `KafkaIT#testSourceKafka` (which exercises `testKafkaSpecificOffsetsToConsole`) passes on the Zeta engine; the remaining Flink/Spark failures in that run are the pre-existing stale-artifact `InvalidClassException`/`NoSuchMethodError` on the local Flink containers, unrelated to this change. -- 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]
