mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-21 13:49:50 +08:00
[ISSUE #3949] not decompress body in proxy
This commit is contained in:
@@ -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 + "]";
|
||||
}
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -78,7 +78,12 @@ public class PullAPIWrapper {
|
||||
this.updatePullFromWhichNode(mq, pullResultExt.getSuggestWhichBrokerId());
|
||||
if (PullStatus.FOUND == pullResult.getPullStatus()) {
|
||||
ByteBuffer byteBuffer = ByteBuffer.wrap(pullResultExt.getMessageBinary());
|
||||
List<MessageExt> msgList = MessageDecoder.decodes(byteBuffer);
|
||||
List<MessageExt> msgList = MessageDecoder.decodesBatch(
|
||||
byteBuffer,
|
||||
this.mQClientFactory.getClientConfig().isDecodeReadBody(),
|
||||
this.mQClientFactory.getClientConfig().isDecodeDecompressBody(),
|
||||
true
|
||||
);
|
||||
|
||||
boolean needDecodeInnerMessage = false;
|
||||
for (MessageExt messageExt: msgList) {
|
||||
|
||||
Reference in New Issue
Block a user