This is an automated email from the ASF dual-hosted git repository.
davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git
The following commit(s) were added to refs/heads/dev by this push:
new 90d5b59bbf [Test][E2E] Retry RocketMQ restore offset checks (#11234)
90d5b59bbf is described below
commit 90d5b59bbf7798383665841941946620291c3a03
Author: Daniel <[email protected]>
AuthorDate: Wed Jul 1 15:58:03 2026 +0800
[Test][E2E] Retry RocketMQ restore offset checks (#11234)
---
.../e2e/connector/rocketmq/RocketMqIT.java | 38 +++++++++++++++++++---
1 file changed, 34 insertions(+), 4 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 9d04f31c8a..67616f4440 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
@@ -88,6 +88,7 @@ import java.util.UUID;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.TimeUnit;
+import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
import static
org.apache.seatunnel.e2e.connector.rocketmq.RocketMqContainer.NAMESRV_PORT;
@@ -653,9 +654,8 @@ public class RocketMqIT extends TestSuiteBase implements
TestResource {
.atMost(5, TimeUnit.MINUTES)
.until(() -> getTopicMaxOffset(sinkTopic) >=
expectedSinkAfterFirstRun + 15);
- Thread.sleep(5000);
- long finalSinkOffset = getTopicMaxOffset(sinkTopic);
long expectedTotal = expectedSinkAfterFirstRun + 15;
+ long finalSinkOffset = awaitTopicMaxOffset(sinkTopic, expectedTotal,
Duration.ofMinutes(1));
Assertions.assertEquals(
expectedTotal,
finalSinkOffset,
@@ -722,11 +722,41 @@ public class RocketMqIT extends TestSuiteBase implements
TestResource {
return result;
}
+ /**
+ * Waits for RocketMQ admin offset visibility and returns the successful
observed offset.
+ *
+ * <p>This keeps the final restore assertion from depending on a single
broker metadata read.
+ */
+ private long awaitTopicMaxOffset(String topicName, long expectedOffset,
Duration timeout) {
+ AtomicLong observedOffset = new AtomicLong();
+ Awaitility.await()
+ .pollInterval(2, TimeUnit.SECONDS)
+ .atMost(timeout)
+ .until(
+ () -> {
+ long current = getTopicMaxOffset(topicName);
+ observedOffset.set(current);
+ return current >= expectedOffset;
+ });
+ return observedOffset.get();
+ }
+
+ /**
+ * Reads topic max offsets with retries because RocketMQ admin queries can
temporarily fail
+ * during restore and broker channel transitions.
+ */
private long getTopicMaxOffset(String topicName) {
try {
List<Map<MessageQueue, TopicOffset>> offsetTopics =
- RocketMqAdminUtil.offsetTopics(
- newConfiguration(), Lists.newArrayList(topicName));
+ RetryUtils.retryWithException(
+ () ->
+ RocketMqAdminUtil.offsetTopics(
+ newConfiguration(),
Lists.newArrayList(topicName)),
+ new RetryUtils.RetryMaterial(
+ Constant.OPERATION_RETRY_TIME,
+ true,
+ exception -> true,
+ Constant.OPERATION_RETRY_SLEEP));
if (offsetTopics.isEmpty()) {
return 0;
}