DanielLeens commented on issue #12138: URL: https://github.com/apache/seatunnel/issues/12138#issuecomment-5714437681
Correction to my own root-cause analysis above: I no longer think the enumerator is at fault, and the evidence in the failing runs points at the test's verification helper instead. Fix proposed in #12375. What changed my mind: 1. In all four failing runs listed here and in #12375 (two more hit unrelated PRs on 2026-09-16/17), the assertion that fails is the *second* count check. The one right before it, `assertEquals(expectedTotal, awaitTopicMaxOffset(sinkTopic, ...))`, passed with exactly 45 every time. That value is the sum of the broker's max offsets for the sink topic, i.e. the number of messages physically stored. If the restored job had re-read a queue from offset 0 as I described, the sink topic would grow past 45 and that first assertion would be the one to fail. 2. The ordering I claimed was unguaranteed is in fact guaranteed in Zeta: `SourceFlowLifeCycle#restoreState` sends `RestoredSplitOperation` and blocks on `.get()`, the reader only reports `READY_START` after that, and `SourceSplitEnumeratorTask` only reaches `STARTING` -> `enumerator.run()` after the coordinator has seen every task `READY_START` and called start. So `addSplitsBack` always precedes `run()` on restore. 3. The surplus comes from `pollMessagesFromOffset` (`assign()` + `seek()` on a `DefaultLitePullConsumer`). In rocketmq-client 4.9.4 the pull task started by `assign()` checks `isCancelled()` before taking the queue lock; if it is past that check when `seek()` cancels it and the replacement task has already consumed the seek offset (`nextPullOffset` resets it to `-1`), the stale batch is still put into the cache and returned by `poll()`. The surplus is bounded by `pullBatchSize`, which the IT sets to 10, matching the `+10`, `+2`, `+1` seen. So this is a test-side artifact, not duplicate consumption by the connector. @1328837476-hug, apologies for pointing you at the enumerator; no connector change is needed for this symptom. I will leave the issue open until #12375 is confirmed on CI. -- 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]
