diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java index 4a4b176cbc..e0d5a4ee20 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java @@ -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 diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java index 6d076490f4..8b54b8b80f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/handler/ReceiveMessageResponseHandler.java @@ -96,10 +96,14 @@ public class ReceiveMessageResponseHandler implements ResponseHandler { + 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 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