mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
[ISSUE #3949] Add topic and queueId for MessageReceiptHandle
This commit is contained in:
@@ -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() {
|
||||
|
||||
+1
-3
@@ -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);
|
||||
|
||||
@@ -134,7 +134,7 @@ public class ReceiptHandleProcessor extends AbstractStartAndShutdown {
|
||||
if (current - messageReceiptHandle.getTimestamp() < messageReceiptHandle.getExpectInvisibleTime()) {
|
||||
CompletableFuture<AckResult> 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<AckResult> 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()
|
||||
);
|
||||
});
|
||||
|
||||
+20
-19
@@ -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()));
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user