首页/新闻资讯/正文详情

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

发布时间:2026/9/23 17:07:00 来源:云帆数科 栏目:资讯中心
Apache Pulsar Canal Source Connector 实战指南:将 MySQL Binlog 实时同步到 Pulsar Topic
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Canal source connector 是 Apache Pulsar 官方提供的 CDCChange 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 开启 binlogbinlog-formatROW后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 配置启用集群模式。truecluster模式。connector 通过zkServers找到实际数据库主机。falsestandalone模式。connector 直连singleHostname和singlePort指定的 Canal server。zkServerstrueNoneZookeeper 地址和端口cluster 模式下 connector 通过它获取实际数据库主机。batchSizefalse1000每次从 Canal 拉取的批量大小。从源码看配置语义对照 CanalSourceConfig.java 中的FieldDoc注解可以进一步确认几点实现细节username与password被标记为sensitive true在配置文件与运行日志中属于敏感信息singlePort与batchSize是int类型cluster是Boolean类型并默认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 正确反序列化为int。zkServers、destination、username、password均为必填键请确保配置完整。使用示例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_PASSWORDcanal -e MYSQL_USERmysqluser -e MYSQL_PASSWORDmysqlpw mysql:5.72. 创建 MySQL 配置文件mysqld.cnfCanal 依赖 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-links0 log-binmysql-bin binlog-formatROW server_id1其中log-binmysql-bin开启 binlogbinlog-formatROW指定行级格式Canal 解析 ROW 格式才能还原每行数据的前后镜像server_id1为从库伪装的唯一 ID。3. 将mysqld.cnf复制进 MySQL 容器$ docker cp mysqld.cnf pulsar-mysql:/etc/mysql/mysql.conf.d/4. 重启 MySQL 使配置生效$ docker restart pulsar-mysql5. 创建测试数据库$ 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.scanfalse -e canal.destinationstest -e canal.instance.master.addresspulsar-mysql:3306 -e canal.instance.dbUsernameroot -e canal.instance.dbPasswordcanal -e canal.instance.connectionCharsetUTF-8 -e canal.instance.tsdb.enabletrue -e canal.instance.gtidonfalse --namepulsar-canal-server -p 8000:8000 -p 2222:2222 -p 11111:11111 -p 11112:11112 -m 4096m canal/canal-server:v1.1.2关键环境变量说明canal.destinationstest指定 destination 名称为test与后续配置文件中的destination一一对应canal.instance.master.addresspulsar-mysql:3306Canal 连接的 MySQL 主机与端口canal.instance.dbUsername/dbPasswordCanal 用于伪装从库访问 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 standalone6650是 Pulsar 客户端连接端口8080是 admin/HTTP 端口--link pulsar-canal-server使容器间可通过主机名pulsar-canal-server互通。8. 修改 connector 配置文件canal-mysql-source-config.yaml将singleHostname指向 Canal 容器主机名destination改为testconfigs: zkServers: batchSize: 5120 destination: test username: password: cluster: false singleHostname: pulsar-canal-server singlePort: 111119. 创建 Pulsar 消费脚本pulsar-client.pyimport pulsar client pulsar.Client(pulsar://localhost:6650) consumer client.subscribe(my-topic, subscription_namemy-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命令参数说明--archiveconnector 的 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.py13. 登录 MySQL 容器另开一个终端窗口$ docker exec -it pulsar-mysql /bin/bash $ mysql -h 127.0.0.1 -uroot -pcanal14. 在 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 ))ENGINEInnoDB DEFAULT CHARSETutf8; mysql INSERT INTO test_table (test_title, test_author, test_date) VALUES(a, b, NOW()); mysql UPDATE test_table SET test_titlec WHERE test_titlea; mysql DELETE FROM test_table WHERE test_titlec;每次执行 INSERT / UPDATE / DELETE第 12 步的pulsar-client.py就会打印出对应的变更消息从而实现 MySQL binlog 到 Pulsar topic 的实时同步。官方还提供了更完整的 CDC 应用场景说明可参考 io-cdc.md 与连接器总览 io-connectors.md。源码级原理消息抓取、转换与 ACK抓取主循环CanalAbstractSource所有 Canal source 共享同一个抓取骨架 CanalAbstractSource.java。其核心流程如下open加载配置后根据cluster字段选择连接方式clustertrue调用CanalConnectors.newClusterConnector(zkServers, destination, username, password)通过 Zookeeper 动态发现 Canal 集群中实际承载该 destination 的数据库主机clusterfalse调用CanalConnectors.newSingleConnector(new InetSocketAddress(singleHostname, singlePort), ...)直连指定 Canal serverCanalAbstractSource.java#L59-L73。process 主循环启动名为canal source thread的独立线程先connector.connect()再connector.subscribe()随后循环执行connector.getWithoutAck(batchSize)批量拉取消息CanalAbstractSource.java#L106-L132空批次退避当批次 ID 为 -1 或条目数为 0 时线程休眠 1 秒后继续轮询避免空转打爆 CPUACK 语义每条记录封装为CanalRecord其ack()回调调用connector.ack(id)向 Canal 确认消费成功实现“拉取后确认”的可靠投递语义CanalAbstractSource.java#L146-L179异常处理线程设置了UncaughtExceptionHandler主循环捕获所有异常并disconnect()清理连接避免进程崩溃。消息转换MessageUtilsCanal 返回的原始Message是 protobuf 结构connector 通过 MessageUtils.messageConverter() 将其转换为更易消费的FlatMessage列表转换要点包括跳过TRANSACTIONBEGIN/TRANSACTIONEND类型条目只保留真正的数据变更为每条变更记录填充database、table、typeINSERT / UPDATE / DELETE、sql、执行时间es与系统时间ts对非 DDL 事件按事件类型选择列集合DELETE 取beforeColumnsListINSERT/UPDATE 取afterColumnsListUPDATE 额外记录old变更前的值通过genColumn把每列组织为包含isKey、isNull、index、mysqlType、columnName、columnValue、updated的结构化 Map。两种输出格式CanalStringSource 与 CanalByteSource连接器提供两个 source 类区别仅在输出类型CanalStringSource.java使用 fastjson 将FlatMessage列表序列化为 JSON 字符串再封装为CanalMessage包含id、message、timestamp三个字段时间戳采用 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.client与canal.protocol版本 1.1.5负责与 Canal server 通信及解析 protobuf 协议fastjson1.2.83负责 FlatMessage 到 JSON 的序列化jackson-dataformat-yaml/jackson-databind负责 YAML / JSON 配置加载spring-core等 Spring 组件为 Canal 客户端运行提供依赖支持通过nifi-nar-maven-plugin打包为 NAR 文件供pulsar-admin source加载运行。常见问题与调优建议用户名密码填错username/password是 Canal server 的账号而非 MySQL 账号。若在非集群模式下使用空账号请确认 Canal server 侧未启用客户端鉴权否则连接会被拒绝。消息延迟或不产出先检查 MySQL 是否已开启binlog-formatROW再确认destination与 Canal 启动时canal.destinations一致最后确认singleHostname能解析到 Canal 容器。批量吞吐调优batchSize决定每次getWithoutAck拉取的条目数默认 1000。在高变更频率场景下可以适当调大如示例中的 5120减少轮询开销但也要注意单批过大带来的内存占用。cluster 模式当 Canal server 以集群方式部署、destination 实际主机通过 Zookeeper 动态分配时将cluster设为true并填写zkServers此时singleHostname/singlePort会被忽略。重复消费getWithoutAck配合ack实现至少一次语义下游消费者应基于消息内容做幂等处理DDL 事件、事务边界事件已被MessageUtils过滤如需原始事务信息应自行扩展。总结Canal source connector 是 Pulsar 生态中打通 MySQL 与消息总线的高性价比方案它复用 Canal 成熟的 binlog 解析能力通过cluster/standalone两种模式适配不同的部署形态并以统一的 source 框架提供可靠的拉取与确认机制。结合 CanalSourceConfig.java、CanalAbstractSource.java 与 MessageUtils.java 的实现你可以按需扩展新的输出格式或自定义过滤逻辑将 MySQL 变更事件稳定、低延迟地送入 Pulsar 主题支撑缓存刷新、数仓同步、搜索索引更新等下游场景。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar Canal Source Connector 实战指南基于 MySQL Binlog 实时同步数据到 Pulsar TopicApache Pulsar Canal Source Connector 实战指南基于 MySQL Binlog 实时同步数据到 Pulsar Topic C消息队列后端流处理Apache Pulsar Flume Source Connector 实战指南将 Flume Agent 日志导入 Pulsar TopicApache Pulsar Flume Source Connector 实战指南将 Flume Agent 日志导入 Pulsar Topic Flume消息队列后端流处理Refine 教程实战结合 Material UI 与 React Router 搭建带主题布局的 CRUD 应用Refine 教程实战结合 Material UI 与 React Router 搭建带主题布局的 CRUD 应用 本篇基于 Refine 官方教程中的 Ma消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

相关推荐

SAP AP应付账款教程:从供应商主数据到发票清账全流程实操指南
SAP AP应付账款教程:从供应商主数据到发票清账全流程实操指南

简介:这份SAP财务系统AP应付账款会计教程面向财务人员、SAP初学者及ERP顾问,帮助其系统掌握应付账款模块的完整业务流程与后台操作。内容围绕供应商主数据维护、AP文档处理、供应商账户余额三大模块展开,涵盖贸易与非贸易供应商的创建、修改、… · 2026/9/23 17:06:59

Python实战:从零实现智能停车管理系统(期末作业完整指南)
Python实战:从零实现智能停车管理系统(期末作业完整指南)

简介:面向Python初学者的智能停车管理系统期末作业源码包,集成OpenCV车牌识别、YOLOOCR深度学习与MySQL数据存储,完整实现车辆模拟入场、出场计费及GUI交互界面,适合K12或高校学生课程设计参考。资源共12个文件,主要由… · 2026/9/23 17:06:53

变换域通信系统TDCS:低截获信号波形设计与MATLAB实现
变换域通信系统TDCS:低截获信号波形设计与MATLAB实现

简介:围绕变换域通信系统(TDCS)的一份Matlab源码包,聚焦低截获(LPI)与抗截获信号设计,适合通信工程学生、安全通信研究者,以及希望实现隐蔽传输的工程师。压缩包内共12个文件&#x… · 2026/9/23 17:06:53

3d全息投影视频源选型避坑:2024速查手册
3d全息投影视频源选型避坑:2024速查手册

3d全息投影视频源选型避坑:2024速查手册 刚把项目里的 three.js 从 r128 升到 r160,跑起来直接白屏?控制台报 WebGL context lost ,检查代码发现 WebGLRenderer… · 2026/9/23 18:40:14

小米网关一二三代怎么选?从Zigbee到Mesh看懂智能家居中枢
小米网关一二三代怎么选?从Zigbee到Mesh看懂智能家居中枢

1. 从“智能家居死机”说起:为什么网关才是全屋智能的命门用了几年智能家居,我最大的感悟是:很多人买设备前纠结传感器买哪家、开关选什么牌子,结果装完发现设备频繁掉线、响应延迟、场景联动像个段子——大概率不是设备本身的问题… · 2026/9/23 18:40:13

薄膜技术应用全景:从光学电子到包装能源医疗的工艺实践指南
薄膜技术应用全景:从光学电子到包装能源医疗的工艺实践指南

1. 薄膜技术到底能用在哪些地方1.1 从手机屏幕到食品包装,薄膜无处不在很多人第一次听到“薄膜”这个词,脑子里浮现的可能是保鲜膜。这没错,保鲜膜确实是最贴近日常生活的薄膜制品之一,但薄膜技术的应用边界远比这宽得多。我在这个… · 2026/9/23 18:40:07

Edge作为嵌入式Web运行时的深度解析与企业级实践
Edge作为嵌入式Web运行时的深度解析与企业级实践

1. 项目概述:这不是一款“替代Chrome”的浏览器,而是一套嵌入式Web体验操作系统 Edge不是Chrome的复刻版,也不是Firefox的轻量分支。它本质上是一套以Chromium内核为底座、但深度重构了渲染管线、进程模型与安全边界的 嵌入式Web体验操作系… · 2026/9/23 18:40:07

HEED分簇协议MATLAB仿真:无线传感器网络能效与生命周期优化指南
HEED分簇协议MATLAB仿真:无线传感器网络能效与生命周期优化指南

简介:基于 MATLAB 的无线传感器网络 HEED 算法实现,面向 WSN 研究者、通信专业学生及算法仿真爱好者,用于解决分簇路由中簇头均衡选举与网络能效优化问题。该算法的核心是根据节点剩余能量与邻居分布动态选举簇头,以延长网络生命周… · 2026/9/23 18:40:00

基于 PaddleHub 的 MSGNet 风格迁移实战:从 Fine-tune 到服务化部署(PaddleFormers 仓库指南)
基于 PaddleHub 的 MSGNet 风格迁移实战:从 Fine-tune 到服务化部署(PaddleFormers 仓库指南)

人工智能预训练微调模型推理服务 【免费下载链接】PaddleFormers PaddleFormers is an easy-to-use library of pre-trained large language model zoo based on PaddlePaddle. 项目地址: https://gitcode.com/gh_mirrors/pa/PaddleFormers 点击查看 免费下载 导读… · 2026/9/23 18:39:54

3招搞定手机怎么下载微信面试难题实战项目解析
3招搞定手机怎么下载微信面试难题实战项目解析

3招搞定手机怎么下载微信面试难题实战项目解析 面试被问“手机怎么下载微信”背后的原理,90%的人答不上来。别笑,这看似弱智的问题,实则是考察你对移动应用分发机制、安全校验及网络协议理解的试金石。我带过不少校招新人,他们背了八股文,却连一个A… · 2026/9/23 0:00:03

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型
你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型

你有新短消息请注意查收:3个新手避坑指南搞定消息系统选型 面试被问“高并发下如何保证消息不丢失”,你张口就是“用Redis”,结果面试官追问“如果Redis宕机了怎么办”,你瞬间卡壳。这种场景太常见了,很多新手在背八股文时,只记住了技术名词… · 2026/9/23 0:00:29

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧
Win7无线热点配置工具源码解析:解决API失效的3个实战技巧

Win7无线热点配置工具源码解析:解决API失效的3个实战技巧 Win7无线热点配置工具在Win10/11上跑不动?不是你的问题,是版本升级后 API 全变了。很多老项目里的 netsh wlan… · 2026/9/23 0:00:36

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码