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()); } /**
