This is an automated email from the ASF dual-hosted git repository.
lizhimins pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/rocketmq-clients.git
The following commit(s) were added to refs/heads/master by this push:
new d463e640 [ISSUE# 1321] [Java] Preserve attempt ID until messages are
received (#1322)
d463e640 is described below
commit d463e6400e9819f95a944fa086877336d2e6aad8
Author: qianye <[email protected]>
AuthorDate: Fri Aug 7 16:17:36 2026 +0800
[ISSUE# 1321] [Java] Preserve attempt ID until messages are received (#1322)
---
.../java/impl/consumer/ProcessQueueImpl.java | 28 +++++------
.../test/client/AttemptIdIntegrationTest.java | 55 ++++++++++++++++++----
2 files changed, 60 insertions(+), 23 deletions(-)
diff --git
a/java/client/src/main/java/org/apache/rocketmq/client/java/impl/consumer/ProcessQueueImpl.java
b/java/client/src/main/java/org/apache/rocketmq/client/java/impl/consumer/ProcessQueueImpl.java
index 18ee864e..e3e31903 100644
---
a/java/client/src/main/java/org/apache/rocketmq/client/java/impl/consumer/ProcessQueueImpl.java
+++
b/java/client/src/main/java/org/apache/rocketmq/client/java/impl/consumer/ProcessQueueImpl.java
@@ -32,7 +32,6 @@ import com.google.common.util.concurrent.ListenableFuture;
import com.google.common.util.concurrent.MoreExecutors;
import com.google.common.util.concurrent.SettableFuture;
import com.google.errorprone.annotations.concurrent.GuardedBy;
-import io.grpc.StatusRuntimeException;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Collections;
@@ -251,6 +250,9 @@ class ProcessQueueImpl implements ProcessQueue {
Futures.addCallback(future, new
FutureCallback<ReceiveMessageResult>() {
@Override
public void onSuccess(ReceiveMessageResult result) {
+ // Rotate the attempt ID based on the raw response
because messages removed by an interceptor
+ // were still delivered by the server for this attempt.
+ final boolean hasReceivedMessages =
!result.getMessageViewImpls().isEmpty();
// Intercept after message reception.
final List<GeneralMessage> generalMessages =
result.getMessageViewImpls().stream()
.map((Function<MessageView, GeneralMessage>)
GeneralMessageImpl::new)
@@ -298,35 +300,29 @@ class ProcessQueueImpl implements ProcessQueue {
// Create new ReceiveMessageResult with
filtered messages.
ReceiveMessageResult filteredResult =
ReceiveMessageResult.createFilteredResult(result, remainingMessages);
- onReceiveMessageResult(filteredResult);
+ onReceiveMessageResult(filteredResult,
request.getAttemptId(), hasReceivedMessages);
} catch (Throwable t) {
// Should never reach here.
log.error("[Bug] Exception raised while
handling receive result, mq={}, endpoints={}, "
+ "clientId={}", mq, endpoints, clientId,
t);
- onReceiveMessageException(t, attemptId);
+ onReceiveMessageException(t,
request.getAttemptId());
}
} else {
// When filtering is disabled, use original result
directly to avoid performance overhead.
try {
- onReceiveMessageResult(result);
+ onReceiveMessageResult(result,
request.getAttemptId(), hasReceivedMessages);
} catch (Throwable t) {
// Should never reach here.
log.error("[Bug] Exception raised while
handling receive result, mq={}, endpoints={}, "
+ "clientId={}", mq, endpoints, clientId,
t);
- onReceiveMessageException(t, attemptId);
+ onReceiveMessageException(t,
request.getAttemptId());
}
}
}
@Override
public void onFailure(Throwable t) {
- String nextAttemptId = null;
- if (t instanceof StatusRuntimeException) {
- StatusRuntimeException exception =
(StatusRuntimeException) t;
- if (io.grpc.Status.DEADLINE_EXCEEDED.getCode() ==
exception.getStatus().getCode()) {
- nextAttemptId = request.getAttemptId();
- }
- }
+ final String nextAttemptId = request.getAttemptId();
// Intercept after message reception.
final MessageInterceptorContextImpl context0 =
new MessageInterceptorContextImpl(context,
MessageHookPointsStatus.ERROR);
@@ -397,7 +393,7 @@ class ProcessQueueImpl implements ProcessQueue {
return cachedMessagesBytes.get();
}
- private void onReceiveMessageResult(ReceiveMessageResult result) {
+ private void onReceiveMessageResult(ReceiveMessageResult result, String
attemptId, boolean hasReceivedMessages) {
final List<MessageViewImpl> messages = result.getMessageViewImpls();
if (!messages.isEmpty()) {
cacheMessages(messages);
@@ -405,7 +401,11 @@ class ProcessQueueImpl implements ProcessQueue {
consumer.getReceivedMessagesQuantity().getAndAdd(messages.size());
consumer.getConsumeService().consume(this, messages);
}
- receiveMessage();
+ if (hasReceivedMessages) {
+ receiveMessage();
+ return;
+ }
+ receiveMessage(attemptId);
}
private void evictCache(MessageViewImpl messageView) {
diff --git
a/java/test/src/test/java/org/apache/rocketmq/test/client/AttemptIdIntegrationTest.java
b/java/test/src/test/java/org/apache/rocketmq/test/client/AttemptIdIntegrationTest.java
index 42435d6b..279bd3cc 100644
---
a/java/test/src/test/java/org/apache/rocketmq/test/client/AttemptIdIntegrationTest.java
+++
b/java/test/src/test/java/org/apache/rocketmq/test/client/AttemptIdIntegrationTest.java
@@ -20,9 +20,19 @@ package org.apache.rocketmq.test.client;
import static org.assertj.core.api.Assertions.assertThat;
import static org.awaitility.Awaitility.await;
+import apache.rocketmq.v2.Digest;
+import apache.rocketmq.v2.DigestType;
+import apache.rocketmq.v2.Encoding;
+import apache.rocketmq.v2.Message;
+import apache.rocketmq.v2.MessageType;
import apache.rocketmq.v2.ReceiveMessageRequest;
import apache.rocketmq.v2.ReceiveMessageResponse;
+import apache.rocketmq.v2.Resource;
+import apache.rocketmq.v2.SystemProperties;
+import com.google.protobuf.ByteString;
+import io.grpc.Status;
import io.grpc.stub.StreamObserver;
+import java.nio.charset.StandardCharsets;
import java.time.temporal.ChronoUnit;
import java.util.Collections;
import java.util.List;
@@ -36,6 +46,7 @@ import org.apache.rocketmq.client.apis.consumer.ConsumeResult;
import org.apache.rocketmq.client.apis.consumer.FilterExpression;
import org.apache.rocketmq.client.apis.consumer.FilterExpressionType;
import org.apache.rocketmq.client.apis.consumer.PushConsumer;
+import org.apache.rocketmq.client.java.message.MessageIdCodec;
import org.apache.rocketmq.test.server.BaseMockServerImpl;
import org.apache.rocketmq.test.server.GrpcServerIntegrationTest;
import org.apache.rocketmq.test.server.MockServer;
@@ -46,7 +57,7 @@ public class AttemptIdIntegrationTest extends
GrpcServerIntegrationTest {
private final String topic = "topic";
private MockServer serverImpl;
- static class MockServerImpl extends BaseMockServerImpl {
+ static class MockServerImpl extends BaseMockServerImpl {
public final List<String> attemptIdList = new CopyOnWriteArrayList<>();
public final AtomicBoolean serverDeadlineFlag = new
AtomicBoolean(true);
@@ -58,7 +69,7 @@ public class AttemptIdIntegrationTest extends
GrpcServerIntegrationTest {
public void receiveMessage(ReceiveMessageRequest request,
StreamObserver<ReceiveMessageResponse> responseObserver) {
// prevent too much request
- if (attemptIdList.size() >= 3) {
+ if (attemptIdList.size() >= 5) {
try {
Thread.sleep(100);
} catch (InterruptedException e) {
@@ -68,10 +79,34 @@ public class AttemptIdIntegrationTest extends
GrpcServerIntegrationTest {
attemptIdList.add(request.getAttemptId());
if (serverDeadlineFlag.compareAndSet(true, false)) {
// timeout
- } else {
-
responseObserver.onNext(ReceiveMessageResponse.newBuilder().setStatus(mockStatus).build());
- responseObserver.onCompleted();
+ return;
}
+ if (attemptIdList.size() == 2) {
+
responseObserver.onError(Status.UNAVAILABLE.asRuntimeException());
+ return;
+ }
+
responseObserver.onNext(ReceiveMessageResponse.newBuilder().setStatus(mockStatus).build());
+ if (attemptIdList.size() == 4) {
+
responseObserver.onNext(ReceiveMessageResponse.newBuilder().setMessage(mockMessage()).build());
+ }
+ responseObserver.onCompleted();
+ }
+
+ private Message mockMessage() {
+ final Digest digest =
Digest.newBuilder().setType(DigestType.CRC32).setChecksum("9EF61F95").build();
+ final SystemProperties systemProperties =
SystemProperties.newBuilder()
+ .setMessageType(MessageType.NORMAL)
+
.setMessageId(MessageIdCodec.getInstance().nextMessageId().toString())
+ .setBornHost("127.0.0.1")
+ .setBodyDigest(digest)
+ .setBodyEncoding(Encoding.IDENTITY)
+ .setReceiptHandle("receipt-handle")
+ .build();
+ return Message.newBuilder()
+ .setTopic(Resource.newBuilder().setName(topic).build())
+ .setBody(ByteString.copyFrom("foobar", StandardCharsets.UTF_8))
+ .setSystemProperties(systemProperties)
+ .build();
}
}
@@ -83,7 +118,7 @@ public class AttemptIdIntegrationTest extends
GrpcServerIntegrationTest {
}
@Test
- public void test() throws Exception {
+ public void testAttemptIdIsRotatedOnlyAfterMessagesAreReceived() throws
Exception {
final ClientServiceProvider provider =
ClientServiceProvider.loadService();
String accessKey = "yourAccessKey";
String secretKey = "yourSecretKey";
@@ -106,11 +141,13 @@ public class AttemptIdIntegrationTest extends
GrpcServerIntegrationTest {
.setMessageListener(messageView -> ConsumeResult.SUCCESS)
.build();
try {
- await().atMost(java.time.Duration.ofSeconds(5)).untilAsserted(()
-> {
+ await().atMost(java.time.Duration.ofSeconds(8)).untilAsserted(()
-> {
List<String> attemptIdList = ((MockServerImpl)
serverImpl).attemptIdList;
- assertThat(attemptIdList.size()).isGreaterThanOrEqualTo(3);
+ assertThat(attemptIdList.size()).isGreaterThanOrEqualTo(5);
assertThat(attemptIdList.get(0)).isEqualTo(attemptIdList.get(1));
-
assertThat(attemptIdList.get(0)).isNotEqualTo(attemptIdList.get(2));
+
assertThat(attemptIdList.get(1)).isEqualTo(attemptIdList.get(2));
+
assertThat(attemptIdList.get(2)).isEqualTo(attemptIdList.get(3));
+
assertThat(attemptIdList.get(3)).isNotEqualTo(attemptIdList.get(4));
});
} finally {
pushConsumer.close();