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]

Reply via email to