ARTICLE DETAIL

资讯详情

深耕网站建设、视觉设计与SEO优化的一线实战洞察。

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

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),仅供参考
返回列表