ps:原有canal同步mysql数据到es的功能,现如今需要将同一份数据再写入其他数据库,单纯的tcp模式已经无法支持,遂改为使用rabbitMQ模式,随记。默认已安装rabbitMQ

一、canal_server(v1.1.8)修改

1.修改./conf/canal.properties
# tcp, kafka, rocketMQ, rabbitMQ, pulsarMQ
canal.serverMode = rabbitMQ


rabbitmq.host = your.rabbitmq.ip
rabbitmq.virtual.host = /
rabbitmq.exchange = canal.exchange
rabbitmq.username = admin
rabbitmq.password = admin
rabbitmq.queue = 
rabbitmq.routingKey = 
rabbitmq.deliveryMode = 
# 1. 端口(默认就是5672,但如果你的RabbitMQ改了端口,这里必须指定)
rabbitmq.port = 5672
# 2. 交换机类型(direct/topic/fanout)。Canal的MQ模式通常使用 `direct` 或 `topic` 以便路由。
rabbitmq.exchange.type = direct
# 3. 是否异步发送(强烈建议为true以提升吞吐量)
rabbitmq.asyncSend = true
二.修改./conf/example/instance.properties
# mq config
canal.mq.topic=canal
canal.mq.partition=0
# 关键配置:指定使用flatMessage格式(即JSON格式)
canal.mq.flatMessage = true

二、rabbitMQ修改

1.创建exchange(和配置中rabbitmq.exchange对应)

在这里插入图片描述

2.创建queue(和配置中rabbitmq.queue对应,如果为空,自定义)

在这里插入图片描述

3.点击进入queue,绑定exchange和routing key(和配置中canal.mq.topic对应)

在这里插入图片描述

4.所有东西操作完成后,可以重启canal_server,按照之前的操作,查看rabbitMQ是否有写入数据。如果能看到流量数据,说明已经写入成功

在这里插入图片描述

三、修改canal_adapter(版本1.1.5以上支持直接接入mq)

1.修改conf/application.yml(和canal_server对应)
canal.conf:
  mode: rabbitMQ #tcp kafka rocketMQ rabbitMQ
  
  consumerProperties:
    # rabbitMQ consumer
    rabbitmq.host: your.rabbitmq.ip
    rabbitmq.port: 5672
    rabbitmq.virtual.host: /
    rabbitmq.username: admin
    rabbitmq.password: admin
    rabbitmq.resource.ownerId:

  canalAdapters:
  - instance: adapter.queue # canal instance Name or mq topic name
2.修改conf/es8/yourTable.yml中destination配置项,注意和conf/application.yml中canalAdapters.instance对应
PS:此时先关闭canal_adapter,重启canal_server,再重启canal_adapter,查看以往流程是否通顺。

四、接入java服务

1.新建新的queue(这里假定是canal.queue),同样绑定到之前的exchange和routing key
2.编写java功能

1.引入maven依赖

        <dependency>
            <groupId>org.springframework.boot</groupId>
            <artifactId>spring-boot-starter-amqp</artifactId>
        </dependency>
        <dependency>
            <groupId>com.alibaba.otter</groupId>
            <artifactId>canal.protocol</artifactId>
            <version>1.1.8</version>
            <scope>compile</scope>
        </dependency>

2.添加配置

spring:  
  rabbitmq:
    host: your.rabbitmq.ip
    port: 5672
    username: admin
    password: admin
    virtual-host: /
    # 建议开启手动确认,防止消息丢失
    listener:
      simple:
        acknowledge-mode: manual
        prefetch: 10 # 根据处理能力调整
canal:
  queueName: canal.queue
  exchange: canal.exchange
  routingKey: canal

3.编写配置类

@Configuration
public class RabbitConfig {

    @Value("${canal.queueName}")
    private String queueName;

    @Bean
    public Queue canalFileQueue() {
        // 创建持久化队列
        return new Queue(queueName, true);
    }

    @Bean
    public Binding bindingFile(Queue canalFileQueue,
                               @Value("${canal.exchange:canal.exchange}") String exchange,
                               @Value("${canal.routingKey:canal}") String routingKey) {
        // 将队列绑定到Canal Server使用的Exchange,并指定routingKey
        return BindingBuilder
                .bind(canalFileQueue)
                .to(new DirectExchange(exchange))
                .with(routingKey);
    }
}

4.编写消费者

@RabbitListener(queues = "${canal.queueName}")
    public void handleMessage(Message message, Channel channel) {
        log.info("消息处理开始");
        try {
            // 1. 解析消息
            String jsonString = new String(message.getBody(), StandardCharsets.UTF_8);
            FlatMessage flatMessage = JSON.parseObject(jsonString, FlatMessage.class);
            long deliveryTag = message.getMessageProperties().getDeliveryTag();
            // 2.执行业务逻辑
            // 3. 任务执行完毕后,确认消息
            channel.basicAck(deliveryTag, false);

            log.info("消息处理完成 - 类型: {}" , flatMessage.getType() );
        } catch (Exception e) {
            log.error("消息处理失败: " + e.getMessage());
            // 拒绝消息,不重新入队(避免死循环)
            try {
                long deliveryTag = message.getMessageProperties().getDeliveryTag();
                channel.basicNack(deliveryTag, false, true);
            } catch (IOException ioException) {
                log.error("Failed to nack message", ioException);
            }
        }
    }
ps:所有工作完成后,需要仔细测试,防止哪一步有遗漏,导致数据丢失
Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐