diff --git a/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java b/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java index 395f5e3005..e8e7e048fa 100644 --- a/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java +++ b/store/src/main/java/org/apache/rocketmq/store/StoreStatsService.java @@ -62,13 +62,13 @@ public class StoreStatsService extends ServiceThread { private volatile long putMessageEntireTimeMax = 0; private volatile long getMessageEntireTimeMax = 0; // for putMessageEntireTimeMax - private ReentrantLock lockPut = new ReentrantLock(); + private ReentrantLock putLock = new ReentrantLock(); // for getMessageEntireTimeMax - private ReentrantLock lockGet = new ReentrantLock(); + private ReentrantLock getLock = new ReentrantLock(); private volatile long dispatchMaxBuffer = 0; - private ReentrantLock lockSampling = new ReentrantLock(); + private ReentrantLock samplingLock = new ReentrantLock(); private long lastPrintTimestamp = System.currentTimeMillis(); public StoreStatsService() { @@ -138,10 +138,10 @@ public class StoreStatsService extends ServiceThread { } if (value > this.putMessageEntireTimeMax) { - this.lockPut.lock(); + this.putLock.lock(); this.putMessageEntireTimeMax = value > this.putMessageEntireTimeMax ? value : this.putMessageEntireTimeMax; - this.lockPut.unlock(); + this.putLock.unlock(); } } @@ -151,10 +151,10 @@ public class StoreStatsService extends ServiceThread { public void setGetMessageEntireTimeMax(long value) { if (value > this.getMessageEntireTimeMax) { - this.lockGet.lock(); + this.getLock.lock(); this.getMessageEntireTimeMax = value > this.getMessageEntireTimeMax ? value : this.getMessageEntireTimeMax; - this.lockGet.unlock(); + this.getLock.unlock(); } } @@ -193,11 +193,11 @@ public class StoreStatsService extends ServiceThread { } public long getPutMessageTimesTotal() { - long rs = 0; - for (LongAdder data : putMessageTopicTimesTotal.values()) { - rs += data.longValue(); - } - return rs; + Map map = putMessageTopicTimesTotal; + return map.values() + .parallelStream() + .mapToLong(LongAdder::longValue) + .sum(); } private String getFormatRuntime() { @@ -217,11 +217,11 @@ public class StoreStatsService extends ServiceThread { } public long getPutMessageSizeTotal() { - long rs = 0; - for (LongAdder data : putMessageTopicSizeTotal.values()) { - rs += data.longValue(); - } - return rs; + Map map = putMessageTopicSizeTotal; + return map.values() + .parallelStream() + .mapToLong(LongAdder::longValue) + .sum(); } private String getPutMessageDistributeTimeStringInfo(Long total) { @@ -315,7 +315,7 @@ public class StoreStatsService extends ServiceThread { private String getPutTps(int time) { String result = ""; - this.lockSampling.lock(); + this.samplingLock.lock(); try { CallSnapshot last = this.putTimesList.getLast(); @@ -325,14 +325,14 @@ public class StoreStatsService extends ServiceThread { } } finally { - this.lockSampling.unlock(); + this.samplingLock.unlock(); } return result; } private String getGetFoundTps(int time) { String result = ""; - this.lockSampling.lock(); + this.samplingLock.lock(); try { CallSnapshot last = this.getTimesFoundList.getLast(); @@ -342,7 +342,7 @@ public class StoreStatsService extends ServiceThread { result += CallSnapshot.getTPS(lastBefore, last); } } finally { - this.lockSampling.unlock(); + this.samplingLock.unlock(); } return result; @@ -350,7 +350,7 @@ public class StoreStatsService extends ServiceThread { private String getGetMissTps(int time) { String result = ""; - this.lockSampling.lock(); + this.samplingLock.lock(); try { CallSnapshot last = this.getTimesMissList.getLast(); @@ -361,14 +361,14 @@ public class StoreStatsService extends ServiceThread { } } finally { - this.lockSampling.unlock(); + this.samplingLock.unlock(); } return result; } private String getGetTotalTps(int time) { - this.lockSampling.lock(); + this.samplingLock.lock(); double found = 0; double miss = 0; try { @@ -392,7 +392,7 @@ public class StoreStatsService extends ServiceThread { } } finally { - this.lockSampling.unlock(); + this.samplingLock.unlock(); } return Double.toString(found + miss); @@ -400,7 +400,7 @@ public class StoreStatsService extends ServiceThread { private String getGetTransferedTps(int time) { String result = ""; - this.lockSampling.lock(); + this.samplingLock.lock(); try { CallSnapshot last = this.transferedMsgCountList.getLast(); @@ -411,7 +411,7 @@ public class StoreStatsService extends ServiceThread { } } finally { - this.lockSampling.unlock(); + this.samplingLock.unlock(); } return result; @@ -469,7 +469,7 @@ public class StoreStatsService extends ServiceThread { } private void sampling() { - this.lockSampling.lock(); + this.samplingLock.lock(); try { this.putTimesList.add(new CallSnapshot(System.currentTimeMillis(), getPutMessageTimesTotal())); if (this.putTimesList.size() > (MAX_RECORDS_OF_SAMPLING + 1)) { @@ -495,7 +495,7 @@ public class StoreStatsService extends ServiceThread { } } finally { - this.lockSampling.unlock(); + this.samplingLock.unlock(); } }