[ISSUE #9375] Make client trace thread can be closed correctly (#9376)

This commit is contained in:
qianye
2025-04-30 10:27:18 +08:00
committed by GitHub
parent 6819b4d684
commit 28370f3a87
2 changed files with 22 additions and 24 deletions
@@ -16,6 +16,19 @@
*/
package org.apache.rocketmq.client.trace;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import org.apache.rocketmq.client.AccessChannel;
import org.apache.rocketmq.client.common.ThreadLocalIndex;
import org.apache.rocketmq.client.exception.MQClientException;
@@ -35,20 +48,6 @@ import org.apache.rocketmq.logging.org.slf4j.Logger;
import org.apache.rocketmq.logging.org.slf4j.LoggerFactory;
import org.apache.rocketmq.remoting.RPCHook;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.HashMap;
import java.util.HashSet;
import java.util.List;
import java.util.Map;
import java.util.Set;
import java.util.concurrent.ArrayBlockingQueue;
import java.util.concurrent.ThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicBoolean;
import java.util.concurrent.atomic.AtomicInteger;
import java.util.concurrent.atomic.AtomicLong;
import static org.apache.rocketmq.client.trace.TraceConstants.TRACE_INSTANCE_NAME;
public class AsyncTraceDispatcher implements TraceDispatcher {
@@ -254,9 +253,10 @@ public class AsyncTraceDispatcher implements TraceDispatcher {
} catch (Throwable e) {
log.error("flushTraceContext error", e);
}
}
if (AsyncTraceDispatcher.this.stopped) {
this.stopped = true;
if (AsyncTraceDispatcher.this.stopped) {
this.stopped = true;
}
}
}
}
@@ -16,6 +16,11 @@
*/
package org.apache.rocketmq.client.impl.consumer;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.TreeMap;
import org.apache.commons.lang3.reflect.FieldUtils;
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
import org.apache.rocketmq.client.exception.MQBrokerException;
@@ -29,12 +34,6 @@ import org.junit.Test;
import org.junit.runner.RunWith;
import org.mockito.junit.MockitoJUnitRunner;
import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
import java.util.TreeMap;
import static org.assertj.core.api.Assertions.assertThat;
import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
@@ -158,7 +157,6 @@ public class ProcessQueueTest {
ProcessQueue processQueue2 = createProcessQueue();
assertEquals(processQueue1.getMsgAccCnt(), processQueue2.getMsgAccCnt());
assertEquals(processQueue1.getTryUnlockTimes(), processQueue2.getTryUnlockTimes());
assertEquals(processQueue1.getLastLockTimestamp(), processQueue2.getLastLockTimestamp());
assertEquals(processQueue1.getLastPullTimestamp(), processQueue2.getLastPullTimestamp());
}