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();

Reply via email to