apache / apache/rocketmq-flink
The message queue is not in assigned list
- Dominant language
- Java
- Stars
- 174
- Forks
- 104
- PR merge metrics
- No merged PRs in 30d
Description
I got the error when i consume rocketmq message in flink job:
```
17:51:50.537 [rmq-pull-thread-1] ERROR org.apache.flink.connector.rocketmq.legacy.common.util.RetryUtil - RuntimeException, retry 2/5
java.lang.RuntimeException: org.apache.rocketmq.client.exception.MQClientException: The message queue is not in assigned list, may be rebalancing, message queue: MessageQueue [topic=yumiao1205, brokerName=broker-a, queueId=1]
For more information, please visit the url, https://rocketmq.apache.org/docs/bestPractice/06FAQ
at org.apache.flink.connector.rocketmq.legacy.RocketMQSourceFunction.lambda$null$1(RocketMQSourceFunction.java:348)
at org.apache.flink.connector.rocketmq.legacy.common.util.RetryUtil.call(RetryUtil.java:56)
at org.apache.flink.connector.rocketmq.legacy.RocketMQSourceFunction.lambda$run$2(RocketMQSourceFunction.java:276)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:750)
Caused by: org.apache.rocketmq.client.exception.MQClientException: The message queue is not in assigned list, may be rebalancing, message queue: MessageQueue [topic=yumiao1205, brokerName=broker-a, queueId=1]
For more information, please visit the url, https://rocketmq.apache.org/docs/bestPractice/06FAQ
at org.apache.rocketmq.client.impl.consumer.DefaultLitePullConsumerImpl.seek(DefaultLitePullConsumerImpl.java:658)
at org.apache.rocketmq.client.consumer.DefaultLitePullConsumer.seek(DefaultLitePullConsumer.java:298)
at org.apache.flink.connector.rocketmq.legacy.RocketMQSourceFunction.lambda$null$1(RocketMQSourceFunction.java:282)
... 5 common frames omitted
17:51:51.342 [rmq-pull-thread-2] ERROR org.apache.flink.connector.rocketmq.legacy.common.util.RetryUtil - RuntimeException, retry 3/5
java.lang.RuntimeException: org.apache.rocketmq.client.exception.MQClientException: The message queue is not in assigned list, may be rebalancing, message queue: MessageQueue [topic=yumiao1205, brokerName=broker-a, queueId=2]
For more information, please visit the url, https://rocketmq.apache.org/docs/bestPractice/06FAQ
at org.apache.flink.connector.rocketmq.legacy.RocketMQSourceFunction.lambda$null$1(RocketMQSourceFunction.java:348)
at org.apache.flink.connector.rocketmq.legacy.common.util.RetryUtil.call(RetryUtil.java:56)
at org.apache.flink.connector.rocketmq.legacy.RocketMQSourceFunction.lambda$run$2(RocketMQSourceFunction.java:276)
at java.util.concurrent.ThreadPoolExecutor.runWorker(ThreadPoolExecutor.java:1149)
at java.util.concurrent.ThreadPoolExecutor$Worker.run(ThreadPoolExecutor.java:624)
at java.lang.Thread.run(Thread.java:750)
Caused by: org.apache.rocketmq.client.exception.MQClientException: The message queue is not in assigned list, may be rebalancing, message queue: MessageQueue [topic=yumiao1205, brokerName=broker-a, queueId=2]
For more information, please visit the url, https://rocketmq.apache.org/docs/bestPractice/06FAQ
at org.apache.rocketmq.client.impl.consumer.DefaultLitePullConsumerImpl.seek(DefaultLitePullConsumerImpl.java:658)
at org.apache.rocketmq.client.consumer.DefaultLitePullConsumer.seek(DefaultLitePullConsumer.java:298)
at org.apache.flink.connector.rocketmq.legacy.RocketMQSourceFunction.lambda$null$1(RocketMQSourceFunction.java:282)
... 5 common frames omitted
```
My rocketmq version: 5.1.1
flink version: 1.15.0
Contributor guide
No contributing guide indexed for this repository
Research direction
Start with RocketMQSourceFunction.java at lines 276, 282, and 348, then inspect RetryUtil.java:56 and DefaultLitePullConsumerImpl.seek at line 658. Reproduce with RocketMQ 5.1.1 and Flink 1.15.0, trace the assigned queues during rebalancing, and confirm the source no longer reports this failure unexpectedly.
Written by the indexing model from the issue text.
Assessment
- Tech stack
- java
- Domain
- stream-processing
- Issue type
- Bug
- Difficulty
- 4/5
- Estimated time
- 3-5 days
- Activity status
- Stale
- Clarity
- Needs clarification
- Newbie friendliness
- 25/100