news 2026/9/23 17:06:54

Apache Pulsar Canal Source Connector 实战指南:将 MySQL Binlog 实时同步到 Pulsar Topic

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar Canal Source Connector 实战指南:将 MySQL Binlog 实时同步到 Pulsar Topic
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

Canal source connector 是 Apache Pulsar 官方提供的 CDC(Change Data Capture,变更数据捕获)连接器,它对接阿里巴巴开源的 Canal 中间件,把 MySQL 的 binlog 变更事件实时拉取并写入 Pulsar topic。阅读本文后,你将掌握 Canal source connector 的全部配置项及其底层含义、两种运行模式(cluster / standalone)的取舍、从零搭建 MySQL + Canal + Pulsar 全链路数据同步的完整步骤,并能通过源码理解消息的抓取、转换与 ACK 机制。

工作原理概述

Canal source connector 的核心职责是“从 Canal 拉数据,向 Pulsar 写数据”:它本身不直接读取 MySQL,而是作为 Canal 的客户端,订阅 Canal server 已经解析好的 binlog 变更消息,再将每条变更记录转换为 Pulsar 消息发布到指定 topic。

  • 数据来源:MySQL 开启 binlog(binlog-format=ROW)后,Canal server 伪装成 MySQL 从库解析 binlog;
  • 中间环节:Canal server 将解析结果以 protobufMessage形式提供给客户端;
  • Pulsar 侧:connector 通过pulsar-admin source以 source connector 形式运行,产出消息进入目标 topic。

连接器的入口实现位于 pulsar-io/canal/src/main/java/org/apache/pulsar/io/canal/,其中CanalStringSource是默认使用的 source 类,服务发现文件 pulsar-io/canal/src/main/resources/META-INF/services/pulsar-io.yaml 中的声明如下:

name: canal description: canal source and read data from mysql sourceClass: org.apache.pulsar.io.canal.CanalStringSource sourceConfigClass: org.apache.pulsar.io.canal.CanalSourceConfig

配置项详解

Canal source connector 的配置由 CanalSourceConfig.java 定义,支持 YAML、JSON 或键值对形式加载(load(String yamlFile)使用 Jackson YAML 解析,load(Map)使用 JSON 序列化后反序列化)。其属性如下:

名称是否必填默认值描述
usernametrueNoneCanal server 的账号(注意:不是 MySQL 账号)。
passwordtrueNoneCanal server 的密码(不是 MySQL 密码)。
destinationtrueNoneCanal source connector 所连接的目标 destination(即 Canal 中配置的实例名称)。
singleHostnamefalseNoneCanal server 地址。
singlePortfalseNoneCanal server 端口。
clustertruefalse是否基于 Canal server 配置启用集群模式。

  • true:cluster模式。connector 通过zkServers找到实际数据库主机。
  • false:standalone模式。connector 直连singleHostnamesinglePort指定的 Canal server。
zkServerstrueNoneZookeeper 地址和端口,cluster 模式下 connector 通过它获取实际数据库主机。
batchSizefalse1000每次从 Canal 拉取的批量大小。

从源码看配置语义

对照 CanalSourceConfig.java 中的@FieldDoc注解,可以进一步确认几点实现细节:

  • usernamepassword被标记为sensitive = true,在配置文件与运行日志中属于敏感信息;
  • singlePortbatchSizeint类型,clusterBoolean类型并默认false,其余字段均为字符串;
  • batchSize的默认值在代码中为1000,与文档一致;该值会直接传给 Canal 客户端的getWithoutAck(batchSize),决定一次拉取的消息条数上限;
  • 字段注释help与官方文档属性表一一对应,说明配置文件、文档与代码三者保持一致。

配置示例

官方提供了可直接使用的示例配置文件 canal-mysql-source-config.yaml。使用连接器前,可通过以下任一方式创建配置文件。

  • JSON 格式
{ "zkServers": "127.0.0.1:2181", "batchSize": "5120", "destination": "example", "username": "", "password": "", "cluster": false, "singleHostname": "127.0.0.1", "singlePort": "11111" }
  • YAML 格式
configs: zkServers: "127.0.0.1:2181" batchSize: 5120 destination: "example" username: "" password: "" cluster: false singleHostname: "127.0.0.1" singlePort: 11111

注意:YAML 与 JSON 两种写法均使用configs作为外层键(JSON 示例中该键隐含于--source-config-file的解析逻辑,实际提交时以 YAML 文件或--source-config字符串为准);singlePort在 YAML 中通常写为整数,在 JSON 中写为字符串也能被 Jackson 正确反序列化为intzkServersdestinationusernamepassword均为必填键,请确保配置完整。

使用示例:MySQL 数据实时同步到 Pulsar

下面以官方文档的完整流程为例,演示如何基于上述配置文件将 MySQL 数据同步到 Pulsar。整个流程包含 MySQL、Canal、Pulsar 三个容器以及一个 Pulsar 消费端脚本。

1. 启动 MySQL 服务器

$ docker pull mysql:5.7 $ docker run -d -it --rm --name pulsar-mysql -p 3306:3306 -e MYSQL_ROOT_PASSWORD=canal -e MYSQL_USER=mysqluser -e MYSQL_PASSWORD=mysqlpw mysql:5.7

2. 创建 MySQL 配置文件mysqld.cnf

Canal 依赖 MySQL 开启 binlog 且格式为 ROW,因此需要如下配置:

[mysqld] pid-file = /var/run/mysqld/mysqld.pid socket = /var/run/mysqld/mysqld.sock datadir = /var/lib/mysql #log-error = /var/log/mysql/error.log # By default we only accept connections from localhost #bind-address = 127.0.0.1 # Disabling symbolic-links is recommended to prevent assorted security risks symbolic-links=0 log-bin=mysql-bin binlog-format=ROW server_id=1

其中log-bin=mysql-bin开启 binlog,binlog-format=ROW指定行级格式(Canal 解析 ROW 格式才能还原每行数据的前后镜像),server_id=1为从库伪装的唯一 ID。

3. 将mysqld.cnf复制进 MySQL 容器

$ docker cp mysqld.cnf pulsar-mysql:/etc/mysql/mysql.conf.d/

4. 重启 MySQL 使配置生效

$ docker restart pulsar-mysql

5. 创建测试数据库

$ docker exec -it pulsar-mysql /bin/bash $ mysql -h 127.0.0.1 -uroot -pcanal -e 'create database test;'

6. 启动 Canal server 并连接 MySQL

$ docker pull canal/canal-server:v1.1.2 $ docker run -d -it --link pulsar-mysql -e canal.auto.scan=false -e canal.destinations=test -e canal.instance.master.address=pulsar-mysql:3306 -e canal.instance.dbUsername=root -e canal.instance.dbPassword=canal -e canal.instance.connectionCharset=UTF-8 -e canal.instance.tsdb.enable=true -e canal.instance.gtidon=false --name=pulsar-canal-server -p 8000:8000 -p 2222:2222 -p 11111:11111 -p 11112:11112 -m 4096m canal/canal-server:v1.1.2

关键环境变量说明:

  • canal.destinations=test:指定 destination 名称为test,与后续配置文件中的destination一一对应;
  • canal.instance.master.address=pulsar-mysql:3306:Canal 连接的 MySQL 主机与端口;
  • canal.instance.dbUsername/dbPassword:Canal 用于伪装从库访问 MySQL 的账号密码;
  • 暴露的11111端口即 Canal 客户端(Pulsar connector)的连接端口,对应配置中的singlePort

7. 启动 Pulsar standalone

$ docker pull apachepulsar/pulsar:2.3.0 $ docker run -d -it --link pulsar-canal-server -p 6650:6650 -p 8080:8080 -v $PWD/data:/pulsar/data --name pulsar-standalone apachepulsar/pulsar:2.3.0 bin/pulsar standalone

6650是 Pulsar 客户端连接端口,8080是 admin/HTTP 端口;--link pulsar-canal-server使容器间可通过主机名pulsar-canal-server互通。

8. 修改 connector 配置文件canal-mysql-source-config.yaml

singleHostname指向 Canal 容器主机名,destination改为test

configs: zkServers: "" batchSize: "5120" destination: "test" username: "" password: "" cluster: false singleHostname: "pulsar-canal-server" singlePort: "11111"

9. 创建 Pulsar 消费脚本pulsar-client.py

import pulsar client = pulsar.Client('pulsar://localhost:6650') consumer = client.subscribe('my-topic', subscription_name='my-sub') while True: msg = consumer.receive() print("Received message: '%s'" % msg.data()) consumer.acknowledge(msg) client.close()

10. 将配置文件和消费脚本复制进 Pulsar 容器

$ docker cp canal-mysql-source-config.yaml pulsar-standalone:/pulsar/conf/ $ docker cp pulsar-client.py pulsar-standalone:/pulsar/

11. 下载 Canal connector 并启动

$ docker exec -it pulsar-standalone /bin/bash $ wget https://archive.apache.org/dist/pulsar/pulsar-2.3.0/connectors/pulsar-io-canal-2.3.0.nar -P connectors $ ./bin/pulsar-admin source localrun \ --archive ./connectors/pulsar-io-canal-2.3.0.nar \ --classname org.apache.pulsar.io.canal.CanalStringSource \ --tenant public \ --namespace default \ --name canal \ --destination-topic-name my-topic \ --source-config-file /pulsar/conf/canal-mysql-source-config.yaml \ --parallelism 1

命令参数说明:

  • --archive:connector 的 NAR 包路径(pulsar-io-canal 模块通过nifi-nar-maven-plugin打包,见 pulsar-io/canal/pom.xml);
  • --classname:指定 source 实现类,此处为CanalStringSource
  • --destination-topic-name my-topic:变更数据将写入my-topic
  • --source-config-file:指向第 8 步修改后的 YAML 配置;
  • localrun模式表示在本地进程内以函数运行时方式运行 connector,便于快速验证。

12. 在 Pulsar 容器内运行消费脚本

$ docker exec -it pulsar-standalone /bin/bash $ python pulsar-client.py

13. 登录 MySQL 容器

另开一个终端窗口:

$ docker exec -it pulsar-mysql /bin/bash $ mysql -h 127.0.0.1 -uroot -pcanal

14. 在 MySQL 中建表并执行增删改

mysql> use test; mysql> show tables; mysql> CREATE TABLE IF NOT EXISTS `test_table`(`test_id` INT UNSIGNED AUTO_INCREMENT,`test_title` VARCHAR(100) NOT NULL, `test_author` VARCHAR(40) NOT NULL, `test_date` DATE,PRIMARY KEY ( `test_id` ))ENGINE=InnoDB DEFAULT CHARSET=utf8; mysql> INSERT INTO test_table (test_title, test_author, test_date) VALUES("a", "b", NOW()); mysql> UPDATE test_table SET test_title='c' WHERE test_title='a'; mysql> DELETE FROM test_table WHERE test_title='c';

每次执行 INSERT / UPDATE / DELETE,第 12 步的pulsar-client.py就会打印出对应的变更消息,从而实现 MySQL binlog 到 Pulsar topic 的实时同步。官方还提供了更完整的 CDC 应用场景说明,可参考 io-cdc.md 与连接器总览 io-connectors.md。

源码级原理:消息抓取、转换与 ACK

抓取主循环:CanalAbstractSource

所有 Canal source 共享同一个抓取骨架 CanalAbstractSource.java。其核心流程如下:

  1. open:加载配置后,根据cluster字段选择连接方式:
    • cluster=true:调用CanalConnectors.newClusterConnector(zkServers, destination, username, password),通过 Zookeeper 动态发现 Canal 集群中实际承载该 destination 的数据库主机;
    • cluster=false:调用CanalConnectors.newSingleConnector(new InetSocketAddress(singleHostname, singlePort), ...),直连指定 Canal server(CanalAbstractSource.java#L59-L73)。
  2. process 主循环:启动名为canal source thread的独立线程,先connector.connect()connector.subscribe(),随后循环执行connector.getWithoutAck(batchSize)批量拉取消息(CanalAbstractSource.java#L106-L132);
  3. 空批次退避:当批次 ID 为 -1 或条目数为 0 时,线程休眠 1 秒后继续轮询,避免空转打爆 CPU;
  4. ACK 语义:每条记录封装为CanalRecord,其ack()回调调用connector.ack(id)向 Canal 确认消费成功,实现“拉取后确认”的可靠投递语义(CanalAbstractSource.java#L146-L179);
  5. 异常处理:线程设置了UncaughtExceptionHandler,主循环捕获所有异常并disconnect()清理连接,避免进程崩溃。

消息转换:MessageUtils

Canal 返回的原始Message是 protobuf 结构,connector 通过 MessageUtils.messageConverter() 将其转换为更易消费的FlatMessage列表,转换要点包括:

  • 跳过TRANSACTIONBEGIN/TRANSACTIONEND类型条目,只保留真正的数据变更;
  • 为每条变更记录填充databasetabletype(INSERT / UPDATE / DELETE)、sql、执行时间es与系统时间ts
  • 对非 DDL 事件,按事件类型选择列集合:DELETE 取beforeColumnsList,INSERT/UPDATE 取afterColumnsList,UPDATE 额外记录old变更前的值;
  • 通过genColumn把每列组织为包含isKeyisNullindexmysqlTypecolumnNamecolumnValueupdated的结构化 Map。

两种输出格式:CanalStringSource 与 CanalByteSource

连接器提供两个 source 类,区别仅在输出类型:

  • CanalStringSource.java:使用 fastjson 将FlatMessage列表序列化为 JSON 字符串,再封装为CanalMessage(包含idmessagetimestamp三个字段,时间戳采用 ISO8601 带时区格式)。其@Connector注解注明“方便 Presto SQL 查询”,即输出为结构化的 JSON 文本,便于后续用 SQL 直接检索(CanalStringSource.java#L38-L42)。这也是pulsar-io.yaml中默认注册的 source 类;
  • CanalByteSource.java:同样先用 fastjson 序列化,但最终输出为byte[]字节数组,适合下游自行解析 JSON 的二进制消费场景。

两者都继承自CanalAbstractSource,因此共享上述连接、拉取、ACK 的全套逻辑,仅extractValue与消息 ID 提取策略不同。

依赖与打包

从 pulsar-io/canal/pom.xml 可以看出该模块的技术栈:

  • canal.clientcanal.protocol(版本 1.1.5):负责与 Canal server 通信及解析 protobuf 协议;
  • fastjson(1.2.83):负责 FlatMessage 到 JSON 的序列化;
  • jackson-dataformat-yaml/jackson-databind:负责 YAML / JSON 配置加载;
  • spring-core等 Spring 组件:为 Canal 客户端运行提供依赖支持;
  • 通过nifi-nar-maven-plugin打包为 NAR 文件,供pulsar-admin source加载运行。

常见问题与调优建议

  1. 用户名密码填错username/password是 Canal server 的账号,而非 MySQL 账号。若在非集群模式下使用空账号,请确认 Canal server 侧未启用客户端鉴权,否则连接会被拒绝。
  2. 消息延迟或不产出:先检查 MySQL 是否已开启binlog-format=ROW;再确认destination与 Canal 启动时canal.destinations一致;最后确认singleHostname能解析到 Canal 容器。
  3. 批量吞吐调优batchSize决定每次getWithoutAck拉取的条目数,默认 1000。在高变更频率场景下可以适当调大(如示例中的 5120),减少轮询开销;但也要注意单批过大带来的内存占用。
  4. cluster 模式:当 Canal server 以集群方式部署、destination 实际主机通过 Zookeeper 动态分配时,将cluster设为true并填写zkServers,此时singleHostname/singlePort会被忽略。
  5. 重复消费getWithoutAck配合ack实现至少一次语义,下游消费者应基于消息内容做幂等处理;DDL 事件、事务边界事件已被MessageUtils过滤,如需原始事务信息应自行扩展。

总结

Canal source connector 是 Pulsar 生态中打通 MySQL 与消息总线的高性价比方案:它复用 Canal 成熟的 binlog 解析能力,通过cluster/standalone两种模式适配不同的部署形态,并以统一的 source 框架提供可靠的拉取与确认机制。结合 CanalSourceConfig.java、CanalAbstractSource.java 与 MessageUtils.java 的实现,你可以按需扩展新的输出格式或自定义过滤逻辑,将 MySQL 变更事件稳定、低延迟地送入 Pulsar 主题,支撑缓存刷新、数仓同步、搜索索引更新等下游场景。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/23 17:02:12

谐波潮流计算与解耦方法详解:从原理到Python实现

简介:面向电力系统谐波分析场景的MATLAB程序包,专注谐波潮流计算与谐波解耦算法,适合电气工程专业学生、科研人员以及从事电能质量治理的工程师使用。当前电网中开关电源、整流器等非线性负载大量接入,谐波畸变已成为影响电能质量…

作者头像 李华
网站建设 2026/9/23 17:01:44

Ude.NET智能编码检测:解决乱码问题的终极方案

1. 编码问题的困扰与解决之道第一次接手遗留系统时,我被满屏的"锟斤拷"和"烫烫烫"震惊了。这些乱码不仅让数据无法正常显示,更导致业务逻辑出现严重错误。后来排查发现,问题出在系统对接第三方数据时没有正确处理字符编码…

作者头像 李华
网站建设 2026/9/23 17:01:41

校园智能外卖配送系统设计与实践

1. 项目背景与需求解析校园外卖配送这个细分市场近年来呈现出爆发式增长。根据我们团队在10所高校的实地调研,平均每所大学每天产生的外卖订单量超过5000单,高峰期配送人员进出校园频次可达200人次/小时。这种高频次的人员流动带来了三个核心痛点&#x…

作者头像 李华
网站建设 2026/9/23 16:59:28

灰狼优化算法实战:SVM超参数自动调优

简介:本资源是面向机器学习初学者与算法实践者的灰狼优化算法(GWO)与支持向量机(SVM)融合实现方案,聚焦SVM核函数参数与惩罚系数的自动寻优难题,适用于分类、回归及异常检测等典型任务。压缩包共…

作者头像 李华
网站建设 2026/9/23 16:56:49

坦克检测数据集1520张VOC+YOLO双格式:从数据校验到YOLOv8训练全流程

简介:这份资源面向目标检测初学者与需要快速验证模型的研究者,提供一套可直接用于YOLO训练的坦克检测数据集,解决单一类别目标检测任务中样本获取与标注成本高的问题。压缩包共2000个文件,以1521个VOC格式xml标注文件和479个txt文…

作者头像 李华