This is an automated email from the ASF dual-hosted git repository.

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-12375-065a3df4b39dc116c3d432b15d5c7c4eadec120c
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit acc7d3f7900861a99d420ba7396ffe2a044712ed
Author: Daniel <[email protected]>
AuthorDate: Fri Sep 18 14:47:00 2026 +0000

    [Test][E2E] Count each stored message once in RocketMqIT's restore 
verification (#12375)
    
    Co-authored-by: DanielLeens <[email protected]>
    Co-authored-by: Claude Fable 5.1 <[email protected]>
---
 .../e2e/connector/rocketmq/RocketMqIT.java         | 32 ++++++++++++++++++++--
 1 file changed, 29 insertions(+), 3 deletions(-)

diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
index 86a3c9c167..8b26dbb550 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-rocketmq-e2e/src/test/java/org/apache/seatunnel/e2e/connector/rocketmq/RocketMqIT.java
@@ -82,6 +82,7 @@ import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collections;
 import java.util.HashMap;
+import java.util.LinkedHashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Set;
@@ -799,8 +800,26 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
                 15, restoreCount, "Expected 15 '_restore_' messages, got: " + 
restoreCount);
     }
 
+    /**
+     * Reads every message stored in the topic from {@code fromOffset}, 
counting each stored message
+     * exactly once.
+     *
+     * <p>Messages are keyed by their physical position (broker, queue id, 
queue offset) because
+     * {@code DefaultLitePullConsumer} in rocketmq-client 4.9.4 can hand the 
same stored message to
+     * {@code poll()} twice around {@code seek()}: the pull task started by 
{@code assign()} may
+     * already be past its cancellation check when {@code seek()} cancels it, 
and once the
+     * replacement task has consumed the seek offset the stale batch (up to 
{@code pullBatchSize}
+     * messages from the pre-seek position) is still put into the consumer's 
cache. That made the
+     * restore test report phantom duplicates while the sink topic's max 
offset proved it held
+     * exactly the expected number of messages. A duplicate really written by 
the connector occupies
+     * its own queue offset, so it is still counted here.
+     *
+     * @param topicName topic to read
+     * @param fromOffset lowest queue offset to read from, clamped to each 
queue's min offset
+     * @return message bodies in first-seen order, one entry per stored message
+     */
     private List<String> pollMessagesFromOffset(String topicName, long 
fromOffset) {
-        List<String> result = new ArrayList<>();
+        Map<String, String> bodyByPosition = new LinkedHashMap<>();
         try {
             DefaultLitePullConsumer consumer =
                     
RocketMqAdminUtil.initDefaultLitePullConsumer(newConfiguration(), false);
@@ -829,14 +848,21 @@ public class RocketMqIT extends TestSuiteBase implements 
TestResource {
                     break;
                 }
                 for (MessageExt msg : messages) {
-                    result.add(new String(msg.getBody(), 
StandardCharsets.UTF_8));
+                    String position =
+                            msg.getBrokerName()
+                                    + "#"
+                                    + msg.getQueueId()
+                                    + "#"
+                                    + msg.getQueueOffset();
+                    bodyByPosition.putIfAbsent(
+                            position, new String(msg.getBody(), 
StandardCharsets.UTF_8));
                 }
             }
             consumer.shutdown();
         } catch (Exception e) {
             log.warn("Failed to poll messages from {}: {}", topicName, 
e.getMessage(), e);
         }
-        return result;
+        return new ArrayList<>(bodyByPosition.values());
     }
 
     /**

Reply via email to