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

  1. Add the missing topic to the source's topics/tables_configs so a TopicTableConfig exists for it
  2. Avoid pattern/wildcard subscriptions that can match unconfigured topics, or filter them
  3. 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

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


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)