From b44d3268c7a5a30d92fa8700af876d1b063df646 Mon Sep 17 00:00:00 2001 From: zhangjidi Date: Tue, 21 Jun 2022 10:16:08 +0800 Subject: [PATCH] [ISSUE #4489]Some trace messages not being sent to the broker in time before producer shutdown. --- .../client/trace/AsyncTraceDispatcher.java | 19 ++++++++++++------- 1 file changed, 12 insertions(+), 7 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java b/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java index 7652ee0e50..410ca1e33f 100644 --- a/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java +++ b/client/src/main/java/org/apache/rocketmq/client/trace/AsyncTraceDispatcher.java @@ -99,12 +99,12 @@ public class AsyncTraceDispatcher implements TraceDispatcher { this.traceTopicName = TopicValidator.RMQ_SYS_TRACE_TOPIC; } this.traceExecutor = new ThreadPoolExecutor(// - 10, // - 20, // - 1000 * 60, // - TimeUnit.MILLISECONDS, // - this.appenderQueue, // - new ThreadFactoryImpl("MQTraceSendThread_")); + 10, // + 20, // + 1000 * 60, // + TimeUnit.MILLISECONDS, // + this.appenderQueue, // + new ThreadFactoryImpl("MQTraceSendThread_")); traceProducer = getAndCreateTraceProducer(rpcHook); } @@ -188,6 +188,11 @@ public class AsyncTraceDispatcher implements TraceDispatcher { // The maximum waiting time for refresh,avoid being written all the time, resulting in failure to return. long end = System.currentTimeMillis() + 500; while (System.currentTimeMillis() <= end) { + synchronized (taskQueueByTopic) { + for (TraceDataSegment taskInfo : taskQueueByTopic.values()) { + taskInfo.sendAllData(); + } + } synchronized (traceContextQueue) { if (traceContextQueue.size() == 0 && appenderQueue.size() == 0) { break; @@ -252,7 +257,7 @@ public class AsyncTraceDispatcher implements TraceDispatcher { while (System.currentTimeMillis() < endTime) { try { TraceContext traceContext = traceContextQueue.poll( - endTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS + endTime - System.currentTimeMillis(), TimeUnit.MILLISECONDS ); if (traceContext != null && !traceContext.getTraceBeans().isEmpty()) {