[ISSUE #3905] Support bname in protocol for 5.0 client (#5334)

* [ISSUE #3905] Support bname in protocol for 5.0 client

* add bname for `CheckTransactionStateRequestHeader`, `ConsumerSendMsgBackRequestHeader`, `EndTransactionRequestHeader`,
`SendMessageRequestHeader`

* * add bname for `AckMessageRequestHeader`, `PeekMessageRequestHeader`, `PopMessageRequestHeader`,
    `ChangeInvisibleTimeRequestHeader`, `NotificationRequestHeader` and `PollingInfoRequestHeader`
This commit is contained in:
Zhouxiang Zhan
2022-10-17 18:28:19 +08:00
committed by GitHub
parent ff60d5c24a
commit a6d341d136
16 changed files with 50 additions and 30 deletions
@@ -56,6 +56,7 @@ public abstract class AbstractTransactionalMessageCheckListener {
checkTransactionStateRequestHeader.setMsgId(msgExt.getUserProperty(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX));
checkTransactionStateRequestHeader.setTransactionId(checkTransactionStateRequestHeader.getMsgId());
checkTransactionStateRequestHeader.setTranStateTableOffset(msgExt.getQueueOffset());
checkTransactionStateRequestHeader.setBname(brokerController.getBrokerConfig().getBrokerName());
msgExt.setTopic(msgExt.getUserProperty(MessageConst.PROPERTY_REAL_TOPIC));
msgExt.setQueueId(Integer.parseInt(msgExt.getUserProperty(MessageConst.PROPERTY_REAL_QUEUE_ID)));
msgExt.setStoreSize(0);
@@ -1451,6 +1451,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback {
public void consumerSendMessageBack(
final String addr,
final String brokerName,
final MessageExt msg,
final String consumerGroup,
final int delayLevel,
@@ -1466,6 +1467,7 @@ public class MQClientAPIImpl implements NameServerUpdateCallback {
requestHeader.setDelayLevel(delayLevel);
requestHeader.setOriginMsgId(msg.getMsgId());
requestHeader.setMaxReconsumeTimes(maxConsumeRetryTimes);
requestHeader.setBname(brokerName);
RemotingCommand response = this.remotingClient.invokeSync(MixAll.brokerVIPChannel(this.clientConfig.isVipChannelEnabled(), addr),
request, timeoutMillis);
@@ -622,8 +622,8 @@ public class DefaultMQPullConsumerImpl implements MQConsumerInner {
consumerGroup = this.defaultMQPullConsumer.getConsumerGroup();
}
this.mQClientFactory.getMQClientAPIImpl().consumerSendMessageBack(brokerAddr, msg, consumerGroup, delayLevel, 3000,
this.defaultMQPullConsumer.getMaxReconsumeTimes());
this.mQClientFactory.getMQClientAPIImpl().consumerSendMessageBack(brokerAddr, brokerName, msg, consumerGroup,
delayLevel, 3000, this.defaultMQPullConsumer.getMaxReconsumeTimes());
} catch (Exception e) {
log.error("sendMessageBack Exception, " + this.defaultMQPullConsumer.getConsumerGroup(), e);
@@ -732,7 +732,7 @@ public class DefaultMQPushConsumerImpl implements MQConsumerInner {
} else {
String brokerAddr = (null != brokerName) ? this.mQClientFactory.findBrokerAddressInPublish(brokerName)
: RemotingHelper.parseSocketAddressAddr(msg.getStoreHost());
this.mQClientFactory.getMQClientAPIImpl().consumerSendMessageBack(brokerAddr, msg,
this.mQClientFactory.getMQClientAPIImpl().consumerSendMessageBack(brokerAddr, brokerName, msg,
this.defaultMQPushConsumer.getConsumerGroup(), delayLevel, 5000, getMaxReconsumeTimes());
}
} catch (Exception e) {
@@ -794,6 +794,7 @@ public class DefaultMQPushConsumerImpl implements MQConsumerInner {
requestHeader.setOffset(queueOffset);
requestHeader.setConsumerGroup(consumerGroup);
requestHeader.setExtraInfo(extraInfo);
requestHeader.setBname(brokerName);
this.mQClientFactory.getMQClientAPIImpl().ackMessageAsync(findBrokerResult.getBrokerAddr(), ASYNC_TIMEOUT, new AckCallback() {
@Override
public void onSuccess(AckResult ackResult) {
@@ -837,6 +838,7 @@ public class DefaultMQPushConsumerImpl implements MQConsumerInner {
requestHeader.setConsumerGroup(consumerGroup);
requestHeader.setExtraInfo(extraInfo);
requestHeader.setInvisibleTime(invisibleTime);
requestHeader.setBname(brokerName);
//here the broker should be polished
this.mQClientFactory.getMQClientAPIImpl().changeInvisibleTimeAsync(brokerName, findBrokerResult.getBrokerAddr(), requestHeader, ASYNC_TIMEOUT, callback);
return;
@@ -379,6 +379,7 @@ public class PullAPIWrapper {
requestHeader.setExpType(expressionType);
requestHeader.setExp(expression);
requestHeader.setOrder(order);
requestHeader.setBname(mq.getBrokerName());
//give 1000 ms for server response
if (poll) {
requestHeader.setPollTime(timeout);
@@ -375,6 +375,7 @@ public class DefaultMQProducerImpl implements MQProducerInner {
thisHeader.setProducerGroup(producerGroup);
thisHeader.setTranStateTableOffset(checkRequestHeader.getTranStateTableOffset());
thisHeader.setFromTransactionCheck(true);
thisHeader.setBname(checkRequestHeader.getBname());
String uniqueKey = message.getProperties().get(MessageConst.PROPERTY_UNIQ_CLIENT_MESSAGE_ID_KEYIDX);
if (uniqueKey == null) {
@@ -835,6 +836,7 @@ public class DefaultMQProducerImpl implements MQProducerInner {
requestHeader.setReconsumeTimes(0);
requestHeader.setUnitMode(this.isUnitMode());
requestHeader.setBatch(msg instanceof MessageBatch);
requestHeader.setBname(brokerName);
if (requestHeader.getTopic().startsWith(MixAll.RETRY_GROUP_TOPIC_PREFIX)) {
String reconsumeTimes = MessageAccessor.getReconsumeTime(msg);
if (reconsumeTimes != null) {
@@ -1365,6 +1367,7 @@ public class DefaultMQProducerImpl implements MQProducerInner {
EndTransactionRequestHeader requestHeader = new EndTransactionRequestHeader();
requestHeader.setTransactionId(transactionId);
requestHeader.setCommitLogOffset(id.getOffset());
requestHeader.setBname(destBrokerName);
switch (localTransactionState) {
case COMMIT_MESSAGE:
requestHeader.setCommitOrRollback(MessageSysFlag.TRANSACTION_COMMIT_TYPE);
@@ -17,11 +17,11 @@
package org.apache.rocketmq.common.protocol.header;
import com.google.common.base.MoreObjects;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.TopicQueueRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class AckMessageRequestHeader implements CommandCustomHeader {
public class AckMessageRequestHeader extends TopicQueueRequestHeader {
@CFNotNull
private String consumerGroup;
@CFNotNull
@@ -17,11 +17,11 @@
package org.apache.rocketmq.common.protocol.header;
import com.google.common.base.MoreObjects;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.TopicQueueRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class ChangeInvisibleTimeRequestHeader implements CommandCustomHeader {
public class ChangeInvisibleTimeRequestHeader extends TopicQueueRequestHeader {
@CFNotNull
private String consumerGroup;
@CFNotNull
@@ -21,11 +21,11 @@
package org.apache.rocketmq.common.protocol.header;
import com.google.common.base.MoreObjects;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.RpcRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class CheckTransactionStateRequestHeader implements CommandCustomHeader {
public class CheckTransactionStateRequestHeader extends RpcRequestHeader {
@CFNotNull
private Long tranStateTableOffset;
@CFNotNull
@@ -18,12 +18,12 @@
package org.apache.rocketmq.common.protocol.header;
import com.google.common.base.MoreObjects;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.RpcRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.annotation.CFNullable;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class ConsumerSendMsgBackRequestHeader implements CommandCustomHeader {
public class ConsumerSendMsgBackRequestHeader extends RpcRequestHeader {
@CFNotNull
private Long offset;
@CFNotNull
@@ -18,13 +18,13 @@
package org.apache.rocketmq.common.protocol.header;
import com.google.common.base.MoreObjects;
import org.apache.rocketmq.common.rpc.RpcRequestHeader;
import org.apache.rocketmq.common.sysflag.MessageSysFlag;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.annotation.CFNullable;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class EndTransactionRequestHeader implements CommandCustomHeader {
public class EndTransactionRequestHeader extends RpcRequestHeader {
@CFNotNull
private String producerGroup;
@CFNotNull
@@ -16,12 +16,12 @@
*/
package org.apache.rocketmq.common.protocol.header;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.TopicQueueRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class NotificationRequestHeader implements CommandCustomHeader {
public class NotificationRequestHeader extends TopicQueueRequestHeader {
@CFNotNull
private String consumerGroup;
@CFNotNull
@@ -70,14 +70,14 @@ public class NotificationRequestHeader implements CommandCustomHeader {
this.topic = topic;
}
public int getQueueId() {
public Integer getQueueId() {
if (queueId < 0) {
return -1;
}
return queueId;
}
public void setQueueId(int queueId) {
public void setQueueId(Integer queueId) {
this.queueId = queueId;
}
@@ -16,11 +16,11 @@
*/
package org.apache.rocketmq.common.protocol.header;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.TopicQueueRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class PeekMessageRequestHeader implements CommandCustomHeader {
public class PeekMessageRequestHeader extends TopicQueueRequestHeader {
@CFNotNull
private String topic;
@CFNotNull
@@ -50,11 +50,11 @@ public class PeekMessageRequestHeader implements CommandCustomHeader {
this.topic = topic;
}
public int getQueueId() {
public Integer getQueueId() {
return queueId;
}
public void setQueueId(int queueId) {
public void setQueueId(Integer queueId) {
this.queueId = queueId;
}
@@ -17,12 +17,12 @@
package org.apache.rocketmq.common.protocol.header;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.TopicQueueRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class PollingInfoRequestHeader implements CommandCustomHeader {
public class PollingInfoRequestHeader extends TopicQueueRequestHeader {
@CFNotNull
private String consumerGroup;
@CFNotNull
@@ -50,14 +50,14 @@ public class PollingInfoRequestHeader implements CommandCustomHeader {
this.topic = topic;
}
public int getQueueId() {
public Integer getQueueId() {
if (queueId < 0) {
return -1;
}
return queueId;
}
public void setQueueId(int queueId) {
public void setQueueId(Integer queueId) {
this.queueId = queueId;
}
@@ -17,11 +17,11 @@
package org.apache.rocketmq.common.protocol.header;
import com.google.common.base.MoreObjects;
import org.apache.rocketmq.remoting.CommandCustomHeader;
import org.apache.rocketmq.common.rpc.TopicQueueRequestHeader;
import org.apache.rocketmq.remoting.annotation.CFNotNull;
import org.apache.rocketmq.remoting.exception.RemotingCommandException;
public class PopMessageRequestHeader implements CommandCustomHeader {
public class PopMessageRequestHeader extends TopicQueueRequestHeader {
@CFNotNull
private String consumerGroup;
@CFNotNull
@@ -102,18 +102,18 @@ public class PopMessageRequestHeader implements CommandCustomHeader {
this.topic = topic;
}
public int getQueueId() {
public Integer getQueueId() {
if (queueId < 0) {
return -1;
}
return queueId;
}
public void setQueueId(int queueId) {
@Override
public void setQueueId(Integer queueId) {
this.queueId = queueId;
}
public int getMaxMsgNums() {
return maxMsgNums;
}
@@ -59,6 +59,8 @@ public class SendMessageRequestHeaderV2 implements CommandCustomHeader, FastCode
@CFNullable
private boolean m; //batch
@CFNullable
private String n; // brokerName
public static SendMessageRequestHeader createSendMessageRequestHeaderV1(final SendMessageRequestHeaderV2 v2) {
SendMessageRequestHeader v1 = new SendMessageRequestHeader();
@@ -75,6 +77,7 @@ public class SendMessageRequestHeaderV2 implements CommandCustomHeader, FastCode
v1.setUnitMode(v2.k);
v1.setMaxReconsumeTimes(v2.l);
v1.setBatch(v2.m);
v1.setBname(v2.n);
return v1;
}
@@ -93,6 +96,7 @@ public class SendMessageRequestHeaderV2 implements CommandCustomHeader, FastCode
v2.k = v1.isUnitMode();
v2.l = v1.getMaxReconsumeTimes();
v2.m = v1.isBatch();
v2.n = v1.getBname();
return v2;
}
@@ -115,6 +119,7 @@ public class SendMessageRequestHeaderV2 implements CommandCustomHeader, FastCode
writeIfNotNull(out, "k", k);
writeIfNotNull(out, "l", l);
writeIfNotNull(out, "m", m);
writeIfNotNull(out, "n", n);
}
@Override
@@ -184,6 +189,11 @@ public class SendMessageRequestHeaderV2 implements CommandCustomHeader, FastCode
if (str != null) {
m = Boolean.parseBoolean(str);
}
str = fields.get("n");
if (str != null) {
n = str;
}
}
public String getA() {
@@ -306,6 +316,7 @@ public class SendMessageRequestHeaderV2 implements CommandCustomHeader, FastCode
.add("k", k)
.add("l", l)
.add("m", m)
.add("n", n)
.toString();
}
}