diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/MessageReceiptHandle.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/MessageReceiptHandle.java index 4c396e71d0..81fcebd082 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/MessageReceiptHandle.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/MessageReceiptHandle.java @@ -17,11 +17,10 @@ package org.apache.rocketmq.proxy.common; -import org.apache.rocketmq.common.message.MessageQueue; - public class MessageReceiptHandle { private final String group; - private final MessageQueue messageQueue; + private final String topic; + private final int queueId; private final String messageId; private final long queueOffset; private final String originalReceiptHandle; @@ -31,10 +30,11 @@ public class MessageReceiptHandle { private String receiptHandle; - public MessageReceiptHandle(String group, MessageQueue messageQueue, String receiptHandle, String messageId, + public MessageReceiptHandle(String group, String topic, int queueId, String receiptHandle, String messageId, long queueOffset, int reconsumeTimes, long expectInvisibleTime) { this.group = group; - this.messageQueue = messageQueue; + this.topic = topic; + this.queueId = queueId; this.receiptHandle = receiptHandle; this.originalReceiptHandle = receiptHandle; this.messageId = messageId; @@ -48,8 +48,12 @@ public class MessageReceiptHandle { return group; } - public MessageQueue getMessageQueue() { - return messageQueue; + public String getTopic() { + return topic; + } + + public int getQueueId() { + return queueId; } public String getReceiptHandle() { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java index 42808961ce..7d0fcaa6fc 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/consumer/ReceiveMessageActivity.java @@ -32,7 +32,6 @@ import org.apache.rocketmq.common.constant.ConsumeInitMode; import org.apache.rocketmq.common.filter.FilterAPI; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageExt; -import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.proxy.common.ContextVariable; import org.apache.rocketmq.proxy.common.MessageReceiptHandle; @@ -121,9 +120,8 @@ public class ReceiveMessageActivity extends AbstractMessingActivity { for (MessageExt messageExt : messageExtList) { String receiptHandle = messageExt.getProperty(MessageConst.PROPERTY_POP_CK); if (receiptHandle != null) { - MessageQueue messageQueue = new MessageQueue(topic, messageExt.getBrokerName(), messageExt.getQueueId()); MessageReceiptHandle messageReceiptHandle = - new MessageReceiptHandle(group, messageQueue, receiptHandle, messageExt.getMsgId(), + new MessageReceiptHandle(group, topic, messageExt.getQueueId(), receiptHandle, messageExt.getMsgId(), messageExt.getQueueOffset(), messageExt.getReconsumeTimes(), requestInvisibleTime); String channelId = proxyContext.getVal(ContextVariable.CHANNEL_KEY); receiptHandleProcessor.addReceiptHandle(channelId, receiptHandle, messageReceiptHandle); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessor.java index 6b2aaca25e..10179a7fa8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessor.java @@ -134,7 +134,7 @@ public class ReceiptHandleProcessor extends AbstractStartAndShutdown { if (current - messageReceiptHandle.getTimestamp() < messageReceiptHandle.getExpectInvisibleTime()) { CompletableFuture future = messagingProcessor.changeInvisibleTime(ProxyContext.create(), handle, messageReceiptHandle.getMessageId(), - messageReceiptHandle.getGroup(), messageReceiptHandle.getMessageQueue().getTopic(), proxyConfig.getRenewSliceTimeMillis()); + messageReceiptHandle.getGroup(), messageReceiptHandle.getTopic(), proxyConfig.getRenewSliceTimeMillis()); future.thenAccept(ackResult -> { if (AckStatus.OK.equals(ackResult.getStatus())) { messageReceiptHandle.update(ackResult.getExtraInfo()); @@ -144,8 +144,7 @@ public class ReceiptHandleProcessor extends AbstractStartAndShutdown { } else { CompletableFuture future = messagingProcessor.changeInvisibleTime(ProxyContext.create(), handle, messageReceiptHandle.getMessageId(), messageReceiptHandle.getGroup(), - messageReceiptHandle.getMessageQueue().getTopic(), - retryPolicy.nextDelayDuration(messageReceiptHandle.getReconsumeTimes(), TimeUnit.MILLISECONDS)); + messageReceiptHandle.getTopic(), retryPolicy.nextDelayDuration(messageReceiptHandle.getReconsumeTimes(), TimeUnit.MILLISECONDS)); future.thenAccept(ackResult -> { if (AckStatus.OK.equals(ackResult.getStatus())) { removeReceiptHandle(key, messageReceiptHandle.getOriginalReceiptHandle()); @@ -187,7 +186,7 @@ public class ReceiptHandleProcessor extends AbstractStartAndShutdown { receiptHandle, value0.getMessageId(), value0.getGroup(), - value0.getMessageQueue().getTopic(), + value0.getTopic(), proxyConfig.getInvisibleTimeMillisWhenClear() ); }); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessorTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessorTest.java index 73c05a6784..c8587a8294 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessorTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/processor/ReceiptHandleProcessorTest.java @@ -23,7 +23,6 @@ import org.apache.rocketmq.broker.client.ConsumerIdsChangeListener; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.client.consumer.AckStatus; import org.apache.rocketmq.common.consumer.ReceiptHandle; -import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.common.subscription.SubscriptionGroupConfig; import org.apache.rocketmq.proxy.common.ContextVariable; import org.apache.rocketmq.proxy.common.MessageReceiptHandle; @@ -38,7 +37,9 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { ProxyContext context = ProxyContext.create(); String group = "group"; - MessageQueue messageQueue = new MessageQueue("topic", "broker", 1); + String topic = "topic"; + String brokerName = "broker"; + int queueId = 1; String messageId = "messageId"; long offset = 123L; long invisibleTime = 100000L; @@ -51,8 +52,8 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { .invisibleTime(invisibleTime) .reviveQueueId(1) .topicType(ReceiptHandle.NORMAL_TOPIC) - .brokerName(messageQueue.getBrokerName()) - .queueId(messageQueue.getQueueId()) + .brokerName(brokerName) + .queueId(queueId) .offset(offset) .commitLogOffset(0L) .build().encode(); @@ -62,7 +63,7 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { context.withVal(ContextVariable.CHANNEL_KEY, "channel-id"); receiptHandleProcessor = new ReceiptHandleProcessor(messagingProcessor); Mockito.doNothing().when(messagingProcessor).registerConsumerListener(Mockito.any(ConsumerIdsChangeListener.class)); - messageReceiptHandle = new MessageReceiptHandle(group, messageQueue, receiptHandle, messageId, offset, + messageReceiptHandle = new MessageReceiptHandle(group, topic, queueId, receiptHandle, messageId, offset, reconsumeTimes, invisibleTime); } @@ -74,7 +75,7 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { receiptHandleProcessor.scheduleRenewTask(); Mockito.verify(messagingProcessor, Mockito.timeout(1000).times(1)) .changeInvisibleTime(Mockito.any(ProxyContext.class), Mockito.any(ReceiptHandle.class), Mockito.eq(messageId), - Mockito.eq(group), Mockito.eq(messageQueue.getTopic()), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis())); + Mockito.eq(group), Mockito.eq(topic), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis())); } @Test @@ -90,8 +91,8 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { .invisibleTime(newInvisibleTime) .reviveQueueId(1) .topicType(ReceiptHandle.NORMAL_TOPIC) - .brokerName(messageQueue.getBrokerName()) - .queueId(messageQueue.getQueueId()) + .brokerName(brokerName) + .queueId(queueId) .offset(offset) .commitLogOffset(0L) .build(); @@ -100,16 +101,16 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { ackResult.setStatus(AckStatus.OK); ackResult.setExtraInfo(newReceiptHandle); Mockito.when(messagingProcessor.changeInvisibleTime(Mockito.any(ProxyContext.class), Mockito.any(ReceiptHandle.class), Mockito.eq(messageId), - Mockito.eq(group), Mockito.eq(messageQueue.getTopic()), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis()))) + Mockito.eq(group), Mockito.eq(topic), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis()))) .thenReturn(CompletableFuture.completedFuture(ackResult)); receiptHandleProcessor.scheduleRenewTask(); Mockito.verify(messagingProcessor, Mockito.timeout(1000).times(1)) .changeInvisibleTime(Mockito.any(ProxyContext.class), Mockito.argThat((r) -> r.getInvisibleTime() == invisibleTime), Mockito.eq(messageId), - Mockito.eq(group), Mockito.eq(messageQueue.getTopic()), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis())); + Mockito.eq(group), Mockito.eq(topic), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis())); receiptHandleProcessor.scheduleRenewTask(); Mockito.verify(messagingProcessor, Mockito.timeout(1000).times(1)) .changeInvisibleTime(Mockito.any(ProxyContext.class), Mockito.argThat((r) -> r.getInvisibleTime() == newInvisibleTime), Mockito.eq(messageId), - Mockito.eq(group), Mockito.eq(messageQueue.getTopic()), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis())); + Mockito.eq(group), Mockito.eq(topic), Mockito.eq(ConfigurationManager.getProxyConfig().getRenewSliceTimeMillis())); } @Test @@ -121,12 +122,12 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { .invisibleTime(newInvisibleTime) .reviveQueueId(1) .topicType(ReceiptHandle.NORMAL_TOPIC) - .brokerName(messageQueue.getBrokerName()) - .queueId(messageQueue.getQueueId()) + .brokerName(brokerName) + .queueId(queueId) .offset(offset) .commitLogOffset(0L) .build().encode(); - messageReceiptHandle = new MessageReceiptHandle(group, messageQueue, receiptHandle, messageId, offset, + messageReceiptHandle = new MessageReceiptHandle(group, topic, queueId, receiptHandle, messageId, offset, reconsumeTimes, newInvisibleTime); String channelId = context.getVal(ContextVariable.CHANNEL_KEY); receiptHandleProcessor.addReceiptHandle(channelId, newReceiptHandle, messageReceiptHandle); @@ -135,7 +136,7 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { receiptHandleProcessor.scheduleRenewTask(); Mockito.verify(messagingProcessor, Mockito.timeout(1000).times(1)) .changeInvisibleTime(Mockito.any(ProxyContext.class), Mockito.any(ReceiptHandle.class), Mockito.eq(messageId), - Mockito.eq(group), Mockito.eq(messageQueue.getTopic()), Mockito.eq(groupConfig.getGroupRetryPolicy().getRetryPolicy().nextDelayDuration(reconsumeTimes, TimeUnit.MILLISECONDS))); + Mockito.eq(group), Mockito.eq(topic), Mockito.eq(groupConfig.getGroupRetryPolicy().getRetryPolicy().nextDelayDuration(reconsumeTimes, TimeUnit.MILLISECONDS))); } @@ -147,12 +148,12 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { .invisibleTime(invisibleTime) .reviveQueueId(1) .topicType(ReceiptHandle.NORMAL_TOPIC) - .brokerName(messageQueue.getBrokerName()) - .queueId(messageQueue.getQueueId()) + .brokerName(brokerName) + .queueId(queueId) .offset(offset) .commitLogOffset(0L) .build().encode(); - messageReceiptHandle = new MessageReceiptHandle(group, messageQueue, newReceiptHandle, messageId, offset, + messageReceiptHandle = new MessageReceiptHandle(group, topic, queueId, newReceiptHandle, messageId, offset, reconsumeTimes, invisibleTime); String channelId = context.getVal(ContextVariable.CHANNEL_KEY); receiptHandleProcessor.addReceiptHandle(channelId, newReceiptHandle, messageReceiptHandle); @@ -187,6 +188,6 @@ public class ReceiptHandleProcessorTest extends BaseProcessorTest { receiptHandleProcessor.scheduleRenewTask(); Mockito.verify(messagingProcessor, Mockito.timeout(1000).times(1)) .changeInvisibleTime(Mockito.any(ProxyContext.class), Mockito.any(ReceiptHandle.class), Mockito.eq(messageId), - Mockito.eq(group), Mockito.eq(messageQueue.getTopic()), Mockito.eq(ConfigurationManager.getProxyConfig().getInvisibleTimeMillisWhenClear())); + Mockito.eq(group), Mockito.eq(topic), Mockito.eq(ConfigurationManager.getProxyConfig().getInvisibleTimeMillisWhenClear())); } } \ No newline at end of file