mirror of
https://github.com/apache/rocketmq.git
synced 2026-08-31 02:20:15 +08:00
This commit is contained in:
+20
-5
@@ -19,12 +19,16 @@ package org.apache.rocketmq.client.impl.consumer;
|
||||
import java.io.ByteArrayOutputStream;
|
||||
import java.lang.reflect.Field;
|
||||
import java.net.InetSocketAddress;
|
||||
import java.util.Collection;
|
||||
import java.util.Collections;
|
||||
import java.util.HashSet;
|
||||
import java.util.List;
|
||||
import java.util.Set;
|
||||
import java.util.concurrent.ConcurrentMap;
|
||||
import java.util.concurrent.CountDownLatch;
|
||||
import java.util.concurrent.atomic.AtomicReference;
|
||||
|
||||
import org.apache.commons.lang3.reflect.FieldUtils;
|
||||
import org.apache.rocketmq.client.consumer.DefaultMQPushConsumer;
|
||||
import org.apache.rocketmq.client.consumer.PullCallback;
|
||||
import org.apache.rocketmq.client.consumer.PullResult;
|
||||
@@ -36,6 +40,7 @@ import org.apache.rocketmq.client.exception.MQBrokerException;
|
||||
import org.apache.rocketmq.client.impl.CommunicationMode;
|
||||
import org.apache.rocketmq.client.impl.FindBrokerResult;
|
||||
import org.apache.rocketmq.client.impl.MQClientAPIImpl;
|
||||
import org.apache.rocketmq.client.impl.MQClientManager;
|
||||
import org.apache.rocketmq.client.impl.factory.MQClientInstance;
|
||||
import org.apache.rocketmq.client.stat.ConsumerStatsManager;
|
||||
import org.apache.rocketmq.common.message.MessageClientExt;
|
||||
@@ -45,10 +50,10 @@ import org.apache.rocketmq.common.message.MessageQueue;
|
||||
import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader;
|
||||
import org.apache.rocketmq.common.stats.StatsItem;
|
||||
import org.apache.rocketmq.common.stats.StatsItemSet;
|
||||
import org.apache.rocketmq.remoting.RPCHook;
|
||||
import org.apache.rocketmq.remoting.exception.RemotingException;
|
||||
import org.junit.After;
|
||||
import org.junit.Before;
|
||||
import org.junit.Ignore;
|
||||
import org.junit.Test;
|
||||
import org.junit.runner.RunWith;
|
||||
import org.mockito.Mock;
|
||||
@@ -81,6 +86,13 @@ public class ConsumeMessageConcurrentlyServiceTest {
|
||||
|
||||
@Before
|
||||
public void init() throws Exception {
|
||||
ConcurrentMap<String, MQClientInstance> factoryTable = (ConcurrentMap<String, MQClientInstance>) FieldUtils.readDeclaredField(MQClientManager.getInstance(), "factoryTable", true);
|
||||
Collection<MQClientInstance> instances = factoryTable.values();
|
||||
for (MQClientInstance instance : instances) {
|
||||
instance.shutdown();
|
||||
}
|
||||
factoryTable.clear();
|
||||
|
||||
consumerGroup = "FooBarGroup" + System.currentTimeMillis();
|
||||
pushConsumer = new DefaultMQPushConsumer(consumerGroup);
|
||||
pushConsumer.setNamesrvAddr("127.0.0.1:9876");
|
||||
@@ -100,12 +112,15 @@ public class ConsumeMessageConcurrentlyServiceTest {
|
||||
field.setAccessible(true);
|
||||
field.set(pushConsumerImpl, rebalancePushImpl);
|
||||
pushConsumer.subscribe(topic, "*");
|
||||
pushConsumer.start();
|
||||
|
||||
mQClientFactory = spy(pushConsumerImpl.getmQClientFactory());
|
||||
// suppress updateTopicRouteInfoFromNameServer
|
||||
pushConsumer.changeInstanceNameToPID();
|
||||
mQClientFactory = MQClientManager.getInstance().getOrCreateMQClientInstance(pushConsumer, (RPCHook) FieldUtils.readDeclaredField(pushConsumerImpl, "rpcHook", true));
|
||||
mQClientFactory = spy(mQClientFactory);
|
||||
field = DefaultMQPushConsumerImpl.class.getDeclaredField("mQClientFactory");
|
||||
field.setAccessible(true);
|
||||
field.set(pushConsumerImpl, mQClientFactory);
|
||||
factoryTable.put(pushConsumer.buildMQClientId(), mQClientFactory);
|
||||
|
||||
field = MQClientInstance.class.getDeclaredField("mQClientAPIImpl");
|
||||
field.setAccessible(true);
|
||||
@@ -117,7 +132,6 @@ public class ConsumeMessageConcurrentlyServiceTest {
|
||||
field.set(pushConsumerImpl, pullAPIWrapper);
|
||||
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().getRebalanceImpl().setmQClientFactory(mQClientFactory);
|
||||
mQClientFactory.registerConsumer(consumerGroup, pushConsumerImpl);
|
||||
|
||||
when(mQClientFactory.getMQClientAPIImpl().pullMessage(anyString(), any(PullMessageRequestHeader.class),
|
||||
anyLong(), any(CommunicationMode.class), nullable(PullCallback.class)))
|
||||
@@ -140,12 +154,13 @@ public class ConsumeMessageConcurrentlyServiceTest {
|
||||
});
|
||||
|
||||
doReturn(new FindBrokerResult("127.0.0.1:10912", false)).when(mQClientFactory).findBrokerAddressInSubscribe(anyString(), anyLong(), anyBoolean());
|
||||
doReturn(false).when(mQClientFactory).updateTopicRouteInfoFromNameServer(anyString());
|
||||
Set<MessageQueue> messageQueueSet = new HashSet<MessageQueue>();
|
||||
messageQueueSet.add(createPullRequest().getMessageQueue());
|
||||
pushConsumer.getDefaultMQPushConsumerImpl().updateTopicSubscribeInfo(topic, messageQueueSet);
|
||||
pushConsumer.start();
|
||||
}
|
||||
|
||||
@Ignore
|
||||
@Test
|
||||
public void testPullMessage_ConsumeSuccess() throws InterruptedException, RemotingException, MQBrokerException, NoSuchFieldException,Exception {
|
||||
final CountDownLatch countDownLatch = new CountDownLatch(1);
|
||||
|
||||
Reference in New Issue
Block a user