[ISSUE #3949] forward message to dlq

This commit is contained in:
kaiyi.lk
2022-07-13 11:29:12 +08:00
committed by zhouxiang
parent 3f07761f60
commit 2cea9131b1
6 changed files with 96 additions and 43 deletions
@@ -171,7 +171,7 @@ public class MQClientAPIExtImpl {
long timeoutMillis) {
CompletableFuture<RemotingCommand> future = new CompletableFuture<>();
try {
RemotingCommand request = RemotingCommand.createResponseCommandWithHeader(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader);
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader);
this.getRemotingClient().invokeAsync(brokerAddr, request, timeoutMillis, responseFuture -> {
RemotingCommand response = responseFuture.getResponseCommand();
if (response != null) {
@@ -19,7 +19,7 @@ package org.apache.rocketmq.common.consumer;
import java.util.Arrays;
import java.util.List;
import org.apache.rocketmq.common.MixAll;
import org.apache.rocketmq.common.KeyBuilder;
import org.apache.rocketmq.common.message.MessageConst;
import org.apache.rocketmq.common.message.MessageExt;
@@ -32,7 +32,7 @@ public class ReceiptHandle {
private final long invisibleTime;
private final long nextVisibleTime;
private final int reviveQueueId;
private final String topic;
private final String topicType;
private final String brokerName;
private final int queueId;
private final long offset;
@@ -40,12 +40,8 @@ public class ReceiptHandle {
private final String receiptHandle;
public String encode() {
String t = NORMAL_TOPIC;
if (topic.startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
t = RETRY_TOPIC;
}
return startOffset + SEPARATOR + retrieveTime + SEPARATOR + invisibleTime + SEPARATOR + reviveQueueId
+ SEPARATOR + t + SEPARATOR + brokerName + SEPARATOR + queueId + SEPARATOR + offset + SEPARATOR
+ SEPARATOR + topicType + SEPARATOR + brokerName + SEPARATOR + queueId + SEPARATOR + offset + SEPARATOR
+ commitLogOffset;
}
@@ -84,7 +80,7 @@ public class ReceiptHandle {
.retrieveTime(retrieveTime)
.invisibleTime(invisibleTime)
.reviveQueueId(reviveQueueId)
.topic(topic)
.topicType(topic)
.brokerName(brokerName)
.queueId(queueId)
.offset(offset)
@@ -94,14 +90,14 @@ public class ReceiptHandle {
}
ReceiptHandle(final long startOffset, final long retrieveTime, final long invisibleTime, final long nextVisibleTime,
final int reviveQueueId, final String topic, final String brokerName, final int queueId, final long offset,
final int reviveQueueId, final String topicType, final String brokerName, final int queueId, final long offset,
final long commitLogOffset, final String receiptHandle) {
this.startOffset = startOffset;
this.retrieveTime = retrieveTime;
this.invisibleTime = invisibleTime;
this.nextVisibleTime = nextVisibleTime;
this.reviveQueueId = reviveQueueId;
this.topic = topic;
this.topicType = topicType;
this.brokerName = brokerName;
this.queueId = queueId;
this.offset = offset;
@@ -115,12 +111,11 @@ public class ReceiptHandle {
private long invisibleTime;
private long nextVisibleTime;
private int reviveQueueId;
private String topic;
private String topicType;
private String brokerName;
private int queueId;
private long offset;
private long commitLogOffset;
private String type;
private String receiptHandle;
ReceiptHandleBuilder() {
@@ -151,8 +146,8 @@ public class ReceiptHandle {
return this;
}
public ReceiptHandle.ReceiptHandleBuilder topic(final String topic) {
this.topic = topic;
public ReceiptHandle.ReceiptHandleBuilder topicType(final String topic) {
this.topicType = topic;
return this;
}
@@ -176,11 +171,6 @@ public class ReceiptHandle {
return this;
}
public ReceiptHandle.ReceiptHandleBuilder type(final String type) {
this.type = type;
return this;
}
public ReceiptHandle.ReceiptHandleBuilder receiptHandle(final String receiptHandle) {
this.receiptHandle = receiptHandle;
return this;
@@ -188,12 +178,12 @@ public class ReceiptHandle {
public ReceiptHandle build() {
return new ReceiptHandle(this.startOffset, this.retrieveTime, this.invisibleTime, this.nextVisibleTime,
this.reviveQueueId, this.topic, this.brokerName, this.queueId, this.offset, this.commitLogOffset, this.receiptHandle);
this.reviveQueueId, this.topicType, this.brokerName, this.queueId, this.offset, this.commitLogOffset, this.receiptHandle);
}
@java.lang.Override
public java.lang.String toString() {
return "ReceiptHandle.ReceiptHandleBuilder(startOffset=" + this.startOffset + ", retrieveTime=" + this.retrieveTime + ", invisibleTime=" + this.invisibleTime + ", nextVisibleTime=" + this.nextVisibleTime + ", reviveQueueId=" + this.reviveQueueId + ", topic=" + this.topic + ", brokerName=" + this.brokerName + ", queueId=" + this.queueId + ", offset=" + this.offset + ", commitLogOffset=" + this.commitLogOffset + ", type=" + this.type + ", receiptHandle=" + this.receiptHandle + ")";
@Override
public String toString() {
return "ReceiptHandle.ReceiptHandleBuilder(startOffset=" + this.startOffset + ", retrieveTime=" + this.retrieveTime + ", invisibleTime=" + this.invisibleTime + ", nextVisibleTime=" + this.nextVisibleTime + ", reviveQueueId=" + this.reviveQueueId + ", topic=" + this.topicType + ", brokerName=" + this.brokerName + ", queueId=" + this.queueId + ", offset=" + this.offset + ", commitLogOffset=" + this.commitLogOffset + ", receiptHandle=" + this.receiptHandle + ")";
}
}
@@ -221,8 +211,8 @@ public class ReceiptHandle {
return this.reviveQueueId;
}
public String getTopic() {
return this.topic;
public String getTopicType() {
return this.topicType;
}
public String getBrokerName() {
@@ -244,4 +234,15 @@ public class ReceiptHandle {
public String getReceiptHandle() {
return this.receiptHandle;
}
public boolean isRetryTopic() {
return RETRY_TOPIC.equals(topicType);
}
public String getRealTopic(String topic, String groupName) {
if (isRetryTopic()) {
return KeyBuilder.buildPopRetryTopic(topic, groupName);
}
return topic;
}
}
@@ -179,7 +179,7 @@ public class Converter {
AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader();
ackMessageRequestHeader.setConsumerGroup(groupName);
ackMessageRequestHeader.setTopic(topicName);
ackMessageRequestHeader.setTopic(handle.getRealTopic(topicName, groupName));
ackMessageRequestHeader.setQueueId(handle.getQueueId());
ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle());
ackMessageRequestHeader.setOffset(handle.getOffset());
@@ -195,7 +195,7 @@ public class Converter {
ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader();
changeInvisibleTimeRequestHeader.setConsumerGroup(groupName);
changeInvisibleTimeRequestHeader.setTopic(topicName);
changeInvisibleTimeRequestHeader.setTopic(handle.getRealTopic(topicName, groupName));
changeInvisibleTimeRequestHeader.setQueueId(handle.getQueueId());
changeInvisibleTimeRequestHeader.setExtraInfo(handle.getReceiptHandle());
changeInvisibleTimeRequestHeader.setOffset(handle.getOffset());
@@ -213,7 +213,7 @@ public class Converter {
ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader();
changeInvisibleTimeRequestHeader.setConsumerGroup(groupName);
changeInvisibleTimeRequestHeader.setTopic(topicName);
changeInvisibleTimeRequestHeader.setTopic(handle.getRealTopic(topicName, groupName));
changeInvisibleTimeRequestHeader.setQueueId(handle.getQueueId());
changeInvisibleTimeRequestHeader.setExtraInfo(handle.getReceiptHandle());
changeInvisibleTimeRequestHeader.setOffset(handle.getOffset());
@@ -233,7 +233,24 @@ public class Converter {
consumerSendMsgBackRequestHeader.setGroup(groupName);
consumerSendMsgBackRequestHeader.setDelayLevel(-1);
consumerSendMsgBackRequestHeader.setOriginMsgId(request.getMessageId());
consumerSendMsgBackRequestHeader.setOriginMsgId(topicName);
consumerSendMsgBackRequestHeader.setOriginTopic(handle.getRealTopic(topicName, groupName));
consumerSendMsgBackRequestHeader.setMaxReconsumeTimes(request.getMaxDeliveryAttempts());
return consumerSendMsgBackRequestHeader;
}
public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader(
NackMessageRequest request) {
String groupName = Converter.getResourceNameWithNamespace(request.getGroup());
String topicName = Converter.getResourceNameWithNamespace(request.getTopic());
String receiptHandleStr = request.getReceiptHandle();
ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr);
ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader();
consumerSendMsgBackRequestHeader.setOffset(handle.getCommitLogOffset());
consumerSendMsgBackRequestHeader.setGroup(groupName);
consumerSendMsgBackRequestHeader.setDelayLevel(-1);
consumerSendMsgBackRequestHeader.setOriginMsgId(request.getMessageId());
consumerSendMsgBackRequestHeader.setOriginTopic(handle.getRealTopic(topicName, groupName));
consumerSendMsgBackRequestHeader.setMaxReconsumeTimes(request.getMaxDeliveryAttempts());
return consumerSendMsgBackRequestHeader;
}
@@ -514,6 +514,7 @@ public class LocalGrpcService implements GrpcForwardService {
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
ChangeInvisibleTimeRequestHeader requestHeader = Converter.buildChangeInvisibleTimeRequestHeader(request);
ReceiptHandle receiptHandle = ReceiptHandle.decode(request.getReceiptHandle());
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CHANGE_MESSAGE_INVISIBLETIME, requestHeader);
command.makeCustomHeaderToNet();
@@ -530,7 +531,7 @@ public class LocalGrpcService implements GrpcForwardService {
.retrieveTime(responseHeader.getPopTime())
.invisibleTime(responseHeader.getInvisibleTime())
.reviveQueueId(responseHeader.getReviveQid())
.topic(Converter.getResourceNameWithNamespace(request.getTopic()))
.topicType(receiptHandle.getTopicType())
.brokerName(brokerController.getBrokerConfig().getBrokerName())
.queueId(requestHeader.getQueueId())
.offset(requestHeader.getOffset())
@@ -33,12 +33,14 @@ import org.apache.rocketmq.common.consumer.ReceiptHandle;
import org.apache.rocketmq.common.message.MessageExt;
import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader;
import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader;
import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader;
import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader;
import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData;
import org.apache.rocketmq.proxy.common.utils.FilterUtil;
import org.apache.rocketmq.proxy.common.utils.FilterUtils;
import org.apache.rocketmq.proxy.common.utils.ProxyUtils;
import org.apache.rocketmq.proxy.config.ConfigurationManager;
import org.apache.rocketmq.proxy.connector.ConnectorManager;
import org.apache.rocketmq.proxy.connector.ForwardProducer;
import org.apache.rocketmq.proxy.connector.ForwardReadConsumer;
import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer;
import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue;
@@ -51,11 +53,13 @@ import java.util.ArrayList;
import java.util.List;
import java.util.concurrent.CompletableFuture;
import java.util.concurrent.TimeUnit;
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
public class ConsumerService extends BaseService {
private final ForwardReadConsumer readConsumer;
private final ForwardWriteConsumer writeConsumer;
private final ForwardProducer producer;
private volatile ReadQueueSelector readQueueSelector;
private volatile ResponseHook<ReceiveMessageRequest, ReceiveMessageResponse> receiveMessageHook = null;
@@ -69,6 +73,7 @@ public class ConsumerService extends BaseService {
super(connectorManager);
this.readConsumer = connectorManager.getForwardReadConsumer();
this.writeConsumer = connectorManager.getForwardWriteConsumer();
this.producer = connectorManager.getForwardProducer();
this.readQueueSelector = new DefaultReadQueueSelector(connectorManager.getTopicRouteCache());
this.delayPolicy = DelayPolicy.build(ConfigurationManager.getProxyConfig().getMessageDelayLevel());
@@ -143,7 +148,7 @@ public class ConsumerService extends BaseService {
List<Message> messages = new ArrayList<>();
for (MessageExt messageExt : result.getMsgFoundList()) {
if (FilterUtil.isTagNotMatched(subscriptionData.getTagsSet(), messageExt.getTags())) {
if (!FilterUtils.isTagMatched(subscriptionData.getTagsSet(), messageExt.getTags())) {
this.ackNoMatchedMessage(ctx, request, messageExt);
continue;
}
@@ -239,10 +244,29 @@ public class ConsumerService extends BaseService {
ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle());
String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName());
ChangeInvisibleTimeRequestHeader requestHeader = this.convertToChangeInvisibleTimeRequestHeader(ctx, request);
CompletableFuture<AckResult> resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader,
ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
resultFuture
if (request.getDeliveryAttempt() >= request.getMaxDeliveryAttempts()) {
CompletableFuture<RemotingCommand> resultFuture = this.producer.sendMessageBack(
brokerAddr,
this.convertToConsumerSendMsgBackToDLQRequestHeader(ctx, request),
ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
resultFuture
.thenAccept(result -> {
try {
future.complete(convertToNackMessageResponse(ctx, request, result));
} catch (Throwable throwable) {
future.completeExceptionally(throwable);
}
})
.exceptionally(throwable -> {
throwable.printStackTrace();
future.completeExceptionally(throwable);
return null;
});
} else {
ChangeInvisibleTimeRequestHeader requestHeader = this.convertToChangeInvisibleTimeRequestHeader(ctx, request);
CompletableFuture<AckResult> resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader,
ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT);
resultFuture
.thenAccept(result -> {
try {
future.complete(convertToNackMessageResponse(ctx, request, result));
@@ -254,6 +278,7 @@ public class ConsumerService extends BaseService {
future.completeExceptionally(throwable);
return null;
});
}
} catch (Throwable t) {
future.completeExceptionally(t);
}
@@ -264,6 +289,10 @@ public class ConsumerService extends BaseService {
return Converter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy);
}
protected ConsumerSendMsgBackRequestHeader convertToConsumerSendMsgBackToDLQRequestHeader(Context ctx, NackMessageRequest request) {
return Converter.buildConsumerSendMsgBackToDLQRequestHeader(request);
}
protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, AckResult ackResult) {
if (AckStatus.OK.equals(ackResult.getStatus())) {
return NackMessageResponse.newBuilder()
@@ -275,8 +304,13 @@ public class ConsumerService extends BaseService {
.build();
}
public void setReadQueueSelector(
ReadQueueSelector readQueueSelector) {
protected NackMessageResponse convertToNackMessageResponse(Context ctx, NackMessageRequest request, RemotingCommand sendMsgBackToDLQResult) {
return NackMessageResponse.newBuilder()
.setCommon(ResponseBuilder.buildCommon(sendMsgBackToDLQResult.getCode(), sendMsgBackToDLQResult.getRemark()))
.build();
}
public void setReadQueueSelector(ReadQueueSelector readQueueSelector) {
this.readQueueSelector = readQueueSelector;
}
@@ -295,7 +295,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
.invisibleTime(1000L)
.nextVisibleTime(1000L)
.reviveQueueId(0)
.topic("topic")
.topicType("topic")
.brokerName("brokerName")
.queueId(0)
.offset(0L)
@@ -328,7 +328,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
.invisibleTime(1000L)
.nextVisibleTime(1000L)
.reviveQueueId(0)
.topic("topic")
.topicType("topic")
.brokerName("brokerName")
.queueId(0)
.offset(0L)
@@ -357,7 +357,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
.invisibleTime(1000L)
.nextVisibleTime(1000L)
.reviveQueueId(0)
.topic("topic")
.topicType("topic")
.brokerName("brokerName")
.queueId(0)
.offset(0L)
@@ -511,7 +511,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest {
.invisibleTime(invisibleTime)
.nextVisibleTime(1000L)
.reviveQueueId(0)
.topic("topic")
.topicType("topic")
.brokerName("brokerName")
.queueId(queueId)
.offset(offset)