From a955e7d2dbf57a049ab8bb81350a63e68debd1e9 Mon Sep 17 00:00:00 2001 From: zjncs <18910855655@163.com> Date: Fri, 11 Sep 2026 14:13:45 +0800 Subject: [PATCH] fix(broker): stop overwriting batch ack uniq key with single ack key in appendAck appendAck set PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX correctly per ack type (genBatchAckUniqueId for BatchAckMsg, genAckUniqueId otherwise), but a later unconditional put overwrote it with genAckUniqueId for every ack. Batch acks therefore landed on the revive topic with a bogus uniq key: the offset segment was the -1 sentinel the batch path assigns and the tag segment was ACK instead of BATCH_ACK, so every batch ack of a pop produced the same non-unique client id and tracing by uniq key was impossible. The buffered path in PopBufferMergeService already writes the batch uniq key without overwriting it. Drop the stray put so the per-type key set just above survives. --- .../broker/processor/AckMessageProcessor.java | 1 - .../processor/AckMessageProcessorTest.java | 49 +++++++++++++++++++ 2 files changed, 49 insertions(+), 1 deletion(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java index 65f5f79aec4..0a22a88f00a 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/AckMessageProcessor.java @@ -298,7 +298,6 @@ private void appendAck(final AckMessageRequestHeader requestHeader, final BatchA msgInner.setBornHost(this.brokerController.getStoreHost()); msgInner.setStoreHost(this.brokerController.getStoreHost()); msgInner.setDeliverTimeMs(popTime + invisibleTime); - msgInner.getProperties().put(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX, PopMessageProcessor.genAckUniqueId(ackMsg)); msgInner.setPropertiesString(MessageDecoder.messageProperties2String(msgInner.getProperties())); if (brokerController.getBrokerConfig().isAppendAckAsync()) { int finalAckCount = ackCount; diff --git a/broker/src/test/java/org/apache/rocketmq/broker/processor/AckMessageProcessorTest.java b/broker/src/test/java/org/apache/rocketmq/broker/processor/AckMessageProcessorTest.java index 1add8bd8d2d..aa5828af2d3 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/processor/AckMessageProcessorTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/processor/AckMessageProcessorTest.java @@ -49,9 +49,11 @@ import org.apache.rocketmq.store.PutMessageStatus; import org.apache.rocketmq.store.config.MessageStoreConfig; import org.apache.rocketmq.store.exception.ConsumeQueueException; +import org.apache.rocketmq.store.pop.BatchAckMsg; import org.junit.Before; import org.junit.Test; import org.junit.runner.RunWith; +import org.mockito.ArgumentCaptor; import org.mockito.Mock; import org.mockito.Mockito; import org.mockito.Spy; @@ -366,4 +368,51 @@ public void testBatchAck_appendAck() throws RemotingCommandException { } } + @Test + public void testBatchAck_appendAck_BatchUniqKeyKept() throws RemotingCommandException { + PopBufferMergeService popBufferMergeService = mock(PopBufferMergeService.class); + when(popBufferMergeService.addAk(anyInt(), any())).thenReturn(false); + when(popMessageProcessor.getPopBufferMergeService()).thenReturn(popBufferMergeService); + PutMessageResult putMessageResult = new PutMessageResult(PutMessageStatus.PUT_OK, null); + ArgumentCaptor msgCaptor = ArgumentCaptor.forClass(MessageExtBrokerInner.class); + when(messageStore.putMessage(msgCaptor.capture())).thenReturn(putMessageResult); + + long popTime = 1666860736757L; + String brokerName = "broker-a"; + BatchAck bAck1 = new BatchAck(); + bAck1.setConsumerGroup(MixAll.DEFAULT_CONSUMER_GROUP); + bAck1.setTopic(topic); + bAck1.setQueueId(0); + bAck1.setReviveQueueId(0); + bAck1.setStartOffset(MIN_OFFSET_IN_QUEUE); + bAck1.setBitSet(new BitSet()); + bAck1.getBitSet().set(1); + bAck1.setRetry("0"); + bAck1.setPopTime(popTime); + bAck1.setInvisibleTime(60000L); + + RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.BATCH_ACK_MESSAGE, null); + BatchAckMessageRequestBody reqBody = new BatchAckMessageRequestBody(); + reqBody.setAcks(Collections.singletonList(bAck1)); + reqBody.setBrokerName(brokerName); + request.setBody(reqBody.encode()); + request.makeCustomHeaderToNet(); + RemotingCommand response = ackMessageProcessor.processRequest(handlerContext, request); + + assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS); + MessageExtBrokerInner ackMessage = msgCaptor.getValue(); + assertThat(ackMessage).isNotNull(); + + BatchAckMsg expected = new BatchAckMsg(); + expected.setConsumerGroup(MixAll.DEFAULT_CONSUMER_GROUP); + expected.setTopic(topic); + expected.setQueueId(0); + expected.setStartOffset(MIN_OFFSET_IN_QUEUE); + expected.setPopTime(popTime); + expected.setBrokerName(brokerName); + expected.getAckOffsetList().add(MIN_OFFSET_IN_QUEUE + 1); + + assertThat(ackMessage.getProperties().get(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX)) + .isEqualTo(PopMessageProcessor.genBatchAckUniqueId(expected)); + } }