mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
* typo fix * add MessageExt.getBrokerName() * add brokerName for toString * get brokerName from mq
This commit is contained in:
@@ -104,6 +104,7 @@ public class PullAPIWrapper {
|
||||
Long.toString(pullResult.getMinOffset()));
|
||||
MessageAccessor.putProperty(msg, MessageConst.PROPERTY_MAX_OFFSET,
|
||||
Long.toString(pullResult.getMaxOffset()));
|
||||
msg.setBrokerName(mq.getBrokerName());
|
||||
}
|
||||
|
||||
pullResultExt.setMsgFoundList(msgListFilterAgain);
|
||||
|
||||
@@ -27,6 +27,8 @@ import org.apache.rocketmq.common.sysflag.MessageSysFlag;
|
||||
public class MessageExt extends Message {
|
||||
private static final long serialVersionUID = 5720810158625748049L;
|
||||
|
||||
private String brokerName;
|
||||
|
||||
private int queueId;
|
||||
|
||||
private int storeSize;
|
||||
@@ -107,6 +109,14 @@ public class MessageExt extends Message {
|
||||
return socketAddress2ByteBuffer(this.storeHost, byteBuffer);
|
||||
}
|
||||
|
||||
public String getBrokerName() {
|
||||
return brokerName;
|
||||
}
|
||||
|
||||
public void setBrokerName(String brokerName) {
|
||||
this.brokerName = brokerName;
|
||||
}
|
||||
|
||||
public int getQueueId() {
|
||||
return queueId;
|
||||
}
|
||||
@@ -235,7 +245,7 @@ public class MessageExt extends Message {
|
||||
|
||||
@Override
|
||||
public String toString() {
|
||||
return "MessageExt [queueId=" + queueId + ", storeSize=" + storeSize + ", queueOffset=" + queueOffset
|
||||
return "MessageExt [brokerName=" + brokerName + ", queueId=" + queueId + ", storeSize=" + storeSize + ", queueOffset=" + queueOffset
|
||||
+ ", sysFlag=" + sysFlag + ", bornTimestamp=" + bornTimestamp + ", bornHost=" + bornHost
|
||||
+ ", storeTimestamp=" + storeTimestamp + ", storeHost=" + storeHost + ", msgId=" + msgId
|
||||
+ ", commitLogOffset=" + commitLogOffset + ", bodyCRC=" + bodyCRC + ", reconsumeTimes="
|
||||
|
||||
Reference in New Issue
Block a user