File tree Expand file tree Collapse file tree 1 file changed +5
-2
lines changed
test/src/test/java/org/apache/rocketmq/test/client/consumer/filter Expand file tree Collapse file tree 1 file changed +5
-2
lines changed Original file line number Diff line number Diff line change @@ -90,14 +90,17 @@ public void testFilterPullConsumer() throws Exception {
9090 DefaultMQPullConsumer consumer = new DefaultMQPullConsumer (group );
9191 consumer .setNamesrvAddr (NAMESRV_ADDR );
9292 consumer .start ();
93- Thread .sleep (3000 );
93+ Set <MessageQueue > mqs = consumer .fetchSubscribeMessageQueues (topic );
94+ consumer .getDefaultMQPullConsumerImpl ()
95+ .getRebalanceImpl ()
96+ .getmQClientFactory ()
97+ .sendHeartbeatToAllBrokerWithLock ();
9498 producer .send ("TagA" , msgSize );
9599 producer .send ("TagB" , msgSize );
96100 producer .send ("TagC" , msgSize );
97101 Assert .assertEquals ("Not all sent succeeded" , msgSize * 3 , producer .getAllUndupMsgBody ().size ());
98102
99103 List <String > receivedMessage = new ArrayList <>(2 );
100- Set <MessageQueue > mqs = consumer .fetchSubscribeMessageQueues (topic );
101104 for (MessageQueue mq : mqs ) {
102105 SINGLE_MQ :
103106 while (true ) {
You can’t perform that action at this time.
0 commit comments