diff --git a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java index 4452bbdfa1..b5ba1cbceb 100644 --- a/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java +++ b/client/src/main/java/org/apache/rocketmq/client/ClientConfig.java @@ -34,6 +34,8 @@ import org.apache.rocketmq.remoting.protocol.LanguageCode; */ public class ClientConfig { public static final String SEND_MESSAGE_WITH_VIP_CHANNEL_PROPERTY = "com.rocketmq.sendMessageWithVIPChannel"; + public static final String DECODE_READ_BODY = "com.rocketmq.read.body"; + public static final String DECODE_DECOMPRESS_BODY = "com.rocketmq.decompress.body"; private String namesrvAddr = NameServerAddressUtils.getNameServerAddresses(); private String clientIP = RemotingUtil.getLocalAddress(); private String instanceName = System.getProperty("rocketmq.client.name", "DEFAULT"); @@ -57,6 +59,8 @@ public class ClientConfig { private long pullTimeDelayMillsWhenException = 1000; private boolean unitMode = false; private String unitName; + private boolean decodeReadBody = Boolean.parseBoolean(System.getProperty(DECODE_READ_BODY, "true")); + private boolean decodeDecompressBody = Boolean.parseBoolean(System.getProperty(DECODE_DECOMPRESS_BODY, "true")); private boolean vipChannelEnabled = Boolean.parseBoolean(System.getProperty(SEND_MESSAGE_WITH_VIP_CHANNEL_PROPERTY, "false")); private boolean useTLS = TlsSystemConfig.tlsEnable; @@ -160,6 +164,8 @@ public class ClientConfig { this.namespace = cc.namespace; this.language = cc.language; this.mqClientApiTimeout = cc.mqClientApiTimeout; + this.decodeReadBody = cc.decodeReadBody; + this.decodeDecompressBody = cc.decodeDecompressBody; } public ClientConfig cloneClientConfig() { @@ -179,6 +185,8 @@ public class ClientConfig { cc.namespace = namespace; cc.language = language; cc.mqClientApiTimeout = mqClientApiTimeout; + cc.decodeReadBody = decodeReadBody; + cc.decodeDecompressBody = decodeDecompressBody; return cc; } @@ -279,6 +287,22 @@ public class ClientConfig { this.language = language; } + public boolean isDecodeReadBody() { + return decodeReadBody; + } + + public void setDecodeReadBody(boolean decodeReadBody) { + this.decodeReadBody = decodeReadBody; + } + + public boolean isDecodeDecompressBody() { + return decodeDecompressBody; + } + + public void setDecodeDecompressBody(boolean decodeDecompressBody) { + this.decodeDecompressBody = decodeDecompressBody; + } + public String getNamespace() { if (namespaceInitialized) { return namespace; @@ -324,6 +348,7 @@ public class ClientConfig { + ", clientCallbackExecutorThreads=" + clientCallbackExecutorThreads + ", pollNameServerInterval=" + pollNameServerInterval + ", heartbeatBrokerInterval=" + heartbeatBrokerInterval + ", persistConsumerOffsetInterval=" + persistConsumerOffsetInterval + ", pullTimeDelayMillsWhenException=" + pullTimeDelayMillsWhenException + ", unitMode=" + unitMode + ", unitName=" + unitName + ", vipChannelEnabled=" - + vipChannelEnabled + ", useTLS=" + useTLS + ", language=" + language.name() + ", namespace=" + namespace + ", mqClientApiTimeout=" + mqClientApiTimeout + "]"; + + vipChannelEnabled + ", useTLS=" + useTLS + ", language=" + language.name() + ", namespace=" + namespace + ", mqClientApiTimeout=" + mqClientApiTimeout + + ", decodeReadBody=" + decodeReadBody + ", decodeDecompressBody=" + decodeDecompressBody + "]"; } } diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java index 765184478f..c7d2f848c7 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java @@ -1034,7 +1034,11 @@ public class MQClientAPIImpl implements NameServerUpdateCallback { case ResponseCode.SUCCESS: popStatus = PopStatus.FOUND; ByteBuffer byteBuffer = ByteBuffer.wrap(response.getBody()); - msgFoundList = MessageDecoder.decodes(byteBuffer); + msgFoundList = MessageDecoder.decodesBatch( + byteBuffer, + clientConfig.isDecodeReadBody(), + clientConfig.isDecodeDecompressBody(), + true); break; case ResponseCode.POLLING_FULL: popStatus = PopStatus.POLLING_FULL; diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullAPIWrapper.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullAPIWrapper.java index 6ce8e261ca..9b0fa8df77 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullAPIWrapper.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/PullAPIWrapper.java @@ -78,7 +78,12 @@ public class PullAPIWrapper { this.updatePullFromWhichNode(mq, pullResultExt.getSuggestWhichBrokerId()); if (PullStatus.FOUND == pullResult.getPullStatus()) { ByteBuffer byteBuffer = ByteBuffer.wrap(pullResultExt.getMessageBinary()); - List msgList = MessageDecoder.decodes(byteBuffer); + List msgList = MessageDecoder.decodesBatch( + byteBuffer, + this.mQClientFactory.getClientConfig().isDecodeReadBody(), + this.mQClientFactory.getClientConfig().isDecodeDecompressBody(), + true + ); boolean needDecodeInnerMessage = false; for (MessageExt messageExt: msgList) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIFactory.java index abdfa5404a..0b813ae608 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/mqclient/MQClientAPIFactory.java @@ -89,6 +89,8 @@ public class MQClientAPIFactory implements StartAndShutdown { protected MQClientAPIExt createAndStart(String instanceName) { ClientConfig clientConfig = new ClientConfig(); clientConfig.setInstanceName(instanceName); + clientConfig.setDecodeReadBody(true); + clientConfig.setDecodeDecompressBody(false); NettyClientConfig nettyClientConfig = new NettyClientConfig(); nettyClientConfig.setDisableCallbackExecutor(true); diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcIT.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcIT.java index 8fcd9f3323..fa0a6ca7e8 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/ClusterGrpcIT.java @@ -69,7 +69,7 @@ public class ClusterGrpcIT extends GrpcBaseIT { String topic = initTopic(); QueryRouteResponse response = blockingStub.queryRoute(buildQueryRouteRequest(topic)); - assertQueryRoute(response, brokerNum * defaultQueueNums); + assertQueryRoute(response, brokerNum * DEFAULT_QUEUE_NUMS); } @Test @@ -87,6 +87,12 @@ public class ClusterGrpcIT extends GrpcBaseIT { super.testTransactionCheckThenCommit(); } + + @Test + public void testSimpleConsumerSendAndRecvBigMessage() throws Exception { + super.testSimpleConsumerSendAndRecvBigMessage(); + } + @Test public void testSimpleConsumerSendAndRecv() throws Exception { super.testSimpleConsumerSendAndRecv(); diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java index f1f79bfcdf..85ae8f1faf 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/GrpcBaseIT.java @@ -27,6 +27,7 @@ import apache.rocketmq.v2.ChangeInvisibleDurationRequest; import apache.rocketmq.v2.ChangeInvisibleDurationResponse; import apache.rocketmq.v2.ClientType; import apache.rocketmq.v2.Code; +import apache.rocketmq.v2.Encoding; import apache.rocketmq.v2.EndTransactionRequest; import apache.rocketmq.v2.EndTransactionResponse; import apache.rocketmq.v2.Endpoints; @@ -101,6 +102,7 @@ import org.apache.rocketmq.proxy.grpc.v2.common.ResponseBuilder; import org.apache.rocketmq.remoting.common.RemotingUtil; import org.apache.rocketmq.test.base.BaseConf; import org.apache.rocketmq.test.util.MQRandomUtils; +import org.apache.rocketmq.test.util.RandomUtils; import org.junit.Rule; import static org.apache.rocketmq.common.message.MessageClientIDSetter.createUniqID; @@ -110,7 +112,7 @@ import static org.awaitility.Awaitility.await; public class GrpcBaseIT extends BaseConf { - protected final int PORT = 8082; + protected final int port = 8082; /** * This rule manages automatic graceful shutdown for the registered servers and channels at the end of test. */ @@ -121,7 +123,7 @@ public class GrpcBaseIT extends BaseConf { protected MessagingServiceGrpc.MessagingServiceStub stub; protected final Metadata header = new Metadata(); - protected static final int defaultQueueNums = 8; + protected static final int DEFAULT_QUEUE_NUMS = 8; public void setUp() throws Exception { brokerController1.getBrokerConfig().setTransactionCheckInterval(3 * 1000); @@ -139,7 +141,7 @@ public class GrpcBaseIT extends BaseConf { System.setProperty(RMQ_PROXY_HOME, mockProxyHome); ConfigurationManager.initEnv(); ConfigurationManager.intConfig(); - ConfigurationManager.getProxyConfig().setGrpcServerPort(PORT); + ConfigurationManager.getProxyConfig().setGrpcServerPort(port); ConfigurationManager.getProxyConfig().setNameSrvAddr(nsAddr); // Set LongPollingReserveTimeInMillis to 500ms to reserve more time for IT ConfigurationManager.getProxyConfig().setLongPollingReserveTimeInMillis(500); @@ -303,6 +305,30 @@ public class GrpcBaseIT extends BaseConf { .build(); } + public void testSimpleConsumerSendAndRecvBigMessage() throws Exception { + String topic = initTopicOnSampleTopicBroker(broker1Name); + String group = MQRandomUtils.getRandomConsumerGroup(); + + int maxDeliveryAttempts = 16; + boolean fifo = false; + int bodySize = 4 * 1024; + + // init consumer offset + this.sendClientSettings(stub, buildSimpleConsumerClientSettings(maxDeliveryAttempts, fifo)).get(); + receiveMessage(blockingStub, topic, group, 1); + + this.sendClientSettings(stub, buildProducerClientSettings(topic)).get(); + String messageId = createUniqID(); + SendMessageResponse sendResponse = blockingStub.sendMessage(buildSendBigMessageRequest(topic, messageId, bodySize)); + assertSendMessage(sendResponse, messageId); + + this.sendClientSettings(stub, buildSimpleConsumerClientSettings(maxDeliveryAttempts, fifo)).get(); + + Message message = assertAndGetReceiveMessage(receiveMessage(blockingStub, topic, group), messageId); + assertThat(message.getSystemProperties().getBodyEncoding()).isEqualTo(Encoding.GZIP); + assertThat(message.getBody().size()).isEqualTo(bodySize); + } + public void testSimpleConsumerSendAndRecv() throws Exception { String topic = initTopicOnSampleTopicBroker(broker1Name); String group = MQRandomUtils.getRandomConsumerGroup(); @@ -432,7 +458,7 @@ public class GrpcBaseIT extends BaseConf { public QueryRouteRequest buildQueryRouteRequest(String topic) { return QueryRouteRequest.newBuilder() - .setEndpoints(buildEndpoints(PORT)) + .setEndpoints(buildEndpoints(port)) .setTopic(Resource.newBuilder() .setName(topic) .build()) @@ -441,7 +467,7 @@ public class GrpcBaseIT extends BaseConf { public QueryAssignmentRequest buildQueryAssignmentRequest(String topic, String group) { return QueryAssignmentRequest.newBuilder() - .setEndpoints(buildEndpoints(PORT)) + .setEndpoints(buildEndpoints(port)) .setTopic(Resource.newBuilder().setName(topic).build()) .setGroup(Resource.newBuilder().setName(group).build()) .build(); @@ -465,6 +491,25 @@ public class GrpcBaseIT extends BaseConf { .build(); } + public SendMessageRequest buildSendBigMessageRequest(String topic, String messageId, int messageSize) { + return SendMessageRequest.newBuilder() + .addMessages(Message.newBuilder() + .setTopic(Resource.newBuilder() + .setName(topic) + .build()) + .setSystemProperties(SystemProperties.newBuilder() + .setMessageId(messageId) + .setQueueId(0) + .setMessageType(MessageType.NORMAL) + .setBodyEncoding(Encoding.GZIP) + .setBornTimestamp(Timestamps.fromMillis(System.currentTimeMillis())) + .setBornHost(StringUtils.defaultString(RemotingUtil.getLocalAddress(), "127.0.0.1:1234")) + .build()) + .setBody(ByteString.copyFromUtf8(RandomUtils.getStringWithCharacter(messageSize))) + .build()) + .build(); + } + public SendMessageRequest buildTransactionSendMessageRequest(String topic, String messageId) { return SendMessageRequest.newBuilder() .addMessages(Message.newBuilder() diff --git a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcIT.java b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcIT.java index 5aa188329e..fc113370bc 100644 --- a/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcIT.java +++ b/test/src/test/java/org/apache/rocketmq/test/grpc/v2/LocalGrpcIT.java @@ -57,7 +57,7 @@ public class LocalGrpcIT extends GrpcBaseIT { String topic = initTopic(); QueryRouteResponse response = blockingStub.queryRoute(buildQueryRouteRequest(topic)); - assertQueryRoute(response, brokerControllerList.size() * defaultQueueNums); + assertQueryRoute(response, brokerControllerList.size() * DEFAULT_QUEUE_NUMS); } @Test @@ -75,6 +75,11 @@ public class LocalGrpcIT extends GrpcBaseIT { super.testTransactionCheckThenCommit(); } + @Test + public void testSimpleConsumerSendAndRecvBigMessage() throws Exception { + super.testSimpleConsumerSendAndRecvBigMessage(); + } + @Test public void testSimpleConsumerSendAndRecv() throws Exception { super.testSimpleConsumerSendAndRecv();