alibaba/canal · error · CanalException
Start RocketMQ producer error
Error message
Start RocketMQ producer error
What it means
Thrown by CanalRocketMQProducer.init when defaultMQProducer.start() raises an MQClientException. By this point the DefaultMQProducer is configured (nameserver, namespace, retry, vip channel); start() fails because the producer cannot register with a nameserver/broker — unreachable nameserver, invalid nameserver address, namespace/access-channel mismatch, or a client-side config error.
Source
Thrown at connector/rocketmq-connector/src/main/java/com/alibaba/otter/canal/connector/rocketmq/producer/CanalRocketMQProducer.java:92
defaultMQProducer = new DefaultMQProducer(rocketMQProperties.getProducerGroup(),
rpcHook,
rocketMQProperties.isEnableMessageTrace(),
rocketMQProperties.getCustomizedTraceTopic());
if (CLOUD_ACCESS_CHANNEL.equals(rocketMQProperties.getAccessChannel())) {
defaultMQProducer.setAccessChannel(AccessChannel.CLOUD);
}
if (!StringUtils.isEmpty(rocketMQProperties.getNamespace())) {
defaultMQProducer.setNamespace(rocketMQProperties.getNamespace());
}
defaultMQProducer.setNamesrvAddr(rocketMQProperties.getNamesrvAddr());
defaultMQProducer.setRetryTimesWhenSendFailed(rocketMQProperties.getRetryTimesWhenSendFailed());
defaultMQProducer.setVipChannelEnabled(rocketMQProperties.isVipChannelEnabled());
logger.info("##Start RocketMQ producer##");
try {
defaultMQProducer.start();
} catch (MQClientException ex) {
throw new CanalException("Start RocketMQ producer error", ex);
}
int parallelPartitionSendThreadSize = mqProperties.getParallelSendThreadSize();
sendPartitionExecutor = new ThreadPoolExecutor(parallelPartitionSendThreadSize,
parallelPartitionSendThreadSize,
0,
TimeUnit.SECONDS,
new ArrayBlockingQueue<>(parallelPartitionSendThreadSize * 2),
new NamedThreadFactory("MQ-Parallel-Sender-Partition"),
new ThreadPoolExecutor.CallerRunsPolicy());
}
private void loadRocketMQProperties(Properties properties) {
RocketMQProducerConfig rocketMQProperties = (RocketMQProducerConfig) this.mqProperties;
// 兼容下<=1.1.4的mq配置
doMoreCompatibleConvert("canal.mq.servers", "rocketmq.namesrv.addr", properties);
doMoreCompatibleConvert("canal.mq.producerGroup", "rocketmq.producer.group", properties);
doMoreCompatibleConvert("canal.mq.namespace", "rocketmq.namespace", properties);View on GitHub (pinned to 87be50e876)
Solutions
- Verify rocketmq.namesrv.addr resolves and the nameserver port (9876) is reachable from the producer host.
- Confirm the producer group name is unique and not used by another active producer instance.
- If using AccessChannel.CLOUD (Aliyun), ensure namespace and cloud credentials are correctly set; otherwise leave access channel as LOCAL.
- Retry start() after a short delay for transient nameserver/broker unavailability.
- Inspect MQClientException.getMessage() — it typically names the exact failing step (e.g. 'No route info of this topic', 'name server error').
Example fix
// before
try {
defaultMQProducer.start();
} catch (MQClientException ex) {
throw new CanalException("Start RocketMQ producer error", ex);
}
// after — retry transient start failures, surface nameserver in error
int attempts = 0;
while (attempts++ < 3) {
try {
defaultMQProducer.start();
break;
} catch (MQClientException ex) {
if (attempts >= 3) throw new CanalException(
"Start RocketMQ producer error [namesrv="
+ rocketMQProperties.getNamesrvAddr() + "]: " + ex.getMessage(), ex);
}
} Defensive patterns
Strategy: retry
Validate before calling
if (StringUtils.isEmpty(rocketMQProperties.getNamesrvAddr()))
throw new IllegalStateException("rocketmq namesrv addr missing"); Try / catch
int attempts = 0;
while (attempts++ < 3) {
try {
defaultMQProducer.start();
break;
} catch (MQClientException ex) {
if (attempts >= 3)
throw new CanalException("Start RocketMQ producer error", ex);
}
} Prevention
- Verify rocketmq.namesrv.addr and the nameserver port (9876) are reachable.
- Use a unique producer group name per producer instance.
- For AccessChannel.CLOUD, ensure cloud namespace and credentials are set.
When it happens
Trigger: defaultMQProducer.start() at line 90 throws MQClientException. Causes: rocketmq.namesrv.addr unreachable or malformed; nameserver returned no broker routing; AccessChannel.CLOUD set without valid cloud credentials/namespace; duplicate producer group already in use; vipChannelEnabled mismatch with the broker.
Common situations: Wrong/unreachable rocketmq.namesrv.addr; firewall blocking the nameserver or broker ports; Aliyun cloud access channel selected without credentials; producer group name collision with another running producer; broker cluster down during producer init.
Related errors
- Start RabbitMQ producer error
- Subscript pulsar consumer error
- Start RabbitMQ producer error
- Pulsar Consumer subscriptName required
- Receive pulsar batch message error
AI-assisted analysis of alibaba/canal@87be50e876 (2026-08-14).
Data as JSON: /api/errors/28ed7b6b6862a68f.
Report an issue: GitHub.