diff --git a/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java b/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java index 7fb6dc0999..8e1c8d15af 100644 --- a/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java +++ b/client/src/main/java/org/apache/rocketmq/client/consumer/DefaultMQPushConsumer.java @@ -640,7 +640,7 @@ public class DefaultMQPushConsumer extends ClientConfig implements MQPushConsume */ @Deprecated public void setSubscription(Map subscription) { - Map subscriptionWithNamespace = new HashMap(); + Map subscriptionWithNamespace = new HashMap(subscription.size(), 1); for (Entry topicEntry : subscription.entrySet()) { subscriptionWithNamespace.put(withNamespace(topicEntry.getKey()), topicEntry.getValue()); } diff --git a/client/src/main/java/org/apache/rocketmq/client/consumer/store/LocalFileOffsetStore.java b/client/src/main/java/org/apache/rocketmq/client/consumer/store/LocalFileOffsetStore.java index d380ba058a..f949b75a81 100644 --- a/client/src/main/java/org/apache/rocketmq/client/consumer/store/LocalFileOffsetStore.java +++ b/client/src/main/java/org/apache/rocketmq/client/consumer/store/LocalFileOffsetStore.java @@ -169,7 +169,7 @@ public class LocalFileOffsetStore implements OffsetStore { @Override public Map cloneOffsetTable(String topic) { - Map cloneOffsetTable = new HashMap(); + Map cloneOffsetTable = new HashMap(this.offsetTable.size(), 1); for (Map.Entry entry : this.offsetTable.entrySet()) { MessageQueue mq = entry.getKey(); if (!UtilAll.isBlank(topic) && !topic.equals(mq.getTopic())) { diff --git a/client/src/main/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStore.java b/client/src/main/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStore.java index 15b5becfd4..409ceab954 100644 --- a/client/src/main/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStore.java +++ b/client/src/main/java/org/apache/rocketmq/client/consumer/store/RemoteBrokerOffsetStore.java @@ -174,7 +174,7 @@ public class RemoteBrokerOffsetStore implements OffsetStore { @Override public Map cloneOffsetTable(String topic) { - Map cloneOffsetTable = new HashMap(); + Map cloneOffsetTable = new HashMap(this.offsetTable.size(), 1); for (Map.Entry entry : this.offsetTable.entrySet()) { MessageQueue mq = entry.getKey(); if (!UtilAll.isBlank(topic) && !topic.equals(mq.getTopic())) { diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java index be70d9f5fa..f2d01897c8 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIImpl.java @@ -388,7 +388,7 @@ public class MQClientAPIImpl { clusterAclVersionInfo.setBrokerAddr(responseHeader.getBrokerAddr()); clusterAclVersionInfo.setAclConfigDataVersion(DataVersion.fromJson(responseHeader.getVersion(), DataVersion.class)); HashMap dataVersionMap = JSON.parseObject(responseHeader.getAllAclFileVersion(), HashMap.class); - Map allAclConfigDataVersion = new HashMap(); + Map allAclConfigDataVersion = new HashMap(dataVersionMap.size(), 1); for (Map.Entry entry : dataVersionMap.entrySet()) { allAclConfigDataVersion.put(entry.getKey(),DataVersion.fromJson(JSON.toJSONString(entry.getValue()), DataVersion.class)); } diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java index a338f7b681..31cd64c01e 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/DefaultMQPushConsumerImpl.java @@ -28,6 +28,7 @@ import java.util.Properties; import java.util.Set; import java.util.concurrent.ConcurrentMap; +import org.apache.commons.collections.CollectionUtils; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.QueryResult; import org.apache.rocketmq.client.Validators; @@ -962,8 +963,8 @@ public class DefaultMQPushConsumerImpl implements MQConsumerInner { throws RemotingException, MQBrokerException, InterruptedException, MQClientException { for (String topic : rebalanceImpl.getSubscriptionInner().keySet()) { Set mqs = rebalanceImpl.getTopicSubscribeInfoTable().get(topic); - Map offsetTable = new HashMap(); - if (mqs != null) { + if (CollectionUtils.isNotEmpty(mqs)) { + Map offsetTable = new HashMap(mqs.size(), 1); for (MessageQueue mq : mqs) { long offset = searchOffset(mq, timeStamp); offsetTable.put(mq, offset); diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java index 7677d8b685..f239d7946b 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/RebalanceImpl.java @@ -117,7 +117,7 @@ public abstract class RebalanceImpl { } private HashMap> buildProcessQueueTableByBrokerName() { - HashMap> result = new HashMap>(); + HashMap> result = new HashMap>(this.processQueueTable.size(), 1); for (MessageQueue mq : this.processQueueTable.keySet()) { Set mqs = result.get(mq.getBrokerName()); if (null == mqs) { diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java b/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java index d5b90979e0..1ba3e32f3c 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/factory/MQClientInstance.java @@ -366,7 +366,7 @@ public class MQClientInstance { * @return newOffsetTable */ public Map parseOffsetTableFromBroker(Map offsetTable, String namespace) { - HashMap newOffsetTable = new HashMap(); + HashMap newOffsetTable = new HashMap(offsetTable.size(), 1); if (StringUtils.isNotEmpty(namespace)) { for (Entry entry : offsetTable.entrySet()) { MessageQueue queue = entry.getKey(); @@ -387,7 +387,7 @@ public class MQClientInstance { try { if (this.lockNamesrv.tryLock(LOCK_TIMEOUT_MILLIS, TimeUnit.MILLISECONDS)) try { - ConcurrentHashMap> updatedTable = new ConcurrentHashMap>(); + ConcurrentHashMap> updatedTable = new ConcurrentHashMap>(this.brokerAddrTable.size(), 1); Iterator>> itBrokerTable = this.brokerAddrTable.entrySet().iterator(); while (itBrokerTable.hasNext()) { @@ -395,7 +395,7 @@ public class MQClientInstance { String brokerName = entry.getKey(); HashMap oneTable = entry.getValue(); - HashMap cloneAddrTable = new HashMap(); + HashMap cloneAddrTable = new HashMap(oneTable.size(), 1); cloneAddrTable.putAll(oneTable); Iterator> it = cloneAddrTable.entrySet().iterator(); 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 86153f5263..7652ee0e50 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 @@ -330,6 +330,7 @@ public class AsyncTraceDispatcher implements TraceDispatcher { traceExecutor.submit(asyncDataSendTask); this.clear(); + } }