apache/seatunnel · critical · RocketMqConnectorException
CONSUME_DATA_FAILED
CONSUME_DATA_FAILED
Error message
No config found for topic: + record.getTopic()
What it means
RocketMqSourceReader.pollNext looks up each consumed MessageExt's topic in the configured topicConfigs map. If a message arrives from a topic that has no matching table config, RocketMqConnectorException(CONSUME_DATA_FAILED) is thrown. This indicates a mismatch between what the consumer subscribed to and the configured tables.
Solutions
- Add the missing topic to the source's topics/tables_configs so a TopicTableConfig exists for it
- Avoid pattern/wildcard subscriptions that can match unconfigured topics, or filter them
- Cancel and restart the job (not restore) after changing tables_configs so subscriptions and topicConfigs are consistent
Example fix
// before topics = "orders" // after (message from topic 'orders-rc' was surfacing) topics = "orders,orders-rc"
Defensive patterns
Strategy: try-catch
Validate before calling
// Before submitting, compare consumer subscription to configured topics
Set<String> configured = topicConfigs.keySet();
Set<String> subscribed = consumerSubscriptionTopics();
if (!configured.containsAll(subscribed)) throw new IllegalStateException("subscribed topics missing config: " + subscribed);
Try / catch
try {
reader.pollNext(...);
} catch (RocketMqConnectorException e) {
if (e.getMessage().startsWith("No config found for topic:")) {
// reconfigure topics/tables_configs and restart job from scratch
} else throw e;
} Prevention
- Avoid regex/wildcard topic subscriptions unless every match has a table config
- Restart jobs fresh after changing tables_configs; don't restore old checkpoints
- Monitor for new topics matching your subscription pattern
When it happens
Trigger: The RocketMQ consumer receives a record whose record.getTopic() is not a key in topicConfigs — e.g. wildcard/regex subscription matching more topics than configured, or topics changed in tables_configs while an old subscription is still active. Raised in pollNext during data polling.
Common situations: Using pattern subscription that picks up newly created topics not present in tables_configs, a broker admin rebinding/rename of topics, or stale consumer state after restoring from a checkpoint taken with an older config.
Understand the failure class
Background: 'Could not be found', 'does not exist', 'not found in database': the resource-not-found family when an ID, slug, key, or URI lookup comes back empty — this error's family across 20 libraries.
Related errors
- API-01
- Bedrock Mantle model refused the request
- checkpoint do not exist or have already been committed.
- Failed to serialize EdgeIngressPacket to JSON
- File discovery failed
AI-assisted analysis of apache/seatunnel@cf67b549a7 (2026-09-10).
Data as JSON: /api/errors/f577a05423468276.
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:153
"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) {
throw new RocketMqConnectorException(
RocketMqConnectorErrorCode
.CONSUME_DATA_FAILED,
"No config found for topic: "
+ record.getTopic());
}
List<String> tags = topicConfig.getTags();
boolean shouldProcess =
tags.isEmpty()
|| tags.contains(
record.getTags());
if (shouldProcess) {
topicConfig
.getDeserializationSchema()
.deserialize(
record.getBody(), output);
lastOffset = record.getQueueOffset();
}
if (Boundedness.BOUNDED.equals(View on GitHub (pinned to cf67b549a7)