{"record":{"id":"b2e745d9cec0d13f","repo":"apache/seatunnel","slug":"rocketmq-consumer-can-not-pull-data-split-sta","errorCode":null,"errorMessage":"Rocketmq consumer can not pull data, split {}, start offset {}, end offset {}","messagePattern":"Rocketmq consumer can not pull data, split (.+?), start offset (.+?), end offset (.+?)","errorType":"console","errorClass":null,"httpStatus":null,"severity":"warning","filePath":"seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/source/RocketMqSourceReader.java","lineNumber":134,"sourceCode":"                sourceSplit -> {\n                    CompletableFuture<Void> completableFuture = new CompletableFuture<>();\n                    try {\n                        RocketMqConsumerThread rocketMqConsumerThread =\n                                consumerThreads.get(sourceSplit.getMessageQueue());\n                        rocketMqConsumerThread\n                                .getTasks()\n                                .put(\n                                        consumer -> {\n                                            try {\n                                                rocketMqConsumerThread.assign(sourceSplit);\n                                                MessageQueue assignedMessageQueue =\n                                                        sourceSplit.getMessageQueue();\n                                                List<MessageExt> records =\n                                                        consumer.poll(\n                                                                metadata.getBaseConfig()\n                                                                        .getPollTimeoutMillis());\n                                                if (records.isEmpty()) {\n                                                    log.warn(\n                                                            \"Rocketmq consumer can not pull data, split {}, start offset {}, end offset {}\",\n                                                            sourceSplit.getMessageQueue(),\n                                                            sourceSplit.getStartOffset(),\n                                                            sourceSplit.getEndOffset());\n                                                }\n                                                List<MessageExt> messages =\n                                                        records.stream()\n                                                                .filter(\n                                                                        record ->\n                                                                                isQueueMatch(\n                                                                                        assignedMessageQueue,\n                                                                                        record))\n                                                                .collect(Collectors.toList());\n                                                long lastOffset = -1;\n                                                for (MessageExt record : messages) {\n                                                    TopicTableConfig topicConfig =\n                                                            topicConfigs.get(record.getTopic());\n                                                    if (topicConfig == null) {","sourceCodeStart":116,"sourceCodeEnd":152,"githubUrl":"https://github.com/apache/seatunnel/blob/cf67b549a7a6c35fa0beb12d83c62892427ea919/seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/source/RocketMqSourceReader.java#L116-L152","documentation":"RocketMqSourceReader.pollNext() polls the RocketMQ consumer with the configured poll timeout; if the returned batch is empty it logs this warning with the split's message queue and start/end offsets. This is a diagnostic warning, not a thrown error — the reader continues, typically retrying the poll within the polling loop.","triggerScenarios":"consumer.poll(pollTimeoutMillis) returns no records for a split whose MessageQueue has no new messages within the poll timeout, or when the start offset already equals the end offset / consumer cannot consume the queue range (e.g. offset out of range or queue reassigned).","commonSituations":"Upstream producers idle so no new messages exist; poll-timeout-ms too small for broker latency; consumer group subscription/queue assignment changed (rebalance); reading a queue range whose messages were already consumed or expired (RocketMQ message retention) so offsets are past the queue's max offset; permission/ACL issues silently returning empty results.","solutions":["Verify producers are actually writing to this topic/queue and check the queue's min/max offsets with mqadmin queryMsgById or consumerProgress.","Increase the poll timeout (RocketMQ base config pollTimeoutMillis) if the broker is slow or remote.","Check consumer group offset progress (mqadmin consumerProgress) — if startOffset >= maxOffset the range is exhausted; the split end condition should handle it.","Confirm the consumer group has ACL permission to read the topic and that the MessageQueue assignment is current (no stuck rebalance).","If messages expired due to retention, restore/reset offsets or re-ingest; RocketMQ will not return deleted data."],"exampleFix":"// before\nList<MessageExt> records = consumer.poll(metadata.getBaseConfig().getPollTimeoutMillis());\nif (records.isEmpty()) {\n    log.warn(\"Rocketmq consumer can not pull data, split {}, ...\", ...);\n}\n// after (config-side)\n# config: increase poll timeout and verify topic activity\npoll.timeout.ms = 10000\n# or in code, back off before repoll\nif (records.isEmpty()) {\n    Thread.sleep(1000); // simple backoff before next poll\n}","handlingStrategy":"validation","validationCode":"// before polling, check the queue has data in range\nlong maxOffset = consumer.maxOffset(messageQueue);\nif (sourceSplit.getStartOffset() >= maxOffset) {\n    // nothing to consume; split is exhausted\n}\n","typeGuard":null,"tryCatchPattern":null,"preventionTips":["Keep producers running or expect empty polls during idle periods","Set poll.timeout.ms generously relative to broker latency","Monitor consumer group lag (mqadmin consumerProgress)","Verify topic ACLs for the consumer group before deployment"],"tags":["rocketmq","kafka-like-source","empty-result","offset"],"backgroundTag":"empty-result-set","analyzedSha":"cf67b549a7a6c35fa0beb12d83c62892427ea919","analyzedAt":"2026-09-10T21:44:55.265Z","contentChangedAt":"2026-09-10T21:44:55.265Z","schemaVersion":2},"datasetVersion":"2026-09-23T08:17:48.524Z"}