mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Implement ForwardMessageToDeadLetterQueue
This commit is contained in:
@@ -27,6 +27,7 @@ import apache.rocketmq.v1.DigestType;
|
||||
import apache.rocketmq.v1.Encoding;
|
||||
import apache.rocketmq.v1.FilterExpression;
|
||||
import apache.rocketmq.v1.FilterType;
|
||||
import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueRequest;
|
||||
import apache.rocketmq.v1.HeartbeatRequest;
|
||||
import apache.rocketmq.v1.Message;
|
||||
import apache.rocketmq.v1.MessageType;
|
||||
@@ -65,6 +66,7 @@ import org.apache.rocketmq.common.message.MessageExt;
|
||||
import org.apache.rocketmq.common.protocol.NamespaceUtil;
|
||||
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.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.ConsumeType;
|
||||
@@ -179,6 +181,23 @@ public class Converter {
|
||||
return changeInvisibleTimeRequestHeader;
|
||||
}
|
||||
|
||||
public static ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackRequestHeader(
|
||||
ForwardMessageToDeadLetterQueueRequest 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.getOffset());
|
||||
consumerSendMsgBackRequestHeader.setGroup(groupName);
|
||||
consumerSendMsgBackRequestHeader.setDelayLevel(-1);
|
||||
consumerSendMsgBackRequestHeader.setOriginMsgId(request.getMessageId());
|
||||
consumerSendMsgBackRequestHeader.setOriginMsgId(topicName);
|
||||
consumerSendMsgBackRequestHeader.setMaxReconsumeTimes(request.getMaxDeliveryAttempts());
|
||||
return consumerSendMsgBackRequestHeader;
|
||||
}
|
||||
|
||||
public static Map<String, String> buildMessageProperty(Message message) {
|
||||
org.apache.rocketmq.common.message.Message messageWithHeader = new org.apache.rocketmq.common.message.Message();
|
||||
// set user properties
|
||||
|
||||
@@ -66,10 +66,12 @@ import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.protocol.RequestCode;
|
||||
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.header.SendMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData;
|
||||
import org.apache.rocketmq.proxy.channel.ChannelManager;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannel;
|
||||
import org.apache.rocketmq.proxy.channel.SimpleChannelHandlerContext;
|
||||
import org.apache.rocketmq.proxy.config.ConfigurationManager;
|
||||
import org.apache.rocketmq.proxy.grpc.adapter.InvocationContext;
|
||||
@@ -252,7 +254,32 @@ public class LocalGrpcService implements GrpcForwardService {
|
||||
@Override
|
||||
public CompletableFuture<ForwardMessageToDeadLetterQueueResponse> forwardMessageToDeadLetterQueue(Context ctx,
|
||||
ForwardMessageToDeadLetterQueueRequest request) {
|
||||
return null;
|
||||
SimpleChannel channel = channelManager.createChannel();
|
||||
SimpleChannelHandlerContext channelHandlerContext = new SimpleChannelHandlerContext(channel);
|
||||
|
||||
ConsumerSendMsgBackRequestHeader requestHeader = Converter.buildConsumerSendMsgBackRequestHeader(request);
|
||||
RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader);
|
||||
command.makeCustomHeaderToNet();
|
||||
|
||||
CompletableFuture<ForwardMessageToDeadLetterQueueResponse> future = new CompletableFuture<>();
|
||||
try {
|
||||
CompletableFuture<RemotingCommand> processorFuture = brokerController.getSendMessageProcessor()
|
||||
.asyncProcessRequest(channelHandlerContext, command);
|
||||
processorFuture.thenAccept(r -> {
|
||||
ForwardMessageToDeadLetterQueueResponse.Builder builder = ForwardMessageToDeadLetterQueueResponse.newBuilder();
|
||||
if (null != r) {
|
||||
builder.setCommon(ResponseBuilder.buildCommon(r.getCode(), r.getRemark()));
|
||||
} else {
|
||||
builder.setCommon(ResponseBuilder.buildCommon(Code.INTERNAL, "Response command is null"));
|
||||
}
|
||||
ForwardMessageToDeadLetterQueueResponse response = builder.build();
|
||||
future.complete(response);
|
||||
});
|
||||
} catch (Exception e) {
|
||||
LOGGER.error("Exception raised when forwardMessageToDeadLetterQueue", e);
|
||||
future.completeExceptionally(e);
|
||||
}
|
||||
return future;
|
||||
}
|
||||
|
||||
@Override public CompletableFuture<EndTransactionResponse> endTransaction(Context ctx, EndTransactionRequest request) {
|
||||
|
||||
Reference in New Issue
Block a user