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
- 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.
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
- 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
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
- checkpoint do not exist or have already been committed.
- CONSUME_DATA_FAILED
- Convert timestamp to LSN offset error
- Convert timestamp to redoLog offset error
- Failed to update offset for skipped event
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)