{"record":{"id":"f577a05423468276","repo":"apache/seatunnel","slug":"consume-data-failed","errorCode":"CONSUME_DATA_FAILED","errorMessage":"No config found for topic: + record.getTopic()","messagePattern":"No config found for topic: \\+ record\\.getTopic\\(\\)","errorType":"error_code","errorClass":"RocketMqConnectorException","httpStatus":null,"severity":"critical","filePath":"seatunnel-connectors-v2/connector-rocketmq/src/main/java/org/apache/seatunnel/connectors/seatunnel/rocketmq/source/RocketMqSourceReader.java","lineNumber":153,"sourceCode":"                                                            \"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) {\n                                                        throw new RocketMqConnectorException(\n                                                                RocketMqConnectorErrorCode\n                                                                        .CONSUME_DATA_FAILED,\n                                                                \"No config found for topic: \"\n                                                                        + record.getTopic());\n                                                    }\n                                                    List<String> tags = topicConfig.getTags();\n                                                    boolean shouldProcess =\n                                                            tags.isEmpty()\n                                                                    || tags.contains(\n                                                                            record.getTags());\n                                                    if (shouldProcess) {\n                                                        topicConfig\n                                                                .getDeserializationSchema()\n                                                                .deserialize(\n                                                                        record.getBody(), output);\n                                                        lastOffset = record.getQueueOffset();\n                                                    }\n                                                    if (Boundedness.BOUNDED.equals(","sourceCodeStart":135,"sourceCodeEnd":171,"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#L135-L171","documentation":"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.","triggerScenarios":"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.","commonSituations":"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.","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"],"exampleFix":"// before\ntopics = \"orders\"\n// after (message from topic 'orders-rc' was surfacing)\ntopics = \"orders,orders-rc\"","handlingStrategy":"try-catch","validationCode":"// Before submitting, compare consumer subscription to configured topics\nSet<String> configured = topicConfigs.keySet();\nSet<String> subscribed = consumerSubscriptionTopics();\nif (!configured.containsAll(subscribed)) throw new IllegalStateException(\"subscribed topics missing config: \" + subscribed);\n","typeGuard":null,"tryCatchPattern":"try {\n    reader.pollNext(...);\n} catch (RocketMqConnectorException e) {\n    if (e.getMessage().startsWith(\"No config found for topic:\")) {\n        // reconfigure topics/tables_configs and restart job from scratch\n    } else throw e;\n}","preventionTips":["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"],"tags":["rocketmq","runtime","subscription"],"backgroundTag":"resource-not-found","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"}