mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-19 10:54:54 +08:00
* Add TimerMessageStore runtime info in ClusterListSubCommand * Modify AdminBrokerProcessorTest to pass the UTs
This commit is contained in:
@@ -2090,6 +2090,19 @@ public class AdminBrokerProcessor implements NettyRequestProcessor {
|
||||
runtimeInfo.put("earliestMessageTimeStamp", String.valueOf(this.brokerController.getMessageStore().getEarliestMessageTime()));
|
||||
runtimeInfo.put("startAcceptSendRequestTimeStamp", String.valueOf(this.brokerController.getBrokerConfig().getStartAcceptSendRequestTimeStamp()));
|
||||
|
||||
if (this.brokerController.getMessageStoreConfig().isTimerWheelEnable()) {
|
||||
runtimeInfo.put("timerReadBehind", String.valueOf(this.brokerController.getMessageStore().getTimerMessageStore().getReadBehind()));
|
||||
runtimeInfo.put("timerOffsetBehind", String.valueOf(this.brokerController.getMessageStore().getTimerMessageStore().getOffsetBehind()));
|
||||
runtimeInfo.put("timerCongestNum", String.valueOf(this.brokerController.getMessageStore().getTimerMessageStore().getALlCongestNum()));
|
||||
runtimeInfo.put("timerEnqueueTps", String.valueOf(this.brokerController.getMessageStore().getTimerMessageStore().getEnqueueTps()));
|
||||
runtimeInfo.put("timerDequeueTps", String.valueOf(this.brokerController.getMessageStore().getTimerMessageStore().getDequeueTps()));
|
||||
} else {
|
||||
runtimeInfo.put("timerReadBehind", "0");
|
||||
runtimeInfo.put("timerOffsetBehind", "0");
|
||||
runtimeInfo.put("timerCongestNum", "0");
|
||||
runtimeInfo.put("timerEnqueueTps", "0.0");
|
||||
runtimeInfo.put("timerDequeueTps", "0.0");
|
||||
}
|
||||
MessageStore messageStore = this.brokerController.getMessageStore();
|
||||
runtimeInfo.put("remainTransientStoreBufferNumbs", String.valueOf(messageStore.remainTransientStoreBufferNumbs()));
|
||||
if (this.brokerController.getMessageStoreConfig().isTransientStorePoolEnable()) {
|
||||
|
||||
+1
-1
@@ -157,6 +157,7 @@ public class AdminBrokerProcessorTest {
|
||||
topic = "FooBar" + System.nanoTime();
|
||||
|
||||
brokerController.getTopicConfigManager().getTopicConfigTable().put(topic, new TopicConfig(topic));
|
||||
brokerController.getMessageStoreConfig().setTimerWheelEnable(false);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -514,7 +515,6 @@ public class AdminBrokerProcessorTest {
|
||||
assertThat(response.getCode()).isEqualTo(ResponseCode.SUCCESS);
|
||||
}
|
||||
|
||||
|
||||
@Test
|
||||
public void testGetTopicConfig() throws Exception {
|
||||
String topic = "foobar";
|
||||
|
||||
@@ -64,6 +64,7 @@ import java.util.concurrent.LinkedBlockingDeque;
|
||||
import java.util.concurrent.ScheduledExecutorService;
|
||||
import java.util.concurrent.TimeUnit;
|
||||
import java.util.concurrent.atomic.AtomicInteger;
|
||||
import org.apache.rocketmq.store.util.PerfCounter;
|
||||
|
||||
public class TimerMessageStore {
|
||||
public static final String TIMER_TOPIC = TopicValidator.SYSTEM_TOPIC_PREFIX + "wheel_timer";
|
||||
@@ -85,6 +86,7 @@ public class TimerMessageStore {
|
||||
public boolean debug = false;
|
||||
|
||||
private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.STORE_LOGGER_NAME);
|
||||
private final PerfCounter.Ticks perfs = new PerfCounter.Ticks(log);
|
||||
private final BlockingQueue<TimerRequest> enqueuePutQueue;
|
||||
private final BlockingQueue<List<TimerRequest>> dequeueGetQueue;
|
||||
private final BlockingQueue<TimerRequest> dequeuePutQueue;
|
||||
@@ -585,12 +587,15 @@ public class TimerMessageStore {
|
||||
try {
|
||||
int i = 0;
|
||||
for (; i < bufferCQ.getSize(); i += ConsumeQueue.CQ_STORE_UNIT_SIZE) {
|
||||
perfs.startTick("enqueue_get");
|
||||
try {
|
||||
long offsetPy = bufferCQ.getByteBuffer().getLong();
|
||||
int sizePy = bufferCQ.getByteBuffer().getInt();
|
||||
bufferCQ.getByteBuffer().getLong(); //tags code
|
||||
MessageExt msgExt = getMessageByCommitOffset(offsetPy, sizePy);
|
||||
if (msgExt != null) {
|
||||
if (null == msgExt) {
|
||||
perfs.getCounter("enqueue_get_miss");
|
||||
} else {
|
||||
lastEnqueueButExpiredTime = System.currentTimeMillis();
|
||||
lastEnqueueButExpiredStoreTime = msgExt.getStoreTimestamp();
|
||||
long delayedTime = Long.parseLong(msgExt.getProperty(TIMER_OUT_MS));
|
||||
@@ -614,6 +619,8 @@ public class TimerMessageStore {
|
||||
holdMomentForUnknownError();
|
||||
throw e;
|
||||
}
|
||||
} finally {
|
||||
perfs.endTick("enqueue_get");
|
||||
}
|
||||
//if broker role changes, ignore last enqueue
|
||||
if (!isRunningEnqueue()) {
|
||||
@@ -700,8 +707,10 @@ public class TimerMessageStore {
|
||||
try {
|
||||
//read the msg one by one
|
||||
while (currOffsetPy != -1) {
|
||||
if (!isRunning())
|
||||
if (!isRunning()) {
|
||||
break;
|
||||
}
|
||||
perfs.startTick("warm_dequeue");
|
||||
if (null == timeSbr || timeSbr.getStartOffset() > currOffsetPy) {
|
||||
timeSbr = timerLog.getWholeBuffer(currOffsetPy);
|
||||
if (null != timeSbr)
|
||||
@@ -735,6 +744,7 @@ public class TimerMessageStore {
|
||||
log.error("Unexpected error in warm", e);
|
||||
} finally {
|
||||
currOffsetPy = prevPos;
|
||||
perfs.endTick("warm_dequeue");
|
||||
}
|
||||
}
|
||||
for (SelectMappedBufferResult sbr : sbrs) {
|
||||
@@ -817,6 +827,7 @@ public class TimerMessageStore {
|
||||
SelectMappedBufferResult timeSbr = null;
|
||||
//read the timer log one by one
|
||||
while (currOffsetPy != -1) {
|
||||
perfs.startTick("dequeue_read_timerlog");
|
||||
if (null == timeSbr || timeSbr.getStartOffset() > currOffsetPy) {
|
||||
timeSbr = timerLog.getWholeBuffer(currOffsetPy);
|
||||
if (null != timeSbr)
|
||||
@@ -846,6 +857,7 @@ public class TimerMessageStore {
|
||||
log.error("Error in dequeue_read_timerlog", e);
|
||||
} finally {
|
||||
currOffsetPy = prevPos;
|
||||
perfs.endTick("dequeue_read_timerlog");
|
||||
}
|
||||
}
|
||||
if (deleteMsgStack.size() == 0 && normalMsgStack.size() == 0) {
|
||||
@@ -1223,12 +1235,14 @@ public class TimerMessageStore {
|
||||
for (TimerRequest req : trs) {
|
||||
req.setLatch(latch);
|
||||
try {
|
||||
perfs.startTick("enqueue_put");
|
||||
if (isMaster() && req.getDelayTime() < currWriteTimeMs) {
|
||||
dequeuePutQueue.put(req);
|
||||
} else {
|
||||
boolean doEnqueueRes = doEnqueue(req.getOffsetPy(), req.getSizePy(), req.getDelayTime(), req.getMsg());
|
||||
req.idempotentRelease(doEnqueueRes || storeConfig.isTimerSkipUnknownError());
|
||||
}
|
||||
perfs.endTick("enqueue_put");
|
||||
} catch (Throwable t) {
|
||||
log.error("Unknown error", t);
|
||||
if (storeConfig.isTimerSkipUnknownError()) {
|
||||
@@ -1329,6 +1343,7 @@ public class TimerMessageStore {
|
||||
break;
|
||||
}
|
||||
try {
|
||||
perfs.startTick("dequeue_put");
|
||||
addMetric(tr.getMsg(), -1);
|
||||
MessageExtBrokerInner msg = convert(tr.getMsg(), tr.getEnqueueTime(), needRoll(tr.getMagic()));
|
||||
doRes = PUT_NEED_RETRY != doPut(msg, needRoll(tr.getMagic()));
|
||||
@@ -1341,6 +1356,7 @@ public class TimerMessageStore {
|
||||
doRes = PUT_NEED_RETRY != doPut(msg, needRoll(tr.getMagic()));
|
||||
Thread.sleep(500 * precisionMs / 1000);
|
||||
}
|
||||
perfs.endTick("dequeue_put");
|
||||
} catch (Throwable t) {
|
||||
log.info("Unknown error", t);
|
||||
if (storeConfig.isTimerSkipUnknownError()) {
|
||||
@@ -1386,6 +1402,7 @@ public class TimerMessageStore {
|
||||
TimerRequest tr = trs.get(i);
|
||||
boolean doRes = false;
|
||||
try {
|
||||
long start = System.currentTimeMillis();
|
||||
MessageExt msgExt = getMessageByCommitOffset(tr.getOffsetPy(), tr.getSizePy());
|
||||
if (null != msgExt) {
|
||||
if (needDelete(tr.getMagic()) && !needRoll(tr.getMagic())) {
|
||||
@@ -1402,6 +1419,7 @@ public class TimerMessageStore {
|
||||
if (null != uniqkey && tr.getDeleteList() != null && tr.getDeleteList().size() > 0 && tr.getDeleteList().contains(uniqkey)) {
|
||||
doRes = true;
|
||||
tr.idempotentRelease();
|
||||
perfs.getCounter("dequeue_delete").flow(1);
|
||||
} else {
|
||||
tr.setMsg(msgExt);
|
||||
while (!isStopped() && !doRes) {
|
||||
@@ -1409,10 +1427,12 @@ public class TimerMessageStore {
|
||||
}
|
||||
}
|
||||
}
|
||||
perfs.getCounter("dequeue_get_msg").flow(System.currentTimeMillis() - start);
|
||||
} else {
|
||||
//the tr will never be processed afterwards, so idempotentRelease it
|
||||
tr.idempotentRelease();
|
||||
doRes = true;
|
||||
perfs.getCounter("dequeue_get_msg_miss").flow(System.currentTimeMillis() - start);
|
||||
}
|
||||
} catch (Throwable e) {
|
||||
log.error("Unknown exception", e);
|
||||
@@ -1555,6 +1575,14 @@ public class TimerMessageStore {
|
||||
return maxOffsetInQueue - tmpQueueOffset;
|
||||
}
|
||||
|
||||
public float getEnqueueTps() {
|
||||
return perfs.getCounter("enqueue_put").getLastTps();
|
||||
}
|
||||
|
||||
public float getDequeueTps() {
|
||||
return perfs.getCounter("dequeue_put").getLastTps();
|
||||
}
|
||||
|
||||
public void prepareTimerCheckPoint() {
|
||||
timerCheckpoint.setLastTimerLogFlushPos(timerLog.getMappedFileQueue().getFlushedWhere());
|
||||
if (isMaster()) {
|
||||
|
||||
+20
-4
@@ -178,9 +178,9 @@ public class ClusterListSubCommand implements SubCommand {
|
||||
}
|
||||
|
||||
private void printClusterBaseInfo(final Set<String> clusterNames,
|
||||
final DefaultMQAdminExt defaultMQAdminExt,
|
||||
final ClusterInfo clusterInfo) {
|
||||
System.out.printf("%-16s %-22s %-4s %-22s %-16s %19s %19s %10s %5s %6s %10s%n",
|
||||
final DefaultMQAdminExt defaultMQAdminExt,
|
||||
final ClusterInfo clusterInfo) {
|
||||
System.out.printf("%-22s %-22s %-4s %-22s %-16s %16s %16s %-22s %-11s %-12s %-8s %-10s%n",
|
||||
"#Cluster Name",
|
||||
"#Broker Name",
|
||||
"#BID",
|
||||
@@ -188,6 +188,7 @@ public class ClusterListSubCommand implements SubCommand {
|
||||
"#Version",
|
||||
"#InTPS(LOAD)",
|
||||
"#OutTPS(LOAD)",
|
||||
"#Timer(Progress)",
|
||||
"#PCWait(ms)",
|
||||
"#Hour",
|
||||
"#SPACE",
|
||||
@@ -217,6 +218,11 @@ public class ClusterListSubCommand implements SubCommand {
|
||||
String pageCacheLockTimeMills = "";
|
||||
String earliestMessageTimeStamp = "";
|
||||
String commitLogDiskRatio = "";
|
||||
long timerReadBehind = 0;
|
||||
long timerOffsetBehind = 0;
|
||||
long timerCongestNum = 0;
|
||||
float timerEnqueueTps = 0.0f;
|
||||
float timerDequeueTps = 0.0f;
|
||||
boolean isBrokerActive = false;
|
||||
try {
|
||||
KVTable kvTable = defaultMQAdminExt.fetchBrokerRuntimeStats(next1.getValue());
|
||||
@@ -235,6 +241,15 @@ public class ClusterListSubCommand implements SubCommand {
|
||||
earliestMessageTimeStamp = kvTable.getTable().get("earliestMessageTimeStamp");
|
||||
commitLogDiskRatio = kvTable.getTable().get("commitLogDiskRatio");
|
||||
|
||||
try {
|
||||
timerReadBehind = Long.parseLong(kvTable.getTable().get("timerReadBehind"));
|
||||
timerOffsetBehind = Long.parseLong(kvTable.getTable().get("timerOffsetBehind"));
|
||||
timerCongestNum = Long.parseLong(kvTable.getTable().get("timerCongestNum"));
|
||||
timerEnqueueTps = Float.parseFloat(kvTable.getTable().get("timerEnqueueTps"));
|
||||
timerDequeueTps = Float.parseFloat(kvTable.getTable().get("timerDequeueTps"));
|
||||
} catch (Throwable ignored) {
|
||||
}
|
||||
|
||||
version = kvTable.getTable().get("brokerVersionDesc");
|
||||
{
|
||||
String[] tpss = putTps.split(" ");
|
||||
@@ -265,7 +280,7 @@ public class ClusterListSubCommand implements SubCommand {
|
||||
space = Double.parseDouble(commitLogDiskRatio);
|
||||
}
|
||||
|
||||
System.out.printf("%-16s %-22s %-4s %-22s %-16s %19s %19s %10s %5s %6s %10s%n",
|
||||
System.out.printf("%-22s %-22s %-4s %-22s %-16s %16s %16s %-22s %11s %-12s %-8s %10s%n",
|
||||
clusterName,
|
||||
brokerName,
|
||||
next1.getKey(),
|
||||
@@ -273,6 +288,7 @@ public class ClusterListSubCommand implements SubCommand {
|
||||
version,
|
||||
String.format("%9.2f(%s,%sms)", in, sendThreadPoolQueueSize, sendThreadPoolQueueHeadWaitTimeMills),
|
||||
String.format("%9.2f(%s,%sms)", out, pullThreadPoolQueueSize, pullThreadPoolQueueHeadWaitTimeMills),
|
||||
String.format("%d-%d(%.1fw, %.1f, %.1f)", timerReadBehind, timerOffsetBehind, timerCongestNum / 10000.0f, timerEnqueueTps, timerDequeueTps),
|
||||
pageCacheLockTimeMills,
|
||||
String.format("%2.2f", hour),
|
||||
String.format("%.4f", space),
|
||||
|
||||
Reference in New Issue
Block a user