DanielLeens opened a new pull request, #12375: URL: https://github.com/apache/seatunnel/pull/12375
## Summary `RocketMqIT#testSourceRocketMqRestore` keeps failing on PRs that do not touch the RocketMQ connector, always with the same assertion: | where | job | message | | --- | --- | --- | | #12365 (Transform-V2 change), fork run `35189291969`, 2026-09-17 | `rocketmq-connector-it (11)` | `Unexpected sink message count after restore. Expected: 45, actual: 55` | | #12350, fork run `35121735231`, 2026-09-16 | `rocketmq-connector-it (8)` | `... Expected: 45, actual: 46` | | #12115 run `33986806422`, 2026-09-05 | `rocketmq-connector-it (8)` | `... Expected: 45, actual: 47` | | #11727, fork run `33994609465`, 2026-09-05 | `rocketmq-connector-it (8)` | `... Expected: 45, actual: 55` | ## Root cause In every one of those runs the assertion that fails is the **second** count check. The first one, which the test runs immediately before it, passed each time: ```java long finalSinkOffset = awaitTopicMaxOffset(sinkTopic, expectedTotal, Duration.ofMinutes(1)); Assertions.assertEquals(expectedTotal, finalSinkOffset, "Sink offset mismatch after restore ..."); // passes: 45 List<String> allSinkMessages = pollMessagesFromOffset(sinkTopic, 0); Assertions.assertEquals(expectedTotal, allSinkMessages.size(), "Unexpected sink message count ..."); // fails: 46 / 47 / 55 ``` `getTopicMaxOffset` sums the broker's max offsets for the sink topic, i.e. the number of messages physically stored in it, and it is exactly 45. The restored job therefore did not write anything twice. What over-counts is the test's own reader, `pollMessagesFromOffset`, which does `consumer.assign(queues)` followed by `consumer.seek(mq, 0)` on a `DefaultLitePullConsumer` and then counts everything `poll()` returns. In rocketmq-client 4.9.4 (`DefaultLitePullConsumerImpl`), that sequence can deliver a stored message twice: 1. `assign()` starts a `PullTaskImpl` that pulls from the consumer group's current position. 2. `seek()` takes the queue lock, clears the cache, calls `tryInterrupt()` on that task, sets the seek offset and schedules a replacement task. 3. The old task checks `isCancelled()` **before** taking the queue lock. If it has already passed that check when `seek()` cancels it, it just waits for the lock. Meanwhile the replacement task consumes the seek offset (`nextPullOffset` resets it to `-1`). 4. The old task then sees `getSeekOffset(mq) == -1` and puts its stale batch into the process queue, so `poll()` hands those messages over in addition to the ones read from the sought position. The surplus is bounded by `pullBatchSize`, which this IT sets to 10 (`newConfiguration().batchSize(10)`), matching the observed `+10`, and is smaller (`+1`, `+2`) when the stale pull started near the end of the queue. This also means the explanation in #12138 (enumerator rediscovering a restored queue) does not match these failures: a second split re-reading from offset 0 would push the sink topic's max offset past 45, and the offset assertion would fail first. Zeta also cannot run the enumerator before `addSplitsBack`: the reader's `restoreState` sends `RestoredSplitOperation` synchronously, the reader only then reports `READY_START`, and `SourceSplitEnumeratorTask` only calls `enumerator.run()` after the coordinator has seen every task ready. ## Fix (test-only) `pollMessagesFromOffset` now keys what it reads by physical position (`brokerName#queueId#queueOffset`) and returns one body per stored message, in first-seen order. A message handed over twice by the pull consumer is counted once; a duplicate really written by the connector occupies its own queue offset, so it is still counted, and the topic max-offset assertion in front of it is unchanged. No assertion, expected value or timeout changes. ## Files - `seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java` ## Test plan Test-helper-only change; verified by this PR's GitHub CI (`rocketmq-connector-it`, both JDK legs), per this initiative's no-local-build policy. The other intermittent failures of this class (`testRocketMqEarliestToConsole`, `testSourceRocketMqTextTagToConsole`) are the topic-route problem tracked in #12322 / #12115 / #12323 and are not addressed here. 🤖 Generated with [Claude Code](https://claude.com/claude-code) -- 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]
