apache/seatunnel · warning

Rocketmq consumer can not pull data, split

Error message

Rocketmq consumer can not pull data, split {}, start offset {}, end offset {}

What it means

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.

Solutions

  1. Verify producers are actually writing to this topic/queue and check the queue's min/max offsets with mqadmin queryMsgById or consumerProgress.
  2. Increase the poll timeout (RocketMQ base config pollTimeoutMillis) if the broker is slow or remote.
  3. Check consumer group offset progress (mqadmin consumerProgress) — if startOffset >= maxOffset the range is exhausted; the split end condition should handle it.
  4. Confirm the consumer group has ACL permission to read the topic and that the MessageQueue assignment is current (no stuck rebalance).
  5. If messages expired due to retention, restore/reset offsets or re-ingest; RocketMQ will not return deleted data.

Example fix

// before
List<MessageExt> records = consumer.poll(metadata.getBaseConfig().getPollTimeoutMillis());
if (records.isEmpty()) {
    log.warn("Rocketmq consumer can not pull data, split {}, ...", ...);
}
// after (config-side)
# config: increase poll timeout and verify topic activity
poll.timeout.ms = 10000
# or in code, back off before repoll
if (records.isEmpty()) {
    Thread.sleep(1000); // simple backoff before next poll
}
Defensive patterns

Strategy: validation

Validate before calling

// before polling, check the queue has data in range
long maxOffset = consumer.maxOffset(messageQueue);
if (sourceSplit.getStartOffset() >= maxOffset) {
    // nothing to consume; split is exhausted
}

Prevention

When it happens

Trigger: 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).

Common situations: 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.

Understand the failure class

Background: EmptyResultError / "no results found": when an API or scraper succeeds but returns zero rows — this error's family across 9 libraries.

Related errors


AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10). Data as JSON: /api/errors/b2e745d9cec0d13f. Report an issue: GitHub.

Appendix: source

Thrown at seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/source/RocketMqSourceReader.java:134

                sourceSplit -> {
                    CompletableFuture<Void> completableFuture = new CompletableFuture<>();
                    try {
                        RocketMqConsumerThread rocketMqConsumerThread =
                                consumerThreads.get(sourceSplit.getMessageQueue());
                        rocketMqConsumerThread
                                .getTasks()
                                .put(
                                        consumer -> {
                                            try {
                                                rocketMqConsumerThread.assign(sourceSplit);
                                                MessageQueue assignedMessageQueue =
                                                        sourceSplit.getMessageQueue();
                                                List<MessageExt> records =
                                                        consumer.poll(
                                                                metadata.getBaseConfig()
                                                                        .getPollTimeoutMillis());
                                                if (records.isEmpty()) {
                                                    log.warn(
                                                            "Rocketmq consumer can not pull data, split {}, start offset {}, end offset {}",
                                                            sourceSplit.getMessageQueue(),
                                                            sourceSplit.getStartOffset(),
                                                            sourceSplit.getEndOffset());
                                                }
                                                List<MessageExt> messages =
                                                        records.stream()
                                                                .filter(
                                                                        record ->
                                                                                isQueueMatch(
                                                                                        assignedMessageQueue,
                                                                                        record))
                                                                .collect(Collectors.toList());
                                                long lastOffset = -1;
                                                for (MessageExt record : messages) {
                                                    TopicTableConfig topicConfig =
                                                            topicConfigs.get(record.getTopic());
                                                    if (topicConfig == null) {

View on GitHub (pinned to cf67b549a7)