canal由tcp模式改为rabbitMQ模式随记
·
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:所有工作完成后,需要仔细测试,防止哪一步有遗漏,导致数据丢失
更多推荐



所有评论(0)