mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
* Fix some NPE when timerWheel disabled * Add UT
This commit is contained in:
@@ -278,19 +278,26 @@ public class PopMessageProcessor implements NettyRequestProcessor {
|
||||
|
||||
if (requestHeader.isTimeoutTooMuch()) {
|
||||
response.setCode(ResponseCode.POLLING_TIMEOUT);
|
||||
response.setRemark(String.format("the broker[%s] poping message is timeout too much",
|
||||
response.setRemark(String.format("the broker[%s] pop message is timeout too much",
|
||||
this.brokerController.getBrokerConfig().getBrokerIP1()));
|
||||
return response;
|
||||
}
|
||||
if (!PermName.isReadable(this.brokerController.getBrokerConfig().getBrokerPermission())) {
|
||||
response.setCode(ResponseCode.NO_PERMISSION);
|
||||
response.setRemark(String.format("the broker[%s] poping message is forbidden",
|
||||
response.setRemark(String.format("the broker[%s] pop message is forbidden",
|
||||
this.brokerController.getBrokerConfig().getBrokerIP1()));
|
||||
return response;
|
||||
}
|
||||
if (requestHeader.getMaxMsgNums() > 32) {
|
||||
response.setCode(ResponseCode.SYSTEM_ERROR);
|
||||
response.setRemark(String.format("the broker[%s] poping message's num is greater than 32",
|
||||
response.setRemark(String.format("the broker[%s] pop message's num is greater than 32",
|
||||
this.brokerController.getBrokerConfig().getBrokerIP1()));
|
||||
return response;
|
||||
}
|
||||
|
||||
if (!brokerController.getMessageStore().getMessageStoreConfig().isTimerWheelEnable()) {
|
||||
response.setCode(ResponseCode.SYSTEM_ERROR);
|
||||
response.setRemark(String.format("the broker[%s] pop message is forbidden because timerWheelEnable is false",
|
||||
this.brokerController.getBrokerConfig().getBrokerIP1()));
|
||||
return response;
|
||||
}
|
||||
|
||||
@@ -600,6 +600,11 @@ public class PopReviveService extends ServiceThread {
|
||||
continue;
|
||||
}
|
||||
|
||||
if (!brokerController.getMessageStore().getMessageStoreConfig().isTimerWheelEnable()) {
|
||||
POP_LOGGER.warn("skip revive topic because timerWheelEnable is false");
|
||||
continue;
|
||||
}
|
||||
|
||||
POP_LOGGER.info("start revive topic={}, reviveQueueId={}", reviveTopic, queueId);
|
||||
ConsumeReviveObj consumeReviveObj = new ConsumeReviveObj();
|
||||
consumeReviveMessage(consumeReviveObj);
|
||||
|
||||
+15
-1
@@ -91,6 +91,7 @@ public class PopMessageProcessorTest {
|
||||
|
||||
@Test
|
||||
public void testProcessRequest_TopicNotExist() throws RemotingCommandException {
|
||||
when(messageStore.getMessageStoreConfig()).thenReturn(new MessageStoreConfig());
|
||||
brokerController.getTopicConfigManager().getTopicConfigTable().remove(topic);
|
||||
final RemotingCommand request = createPopMsgCommand();
|
||||
RemotingCommand response = popMessageProcessor.processRequest(handlerContext, request);
|
||||
@@ -102,6 +103,7 @@ public class PopMessageProcessorTest {
|
||||
@Test
|
||||
public void testProcessRequest_Found() throws RemotingCommandException, InterruptedException {
|
||||
GetMessageResult getMessageResult = createGetMessageResult(1);
|
||||
when(messageStore.getMessageStoreConfig()).thenReturn(new MessageStoreConfig());
|
||||
when(messageStore.getMessageAsync(anyString(), anyString(), anyInt(), anyLong(), anyInt(), any())).thenReturn(CompletableFuture.completedFuture(getMessageResult));
|
||||
|
||||
final RemotingCommand request = createPopMsgCommand();
|
||||
@@ -115,6 +117,7 @@ public class PopMessageProcessorTest {
|
||||
public void testProcessRequest_MsgWasRemoving() throws RemotingCommandException {
|
||||
GetMessageResult getMessageResult = createGetMessageResult(1);
|
||||
getMessageResult.setStatus(GetMessageStatus.MESSAGE_WAS_REMOVING);
|
||||
when(messageStore.getMessageStoreConfig()).thenReturn(new MessageStoreConfig());
|
||||
when(messageStore.getMessageAsync(anyString(), anyString(), anyInt(), anyLong(), anyInt(), any())).thenReturn(CompletableFuture.completedFuture(getMessageResult));
|
||||
|
||||
final RemotingCommand request = createPopMsgCommand();
|
||||
@@ -128,6 +131,7 @@ public class PopMessageProcessorTest {
|
||||
public void testProcessRequest_NoMsgInQueue() throws RemotingCommandException {
|
||||
GetMessageResult getMessageResult = createGetMessageResult(0);
|
||||
getMessageResult.setStatus(GetMessageStatus.NO_MESSAGE_IN_QUEUE);
|
||||
when(messageStore.getMessageStoreConfig()).thenReturn(new MessageStoreConfig());
|
||||
when(messageStore.getMessageAsync(anyString(), anyString(), anyInt(), anyLong(), anyInt(), any())).thenReturn(CompletableFuture.completedFuture(getMessageResult));
|
||||
|
||||
final RemotingCommand request = createPopMsgCommand();
|
||||
@@ -135,6 +139,17 @@ public class PopMessageProcessorTest {
|
||||
assertThat(response).isNull();
|
||||
}
|
||||
|
||||
@Test
|
||||
public void testProcessRequest_whenTimerWheelIsFalse() throws RemotingCommandException {
|
||||
MessageStoreConfig messageStoreConfig = new MessageStoreConfig();
|
||||
messageStoreConfig.setTimerWheelEnable(false);
|
||||
when(messageStore.getMessageStoreConfig()).thenReturn(messageStoreConfig);
|
||||
final RemotingCommand request = createPopMsgCommand();
|
||||
RemotingCommand response = popMessageProcessor.processRequest(handlerContext, request);
|
||||
assertThat(response).isNotNull();
|
||||
assertThat(response.getCode()).isEqualTo(ResponseCode.SYSTEM_ERROR);
|
||||
assertThat(response.getRemark()).contains("pop message is forbidden because timerWheelEnable is false");
|
||||
}
|
||||
|
||||
private RemotingCommand createPopMsgCommand() {
|
||||
PopMessageRequestHeader requestHeader = new PopMessageRequestHeader();
|
||||
@@ -152,7 +167,6 @@ public class PopMessageProcessorTest {
|
||||
return request;
|
||||
}
|
||||
|
||||
|
||||
private GetMessageResult createGetMessageResult(int msgCnt) {
|
||||
GetMessageResult getMessageResult = new GetMessageResult();
|
||||
getMessageResult.setStatus(GetMessageStatus.FOUND);
|
||||
|
||||
+52
-50
@@ -116,57 +116,59 @@ public class DefaultStoreMetricsManager {
|
||||
measurement.record(System.currentTimeMillis() - earliestMessageTime, newAttributesBuilder().build());
|
||||
});
|
||||
|
||||
timerEnqueueLag = meter.gaugeBuilder(GAUGE_TIMER_ENQUEUE_LAG)
|
||||
.setDescription("Timer enqueue messages lag")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getEnqueueBehindMessages(), newAttributesBuilder().build());
|
||||
});
|
||||
if (messageStore.getMessageStoreConfig().isTimerWheelEnable()) {
|
||||
timerEnqueueLag = meter.gaugeBuilder(GAUGE_TIMER_ENQUEUE_LAG)
|
||||
.setDescription("Timer enqueue messages lag")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getEnqueueBehindMessages(), newAttributesBuilder().build());
|
||||
});
|
||||
|
||||
timerEnqueueLatency = meter.gaugeBuilder(GAUGE_TIMER_ENQUEUE_LATENCY)
|
||||
.setDescription("Timer enqueue latency")
|
||||
.setUnit("milliseconds")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getEnqueueBehindMillis(), newAttributesBuilder().build());
|
||||
});
|
||||
timerDequeueLag = meter.gaugeBuilder(GAUGE_TIMER_DEQUEUE_LAG)
|
||||
.setDescription("Timer dequeue messages lag")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getDequeueBehindMessages(), newAttributesBuilder().build());
|
||||
});
|
||||
timerDequeueLatency = meter.gaugeBuilder(GAUGE_TIMER_DEQUEUE_LATENCY)
|
||||
.setDescription("Timer dequeue latency")
|
||||
.setUnit("milliseconds")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getDequeueBehind(), newAttributesBuilder().build());
|
||||
});
|
||||
timingMessages = meter.gaugeBuilder(GAUGE_TIMING_MESSAGES)
|
||||
.setDescription("Current message number in timing")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
timerMessageStore.getTimerMetrics()
|
||||
.getTimingCount()
|
||||
.forEach((topic, metric) -> {
|
||||
measurement.record(
|
||||
metric.getCount().get(),
|
||||
newAttributesBuilder().put(LABEL_TOPIC, topic).build()
|
||||
);
|
||||
});
|
||||
});
|
||||
timerDequeueTotal = meter.counterBuilder(COUNTER_TIMER_DEQUEUE_TOTAL)
|
||||
.setDescription("Total number of timer dequeue")
|
||||
.build();
|
||||
timerEnqueueTotal = meter.counterBuilder(COUNTER_TIMER_ENQUEUE_TOTAL)
|
||||
.setDescription("Total number of timer enqueue")
|
||||
.build();
|
||||
timerEnqueueLatency = meter.gaugeBuilder(GAUGE_TIMER_ENQUEUE_LATENCY)
|
||||
.setDescription("Timer enqueue latency")
|
||||
.setUnit("milliseconds")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getEnqueueBehindMillis(), newAttributesBuilder().build());
|
||||
});
|
||||
timerDequeueLag = meter.gaugeBuilder(GAUGE_TIMER_DEQUEUE_LAG)
|
||||
.setDescription("Timer dequeue messages lag")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getDequeueBehindMessages(), newAttributesBuilder().build());
|
||||
});
|
||||
timerDequeueLatency = meter.gaugeBuilder(GAUGE_TIMER_DEQUEUE_LATENCY)
|
||||
.setDescription("Timer dequeue latency")
|
||||
.setUnit("milliseconds")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
measurement.record(timerMessageStore.getDequeueBehind(), newAttributesBuilder().build());
|
||||
});
|
||||
timingMessages = meter.gaugeBuilder(GAUGE_TIMING_MESSAGES)
|
||||
.setDescription("Current message number in timing")
|
||||
.ofLongs()
|
||||
.buildWithCallback(measurement -> {
|
||||
TimerMessageStore timerMessageStore = messageStore.getTimerMessageStore();
|
||||
timerMessageStore.getTimerMetrics()
|
||||
.getTimingCount()
|
||||
.forEach((topic, metric) -> {
|
||||
measurement.record(
|
||||
metric.getCount().get(),
|
||||
newAttributesBuilder().put(LABEL_TOPIC, topic).build()
|
||||
);
|
||||
});
|
||||
});
|
||||
timerDequeueTotal = meter.counterBuilder(COUNTER_TIMER_DEQUEUE_TOTAL)
|
||||
.setDescription("Total number of timer dequeue")
|
||||
.build();
|
||||
timerEnqueueTotal = meter.counterBuilder(COUNTER_TIMER_ENQUEUE_TOTAL)
|
||||
.setDescription("Total number of timer enqueue")
|
||||
.build();
|
||||
}
|
||||
}
|
||||
|
||||
public static void incTimerDequeueCount(String topic) {
|
||||
|
||||
Reference in New Issue
Block a user