DanielLeens opened a new issue, #12138: URL: https://github.com/apache/seatunnel/issues/12138
### Search before asking - [x] I had searched in the [issues](https://github.com/apache/seatunnel/issues?q=is%3Aissue) and found no similar issues. ### What happened `RocketMqIT#testSourceRocketMqRestore` (Zeta, `rocketmq_source_restore.conf`, `start.mode = CONSUME_FROM_FIRST_OFFSET`) intermittently sees **duplicate messages after a savepoint restore**: - apache/seatunnel run 33986806422 (PR #12115, `rocketmq-connector-it (8)`, job 101363395159): `Unexpected sink message count after restore. Expected: 45, actual: 47` - PR #11727, fork run `abdessalems/seatunnel` 33994609465, job 101383187178 (`rocketmq-connector-it (8)`): `Expected: 45, actual: 55` The first run wrote exactly 30 messages to the sink before the savepoint (the test logs `[Restore] polling sinkTopic offset: current=30, expected=30` before it triggers the savepoint), the restored job then produced 15 new messages plus 2, respectively 10, messages it had already consumed before the savepoint. The rest of the class passed on the same runs (the broker double-identity problem from #12115 is not involved; the `TBW102` route only ever listed `broker-a` in both runs). ### Root cause (from `dev` source) The restore contract between Zeta and `RocketMqSourceSplitEnumerator` is order-dependent, and the order is not guaranteed: 1. `RocketMqSourceSplitEnumerator(…, RocketMqSourceState sourceState)` deliberately ignores the checkpointed state (`sourceState` "intentionally unused", `isRestored = true` only). The Javadoc says restoration happens through `addSplitsBack(List, int)`, "which the engine calls before `run()`". 2. In Zeta the reader restores its splits in `SourceFlowLifeCycle#restoreState` (`reader.addSplits(splits)`, then `RestoredSplitOperation` -> `SourceSplitEnumeratorTask#addSplitsBack`). The enumerator task, however, starts `enumerator.run()` as soon as `restoreComplete.isDone() && readerRegisterComplete` (`SourceSplitEnumeratorTask#stateProcess`, WAITING_RESTORE -> READY_START -> STARTING). Reader registration happens independently of the reader's restore notification, so `run()` can execute before the `RestoredSplitOperation` arrives. 3. When that happens, `run()` -> `fetchPendingPartitionSplit()` rediscovers the queue (it is in neither `assignedSplit` nor `restoredSplits`), `setPartitionStartOffset()` assigns it the start-mode offset (`CONSUME_FROM_FIRST_OFFSET` -> the queue's min offset), `restoredSplits.clear()`, and `assignSplit()` dispatches this second split to the reader. The later `addSplitsBack` call sees `initialized == true`, goes through `convertToNextSplit()` and is then skipped by `assignSplit()` because the queue is already in `assignedSplit`. 4. The reader now holds two `RocketMqSourceSplit`s for the same `MessageQueue`: the one it restored itself (start offset 30) and the rediscovered one (start offset 0). `RocketMqConsumerThread#assign` seeks whenever `lastPolledOffset != startOffset - 1`, so every `pollNext` cycle alternates between the two positions and re-emits already-consumed messages; how many depends on how many poll batches run before the test counts (`batchSize = 10` in the test config: +10 in one run, +2 in the other). The `checkpoint 4 do not exist or have already been committed` warning that `RocketMqSourceReader#notifyCheckpointComplete` logs right after the restore in both runs is the visible trace of the fresh reader instance receiving the savepoint's completion notification; it is harmless by itself but confirms the timeline. ### What you expected to happen Restoring from a savepoint must not re-emit messages the job already produced: the enumerator should not dispatch a second split for a queue that the checkpoint state says a reader already owns. Concretely, `RocketMqSourceSplitEnumerator` should use the `sourceState` it is given (the checkpointed `assignedSplit` set) to mark those queues as restored, so `run()` neither resets their offsets nor assigns them again, instead of relying on `addSplitsBack` racing `run()`. ### SeaTunnel Version dev (`af0a647d`), Zeta, `apache/rocketmq:4.9.4` e2e image. ### Engine Zeta ### Additional context Seen twice today on unrelated PRs (#12115 and #11727); a third failure mode of the same test (sink route dropped after the job stop) was separately handled in #12115. Filed while triaging CI for #11727; that PR does not touch connector-rocketmq. -- 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]
