From 786223bd5747159a0cb677d6d837dd371476e812 Mon Sep 17 00:00:00 2001 From: yaozichen2025 Date: Tue, 11 Aug 2026 11:50:45 +0800 Subject: [PATCH] fix(proxy): tolerate malformed retry counters --- .../proxy/processor/ProducerProcessor.java | 12 ++++++++++-- .../proxy/processor/ProducerProcessorTest.java | 17 +++++++++++++++++ 2 files changed, 27 insertions(+), 2 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java index 8c4907c588a..caf1fab237e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ProducerProcessor.java @@ -211,13 +211,21 @@ protected SendMessageRequestHeader buildSendMessageRequestHeader(List m if (requestHeader.getTopic().startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) { String reconsumeTimes = MessageAccessor.getReconsumeTime(message); if (reconsumeTimes != null) { - requestHeader.setReconsumeTimes(Integer.valueOf(reconsumeTimes)); + try { + requestHeader.setReconsumeTimes(Integer.valueOf(reconsumeTimes)); + } catch (NumberFormatException e) { + log.warn("parse reconsume times error, with value:{}", reconsumeTimes); + } MessageAccessor.clearProperty(message, MessageConst.PROPERTY_RECONSUME_TIME); } String maxReconsumeTimes = MessageAccessor.getMaxReconsumeTimes(message); if (maxReconsumeTimes != null) { - requestHeader.setMaxReconsumeTimes(Integer.valueOf(maxReconsumeTimes)); + try { + requestHeader.setMaxReconsumeTimes(Integer.valueOf(maxReconsumeTimes)); + } catch (NumberFormatException e) { + log.warn("parse max reconsume times error, with value:{}", maxReconsumeTimes); + } MessageAccessor.clearProperty(message, MessageConst.PROPERTY_MAX_RECONSUME_TIMES); } } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java index e6a90df36be..9bdb9265350 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ProducerProcessorTest.java @@ -19,6 +19,7 @@ import java.nio.ByteBuffer; import java.util.ArrayList; +import java.util.Collections; import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CompletionException; @@ -53,6 +54,7 @@ import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertNotNull; +import static org.junit.Assert.assertNull; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; import static org.mockito.ArgumentMatchers.anyLong; @@ -188,6 +190,21 @@ public void testSendRetryMessage() throws Throwable { assertEquals(16, requestHeader.getMaxReconsumeTimes().intValue()); } + @Test + public void testBuildSendMessageRequestHeaderWithMalformedRetryCounters() { + Message message = createMessageExt(MixAll.getRetryTopic(CONSUMER_GROUP), "tag", 0, 0); + MessageAccessor.putProperty(message, MessageConst.PROPERTY_RECONSUME_TIME, "not-a-number"); + MessageAccessor.putProperty(message, MessageConst.PROPERTY_MAX_RECONSUME_TIMES, "also-invalid"); + + SendMessageRequestHeader requestHeader = producerProcessor.buildSendMessageRequestHeader( + Collections.singletonList(message), PRODUCER_GROUP, 0, 0); + + assertEquals(0, requestHeader.getReconsumeTimes().intValue()); + assertNull(requestHeader.getMaxReconsumeTimes()); + assertNull(MessageAccessor.getReconsumeTime(message)); + assertNull(MessageAccessor.getMaxReconsumeTimes(message)); + } + @Test public void testForwardMessageToDeadLetterQueue() throws Throwable { ArgumentCaptor requestHeaderArgumentCaptor = ArgumentCaptor.forClass(ConsumerSendMsgBackRequestHeader.class);