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

  1. Verify rocketmq.namesrv.addr resolves and the nameserver port (9876) is reachable from the producer host.
  2. Confirm the producer group name is unique and not used by another active producer instance.
  3. If using AccessChannel.CLOUD (Aliyun), ensure namespace and cloud credentials are correctly set; otherwise leave access channel as LOCAL.
  4. Retry start() after a short delay for transient nameserver/broker unavailability.
  5. 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

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


AI-assisted analysis of alibaba/canal@87be50e876 (2026-08-14). Data as JSON: /api/errors/28ed7b6b6862a68f. Report an issue: GitHub.