mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] use the right exception.
This commit is contained in:
+13
-5
@@ -74,8 +74,12 @@ public class TransactionId {
|
||||
commitLogOffset, sendResult.getQueueOffset());
|
||||
}
|
||||
|
||||
public static TransactionId genByBrokerTransactionId(SocketAddress brokerAddr, String orgTransactionId,
|
||||
long commitLogOffset, long tranStateTableOffset) {
|
||||
public static TransactionId genByBrokerTransactionId(
|
||||
SocketAddress brokerAddr,
|
||||
String orgTransactionId,
|
||||
long commitLogOffset,
|
||||
long tranStateTableOffset
|
||||
) {
|
||||
byte[] orgTransactionIdByte = new byte[0];
|
||||
if (StringUtils.isNotBlank(orgTransactionId)) {
|
||||
orgTransactionIdByte = orgTransactionId.getBytes(StandardCharsets.UTF_8);
|
||||
@@ -127,9 +131,13 @@ public class TransactionId {
|
||||
.build();
|
||||
}
|
||||
|
||||
public static long generateCommitLogOffset(String messageId) throws UnknownHostException {
|
||||
MessageId id = MessageDecoder.decodeMessageId(messageId);
|
||||
return id.getOffset();
|
||||
public static long generateCommitLogOffset(String messageId) throws IllegalArgumentException {
|
||||
try {
|
||||
MessageId id = MessageDecoder.decodeMessageId(messageId);
|
||||
return id.getOffset();
|
||||
} catch (UnknownHostException e) {
|
||||
throw new IllegalArgumentException("illegal messageId: " + messageId);
|
||||
}
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+19
-10
@@ -96,10 +96,14 @@ public class ReceiveMessageResponseHandler implements ResponseHandler<ReceiveMes
|
||||
// find pop ck offset
|
||||
String key = messageExt.getTopic() + messageExt.getQueueId();
|
||||
if (!map.containsKey(messageExt.getTopic() + messageExt.getQueueId())) {
|
||||
String extraInfo = ExtraInfoUtil.buildExtraInfo(messageExt.getQueueOffset(),
|
||||
responseHeader.getPopTime(), responseHeader.getInvisibleTime(),
|
||||
responseHeader.getReviveQid(), messageExt.getTopic(), brokerName,
|
||||
messageExt.getQueueId());
|
||||
String extraInfo = ExtraInfoUtil.buildExtraInfo(
|
||||
messageExt.getQueueOffset(),
|
||||
responseHeader.getPopTime(),
|
||||
responseHeader.getInvisibleTime(),
|
||||
responseHeader.getReviveQid(),
|
||||
messageExt.getTopic(), brokerName,
|
||||
messageExt.getQueueId()
|
||||
);
|
||||
map.put(key, extraInfo);
|
||||
}
|
||||
messageExt.getProperties().put(MessageConst.PROPERTY_POP_CK,
|
||||
@@ -113,10 +117,16 @@ public class ReceiveMessageResponseHandler implements ResponseHandler<ReceiveMes
|
||||
log.warn("Queue offset[{}] of msg is strange, not equal to the stored in msg, {}",
|
||||
msgQueueOffset, messageExt);
|
||||
}
|
||||
String extraInfo = ExtraInfoUtil.buildExtraInfo(startOffsetInfo.get(key),
|
||||
responseHeader.getPopTime(), responseHeader.getInvisibleTime(),
|
||||
responseHeader.getReviveQid(), messageExt.getTopic(),
|
||||
brokerName, messageExt.getQueueId(), msgQueueOffset);
|
||||
String extraInfo = ExtraInfoUtil.buildExtraInfo(
|
||||
startOffsetInfo.get(key),
|
||||
responseHeader.getPopTime(),
|
||||
responseHeader.getInvisibleTime(),
|
||||
responseHeader.getReviveQid(),
|
||||
messageExt.getTopic(),
|
||||
brokerName,
|
||||
messageExt.getQueueId(),
|
||||
msgQueueOffset
|
||||
);
|
||||
messageExt.setQueueOffset(msgQueueOffset);
|
||||
messageExt.getProperties().put(MessageConst.PROPERTY_POP_CK, extraInfo);
|
||||
if (fifo && orderCountInfo != null) {
|
||||
@@ -142,8 +152,7 @@ public class ReceiveMessageResponseHandler implements ResponseHandler<ReceiveMes
|
||||
}
|
||||
response = builder.build();
|
||||
long elapsed = stopWatch.stop().elapsed(TimeUnit.MILLISECONDS);
|
||||
log.debug("Translating remoting response to gRPC response costs {}ms. Duration request received: {}",
|
||||
elapsed, popCosts);
|
||||
log.debug("Translating remoting response to gRPC response costs {}ms. Duration request received: {}", elapsed, popCosts);
|
||||
future.complete(response);
|
||||
} catch (Exception e) {
|
||||
log.error("Unexpected exception raised when handle pop remoting command", e);
|
||||
|
||||
+14
-5
@@ -20,8 +20,8 @@ package org.apache.rocketmq.proxy.grpc.v2.adapter.handler;
|
||||
import apache.rocketmq.v2.SendMessageRequest;
|
||||
import apache.rocketmq.v2.SendMessageResponse;
|
||||
import apache.rocketmq.v2.SendReceipt;
|
||||
import java.net.UnknownHostException;
|
||||
import org.apache.commons.lang3.StringUtils;
|
||||
import org.apache.rocketmq.common.constant.LoggerName;
|
||||
import org.apache.rocketmq.common.protocol.ResponseCode;
|
||||
import org.apache.rocketmq.common.protocol.header.SendMessageResponseHeader;
|
||||
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
|
||||
@@ -30,8 +30,11 @@ import org.apache.rocketmq.proxy.connector.transaction.TransactionId;
|
||||
import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder;
|
||||
import org.apache.rocketmq.remoting.common.RemotingUtil;
|
||||
import org.apache.rocketmq.remoting.protocol.RemotingCommand;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public class SendMessageResponseHandler implements ResponseHandler<SendMessageRequest, SendMessageResponse> {
|
||||
private static final Logger log = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME);
|
||||
private final String messageId;
|
||||
private final int sysFlag;
|
||||
private final String localAddress;
|
||||
@@ -42,7 +45,8 @@ public class SendMessageResponseHandler implements ResponseHandler<SendMessageRe
|
||||
this.localAddress = localAddress;
|
||||
}
|
||||
|
||||
@Override public void handle(RemotingCommand responseCommand,
|
||||
@Override
|
||||
public void handle(RemotingCommand responseCommand,
|
||||
InvocationContext<SendMessageRequest, SendMessageResponse> context) {
|
||||
// If responseCommand equals to null, then the response has been written to channel.
|
||||
// org.apache.rocketmq.broker.processor.SendMessageProcessor#handlePutMessageResult
|
||||
@@ -55,11 +59,16 @@ public class SendMessageResponseHandler implements ResponseHandler<SendMessageRe
|
||||
long commitLogOffset = 0L;
|
||||
try {
|
||||
commitLogOffset = TransactionId.generateCommitLogOffset(responseHeader.getMsgId());
|
||||
} catch (UnknownHostException e) {
|
||||
} catch (IllegalArgumentException e) {
|
||||
log.warn("illegal messageId:{}", responseHeader.getMsgId());
|
||||
e.printStackTrace();
|
||||
}
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(RemotingUtil.string2SocketAddress(localAddress),
|
||||
responseHeader.getTransactionId(), commitLogOffset, responseHeader.getQueueOffset());
|
||||
TransactionId transactionId = TransactionId.genByBrokerTransactionId(
|
||||
RemotingUtil.string2SocketAddress(localAddress),
|
||||
responseHeader.getTransactionId(),
|
||||
commitLogOffset,
|
||||
responseHeader.getQueueOffset()
|
||||
);
|
||||
transactionIdString = transactionId.getProxyTransactionId();
|
||||
}
|
||||
SendMessageResponse response = SendMessageResponse.newBuilder()
|
||||
|
||||
Reference in New Issue
Block a user