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]

Reply via email to