diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java index 8b0463bca6..a9474fd37e 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExtImpl.java @@ -171,7 +171,7 @@ public class MQClientAPIExtImpl { long timeoutMillis) { CompletableFuture 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) { diff --git a/common/src/main/java/org/apache/rocketmq/common/consumer/ReceiptHandle.java b/common/src/main/java/org/apache/rocketmq/common/consumer/ReceiptHandle.java index 20f573e514..ee7f9a83bf 100644 --- a/common/src/main/java/org/apache/rocketmq/common/consumer/ReceiptHandle.java +++ b/common/src/main/java/org/apache/rocketmq/common/consumer/ReceiptHandle.java @@ -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; + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java index e1bd1f35d6..d6cab5b5c3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/common/Converter.java @@ -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; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java index 0fa164afe4..89be32779e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcService.java @@ -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()) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java index 21b8edc980..d6975e6892 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java @@ -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 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 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 resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader, - ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); - resultFuture + if (request.getDeliveryAttempt() >= request.getMaxDeliveryAttempts()) { + CompletableFuture 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 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; } diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java index 6731ed8459..e473b039c9 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java @@ -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)