mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-01 15:47:01 +08:00
@@ -258,8 +258,8 @@ public class PeekMessageProcessor implements NettyRequestProcessor {
|
||||
BrokerMetricsManager.throughputOutTotal.add(getMessageResult.getBufferTotalSize(), attributes);
|
||||
}
|
||||
|
||||
for (SelectMappedBufferResult mapedBuffer : getMessageTmpResult.getMessageMapedList()) {
|
||||
getMessageResult.addMessage(mapedBuffer);
|
||||
for (SelectMappedBufferResult mappedBuffer : getMessageTmpResult.getMessageMapedList()) {
|
||||
getMessageResult.addMessage(mappedBuffer);
|
||||
}
|
||||
}
|
||||
return restNum;
|
||||
|
||||
+3
-3
@@ -115,10 +115,10 @@ public class ReplyMessageProcessor extends AbstractSendMessageProcessor {
|
||||
response.addExtField(MessageConst.PROPERTY_TRACE_SWITCH, String.valueOf(this.brokerController.getBrokerConfig().isTraceOn()));
|
||||
|
||||
log.debug("receive SendReplyMessage request command, {}", request);
|
||||
final long startTimstamp = this.brokerController.getBrokerConfig().getStartAcceptSendRequestTimeStamp();
|
||||
if (this.brokerController.getMessageStore().now() < startTimstamp) {
|
||||
final long startTimestamp = this.brokerController.getBrokerConfig().getStartAcceptSendRequestTimeStamp();
|
||||
if (this.brokerController.getMessageStore().now() < startTimestamp) {
|
||||
response.setCode(ResponseCode.SYSTEM_ERROR);
|
||||
response.setRemark(String.format("broker unable to service, until %s", UtilAll.timeMillisToHumanString2(startTimstamp)));
|
||||
response.setRemark(String.format("broker unable to service, until %s", UtilAll.timeMillisToHumanString2(startTimestamp)));
|
||||
return response;
|
||||
}
|
||||
|
||||
|
||||
@@ -1619,10 +1619,10 @@ public class MQClientAPIImpl implements NameServerUpdateCallback, StartAndShutdo
|
||||
final QueryMessageRequestHeader requestHeader,
|
||||
final long timeoutMillis,
|
||||
final InvokeCallback invokeCallback,
|
||||
final Boolean isUnqiueKey
|
||||
final Boolean isUniqueKey
|
||||
) throws RemotingException, MQBrokerException, InterruptedException {
|
||||
RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.QUERY_MESSAGE, requestHeader);
|
||||
request.addExtField(MixAll.UNIQUE_MSG_QUERY_FLAG, isUnqiueKey.toString());
|
||||
request.addExtField(MixAll.UNIQUE_MSG_QUERY_FLAG, isUniqueKey.toString());
|
||||
this.remotingClient.invokeAsync(MixAll.brokerVIPChannel(this.clientConfig.isVipChannelEnabled(), addr), request, timeoutMillis,
|
||||
invokeCallback);
|
||||
}
|
||||
|
||||
@@ -26,8 +26,8 @@ public class RMQBroadCastConsumer extends RMQNormalConsumer {
|
||||
private static Logger logger = LoggerFactory.getLogger(RMQBroadCastConsumer.class);
|
||||
|
||||
public RMQBroadCastConsumer(String nsAddr, String topic, String subExpression,
|
||||
String consumerGroup, AbstractListener listner) {
|
||||
super(nsAddr, topic, subExpression, consumerGroup, listner);
|
||||
String consumerGroup, AbstractListener listener) {
|
||||
super(nsAddr, topic, subExpression, consumerGroup, listener);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
@@ -69,8 +69,8 @@ public abstract class AbstractMQConsumer implements MQConsumer {
|
||||
return listener;
|
||||
}
|
||||
|
||||
public void setListener(AbstractListener listner) {
|
||||
this.listener = listner;
|
||||
public void setListener(AbstractListener listener) {
|
||||
this.listener = listener;
|
||||
}
|
||||
|
||||
public String getNsAddr() {
|
||||
|
||||
Reference in New Issue
Block a user