diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java index b27117fbbf..cedbbdb6ca 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java @@ -223,25 +223,10 @@ public class DefaultMQProducerImpl implements MQProducerInner { this.mQClientFactory.sendHeartbeatToAllBrokerWithLock(); - this.startScheduledTask(); + RequestFutureHolder.getInstance().startScheduledTask(this); } - private void startScheduledTask() { - if (RequestFutureHolder.getInstance().getProducerNum().incrementAndGet() == 1) { - RequestFutureHolder.getInstance().getScheduledExecutorService().scheduleAtFixedRate(new Runnable() { - @Override - public void run() { - try { - RequestFutureHolder.getInstance().scanExpiredRequest(); - } catch (Throwable e) { - log.error("scan RequestFutureTable exception", e); - } - } - }, 1000 * 3, 1000, TimeUnit.MILLISECONDS); - } - } - private void checkConfig() throws MQClientException { Validators.checkGroup(this.defaultMQProducer.getProducerGroup()); @@ -269,9 +254,7 @@ public class DefaultMQProducerImpl implements MQProducerInner { if (shutdownFactory) { this.mQClientFactory.shutdown(); } - if (RequestFutureHolder.getInstance().getProducerNum().decrementAndGet() == 0) { - RequestFutureHolder.getInstance().getScheduledExecutorService().shutdown(); - } + RequestFutureHolder.getInstance().shutdown(this); log.info("the producer [{}] shutdown OK", this.defaultMQProducer.getProducerGroup()); this.serviceState = ServiceState.SHUTDOWN_ALREADY; break; diff --git a/client/src/main/java/org/apache/rocketmq/client/producer/RequestFutureHolder.java b/client/src/main/java/org/apache/rocketmq/client/producer/RequestFutureHolder.java index 24b3a906a5..8fe9abcc69 100644 --- a/client/src/main/java/org/apache/rocketmq/client/producer/RequestFutureHolder.java +++ b/client/src/main/java/org/apache/rocketmq/client/producer/RequestFutureHolder.java @@ -17,39 +17,36 @@ package org.apache.rocketmq.client.producer; +import java.util.HashSet; import java.util.Iterator; import java.util.LinkedList; import java.util.List; import java.util.Map; +import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.ThreadFactory; -import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.TimeUnit; import org.apache.rocketmq.client.common.ClientErrorCode; import org.apache.rocketmq.client.exception.RequestTimeoutException; +import org.apache.rocketmq.client.impl.producer.DefaultMQProducerImpl; import org.apache.rocketmq.client.log.ClientLogger; +import org.apache.rocketmq.common.ThreadFactoryImpl; import org.apache.rocketmq.logging.InternalLogger; public class RequestFutureHolder { private static InternalLogger log = ClientLogger.getLog(); private static final RequestFutureHolder INSTANCE = new RequestFutureHolder(); private ConcurrentHashMap requestFutureTable = new ConcurrentHashMap(); - private final AtomicInteger producerNum = new AtomicInteger(0); - private final ScheduledExecutorService scheduledExecutorService = Executors - .newSingleThreadScheduledExecutor(new ThreadFactory() { - @Override - public Thread newThread(Runnable r) { - return new Thread(r, "RequestHouseKeepingService"); - } - }); + private final Set producerSet = new HashSet<>(); + private ScheduledExecutorService scheduledExecutorService = null; public ConcurrentHashMap getRequestFutureTable() { return requestFutureTable; } - public void scanExpiredRequest() { + private void scanExpiredRequest() { final List rfList = new LinkedList(); Iterator> it = requestFutureTable.entrySet().iterator(); while (it.hasNext()) { @@ -74,16 +71,35 @@ public class RequestFutureHolder { } } - private RequestFutureHolder() { + public synchronized void startScheduledTask(DefaultMQProducerImpl producer) { + this.producerSet.add(producer); + if (null == scheduledExecutorService) { + this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor(new ThreadFactoryImpl("RequestHouseKeepingService")); + + this.scheduledExecutorService.scheduleAtFixedRate(new Runnable() { + @Override + public void run() { + try { + RequestFutureHolder.getInstance().scanExpiredRequest(); + } catch (Throwable e) { + log.error("scan RequestFutureTable exception", e); + } + } + }, 1000 * 3, 1000, TimeUnit.MILLISECONDS); + + } } - public AtomicInteger getProducerNum() { - return producerNum; + public synchronized void shutdown(DefaultMQProducerImpl producer) { + this.producerSet.remove(producer); + if (this.producerSet.size() <= 0 && null != this.scheduledExecutorService) { + ScheduledExecutorService executorService = this.scheduledExecutorService; + this.scheduledExecutorService = null; + executorService.shutdown(); + } } - public ScheduledExecutorService getScheduledExecutorService() { - return scheduledExecutorService; - } + private RequestFutureHolder() {} public static RequestFutureHolder getInstance() { return INSTANCE;