☰
Apache Pulsar CDC 连接器实战指南:基于 Canal 与 Debezium 的数据库变更捕获方案
2026/9/25 3:21:52 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

CDC(Change Data Capture,变更数据捕获)是构建实时数据管道的关键技术,Apache Pulsar 通过内置的 Canal 与 Debezium 两类 source connector,将 MySQL、MongoDB、PostgreSQL 等数据库的日志变更持续流入 Pulsar 集群,并以持久化、复制、分区的 topic 形式对外提供数据。本文以仓库文档 io-cdc.md 为主线,结合 io-cdc-debezium.md 的完整配置与实战示例,以及pulsar-io/canal、pulsar-io/debezium模块的源码实现,系统讲解 Pulsar CDC 连接器的选型、配置、部署与排错,读者可据此在自有集群中快速搭建数据库到 Pulsar 的实时同步链路。

一、Pulsar CDC 连接器总览

CDC source connector 的核心职责是捕获数据库的日志变更(如 MySQL binlog、MongoDB oplog、PostgreSQL WAL),并把变更事件写入 Pulsar。官方文档明确指出:

CDC source connectors are built on top of Canal and Debezium and store all data into Pulsar cluster in a persistent, replicated, and partitioned way.

也就是说,Pulsar 并未重新发明一套 CDC 引擎,而是把业界成熟的 Canal 与 Debezium 能力包装成 Pulsar IO 框架下的 source connector,从而天然继承 Pulsar 的多租户、持久化、复制与分区特性。

当前 Pulsar 提供的 CDC 连接器如下表(来源:io-cdc.md):

名称Java 类
Canal source connectororg.apache.pulsar.io.canal.CanalStringSource
Debezium source connectororg.apache.pulsar.io.debezium.DebeziumSource、org.apache.pulsar.io.debezium.mysql.DebeziumMysqlSource、org.apache.pulsar.io.debezium.postgres.DebeziumPostgresSource

需要说明的是,从当前仓库源码结构看,Debezium 连接器家族实际比表格列出的更完整:在 pulsar-io/debezium 目录下,除 MySQL、PostgreSQL 外,还实现了 MongoDB(DebeziumMongoDbSource.java)、MSSQL(DebeziumMsSqlSource.java)与 Oracle(DebeziumOracleSource.java)等实现,各数据库分支均有对应的 NAR 归档与示例配置文件(如 debezium-mysql-source-config.yaml)。

两条技术路线的详细使用指南分别见 CDC Canal Connector 文档 与 Debezium source connector 文档。

二、Canal source connector:MySQL binlog 同步

Canal 是阿里巴巴开源的中继组件,通过模拟 MySQL 主从复制协议读取 binlog。Pulsar 的 Canal source connector 位于 pulsar-io/canal 模块,提供两个入口类:

  • CanalStringSource:将变更事件序列化为 JSON 字符串,适合与 Pulsar SQL/Presto 联合做 SQL 检索;
  • CanalByteSource:输出原始字节数组。

2.1 源码架构:PushSource 之上的拉取循环

从源码看,两个入口类都继承自 CanalAbstractSource.java,这是一个继承PushSource<V>的抽象类,其open()方法完成 Canal 连接器初始化:若cluster=true,则通过CanalConnectors.newClusterConnector(zkServers, ...)走 ZooKeeper 集群模式;否则通过newSingleConnector(InetSocketAddress(singleHostname, singlePort), ...)直连单机 Canal server。

随后后台线程process()进入核心拉取循环:

  1. connector.connect()建立连接,connector.subscribe()订阅目标;
  2. connector.getWithoutAck(batchSize)拉取一批 binlog 消息;
  3. 通过 MessageUtils.messageConverter 将 protobuf 原消息转换为FlatMessage(扁平化后的列结构,包含isKey、isNull、mysqlType、columnName、columnValue等字段);
  4. 封装为CanalRecord交给 Pulsar 框架消费;在CanalRecord.ack()中回调connector.ack(batchId),实现精确一次/至少一次的语义配合(batchId 为 -1 或空批次时 sleep 1 秒再继续)。

以 CanalStringSource.java 为例,其输出消息结构为CanalMessage:{ id, message, timestamp },其中message是JSON.toJSONString(flatMessages, WriteMapNullValue)的结果,timestamp为 ISO8601 格式,方便下游按时间检索。

2.2 Canal 配置参数

配置项在 CanalSourceConfig.java 中定义,核心字段如下:

参数必填默认值说明
username是空连接 MySQL 的用户名(敏感项)
password是空连接 MySQL 的密码(敏感项)
destination是空Canal 实例名(Canal destination),即 Canal source connector 要连接的目标
singleHostname否空单机模式下 MySQL/Canal server 主机名
singlePort否空单机模式下端口
cluster否falsetrue时通过zkServers发现真实数据库主机(集群模式);false时直连singleHostname:singlePort
zkServers是空集群模式下使用的 ZooKeeper 地址,用于发现 Canal server 列表
batchSize否1000每次从 Canal 拉取的批大小

字段上的@FieldDoc注解同时会被 Pulsar Admin 工具读取,用于生成连接器配置帮助信息。

三、Debezium source connector:多数据库统一方案

Debezium 是 Red Hat 主导的分布式 CDC 框架。Pulsar 的 Debezium source connector 位于 pulsar-io/debezium,通过把 Debezium 的 Kafka Connect 任务包装成 Pulsar source 运行,将 MySQL、PostgreSQL、MongoDB 等数据库的变更消息直接写入 Pulsar topic。

3.1 配置参数全表

以下参数表完整继承自 io-cdc-debezium.md,是配置 Debezium source connector 的依据:

参数必填默认值说明
task.class是nullDebezium 实现的具体 source task 类
database.hostname是null数据库服务器地址
database.port是null数据库服务器端口
database.user是null具备所需权限的数据库用户名
database.password是null对应密码
database.server.id是null连接器标识,必须在数据库集群内唯一,类似 MySQLserver-id
database.server.name是null数据库服务器/集群的逻辑名称,构成命名空间,用于 Kafka topic 名、Kafka Connect schema 名及 Avro schema 命名空间
database.whitelist否null该服务器上被连接器监控的数据库列表;可选,另有其他属性可控制库表的包含/排除
key.converter是nullKafka Connect 提供的记录 key 转换器
value.converter是nullKafka Connect 提供的记录 value 转换器
database.history是null数据库历史类名
database.history.pulsar.topic是null连接器写入并恢复 DDL 语句的数据库历史 topic。注意:该 topic 仅供内部使用,消费者不应使用
database.history.pulsar.service.url是null历史 topic 使用的 Pulsar 集群服务地址
pulsar.service.url是nullDebezium 偏移量 topic 使用的 Pulsar 集群服务地址,可用bin/pulsar-admin --admin-url http://pulsar:8080 sources localrun --source-config-file configs/pg-pulsar-config.yaml指定目标集群
offset.storage.topic是null记录连接器已成功提交的最近偏移量
mongodb.hosts是nullMongoDB 副本集主机端口逗号分隔列表(host或host:port形式)
mongodb.name是null标识连接器及其监控的 MongoDB 副本集/共享集群的唯一名称,每个服务器至多由一个 Debezium 连接器监控
mongodb.user是null连接 MongoDB 的数据库用户名(仅在启用认证时需要)
mongodb.password是null连接 MongoDB 的密码(仅在启用认证时需要)
mongodb.task.id是nullMongoDB 连接器 taskId,用于为每个副本集分配独立 task

3.2 源码要点:默认值与内部 topic 的自动推导

打开 DebeziumSource.java 可以看到连接器启动时的一连串默认值填充逻辑:

  • key.converter/value.converter缺省时统一置为org.apache.kafka.connect.json.JsonConverter;
  • database.history缺省时置为 Pulsar 自研的org.apache.pulsar.io.debezium.PulsarDatabaseHistory;
  • 未显式提供历史 topic 时,按{tenant}/{namespace}/{sourceName}-debezium-history-topic自动生成;
  • 未显式提供偏移量 topic 时,按{tenant}/{namespace}/{sourceName}-debezium-offset-topic自动生成;
  • 若未配置database.history.pulsar.service.url,则会序列化 Pulsar ClientBuilder 传入历史实现,从而复用 source 运行实例所在集群的客户端;
  • database.user、database.password支持从 Pulsar Functions 的 Secret 中加载(tryLoadingConfigSecret)。

这里特别值得关注的是 PulsarDatabaseHistory.java:它实现了 Debezium 的DatabaseHistorySPI,把数据库的 schema 变更(DDL)以普通 Pulsar 消息的形式写入指定 topic,并能在连接器重启时通过 reader 重放该 topic 恢复历史。这正是在 Pulsar 之上替换 Kafka 的FileDatabaseHistory/KafkaDatabaseHistory的关键所在——schema 历史不再依赖文件或 Kafka,而是落在 Pulsar 自身的持久化存储上,配合offset.storage.topic记录消费位点,共同保证重启后增量续传与 schema 一致性。

3.3 MySQL 实战示例

配置文件(JSON 与 YAML 两种形态)

JSON 形态:

{ "configs": { "database.hostname": "localhost", "database.port": "3306", "database.user": "debezium", "database.password": "dbz", "database.server.id": "184054", "database.server.name": "dbserver1", "database.whitelist": "inventory", "database.history": "org.apache.pulsar.io.debezium.PulsarDatabaseHistory", "database.history.pulsar.topic": "history-topic", "database.history.pulsar.service.url": "pulsar://127.0.0.1:6650", "key.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "pulsar.service.url": "pulsar://127.0.0.1:6650", "offset.storage.topic": "offset-topic" } }

YAML 形态(对应仓库中的 debezium-mysql-source-config.yaml):

tenant: "public" namespace: "default" name: "debezium-mysql-source" topicName: "debezium-mysql-topic" archive: "connectors/pulsar-io-debezium-mysql-@pulsar:version@.nar" parallelism: 1 configs: ## config for mysql, docker image: debezium/example-mysql:0.8 database.hostname: "localhost" database.port: "3306" database.user: "debezium" database.password: "dbz" database.server.id: "184054" database.server.name: "dbserver1" database.whitelist: "inventory" database.history: "org.apache.pulsar.io.debezium.PulsarDatabaseHistory" database.history.pulsar.topic: "history-topic" database.history.pulsar.service.url: "pulsar://127.0.0.1:6650" ## KEY_CONVERTER_CLASS_CONFIG, VALUE_CONVERTER_CLASS_CONFIG key.converter: "org.apache.kafka.connect.json.JsonConverter" value.converter: "org.apache.kafka.connect.json.JsonConverter" ## PULSAR_SERVICE_URL_CONFIG pulsar.service.url: "pulsar://127.0.0.1:6650" ## OFFSET_STORAGE_TOPIC_CONFIG offset.storage.topic: "offset-topic"
启动步骤
  1. 启动带示例数据库的 MySQL 容器(Debezium 官方镜像自带 inventory 库):

    $ docker run -it --rm \ --name mysql \ -p 3306:3306 \ -e MYSQL_ROOT_PASSWORD=debezium \ -e MYSQL_USER=mysqluser \ -e MYSQL_PASSWORD=mysqlpw debezium/example-mysql:0.8
  2. 本地以 standalone 模式启动 Pulsar:

    $ bin/pulsar standalone
  3. 以 localrun 模式启动连接器。两种方式任选其一:

    • 使用JSON配置(需确保connectors/pulsar-io-debezium-mysql-@pulsar:version@.nar存在):

      $ bin/pulsar-admin source localrun \ --archive connectors/pulsar-io-debezium-mysql-@pulsar:version@.nar \ --name debezium-mysql-source --destination-topic-name debezium-mysql-topic \ --tenant public \ --namespace default \ --source-config '{"database.hostname": "localhost","database.port": "3306","database.user": "debezium","database.password": "dbz","database.server.id": "184054","database.server.name": "dbserver1","database.whitelist": "inventory","database.history": "org.apache.pulsar.io.debezium.PulsarDatabaseHistory","database.history.pulsar.topic": "history-topic","database.history.pulsar.service.url": "pulsar://127.0.0.1:6650","key.converter": "org.apache.kafka.connect.json.JsonConverter","value.converter": "org.apache.kafka.connect.json.JsonConverter","pulsar.service.url": "pulsar://127.0.0.1:6650","offset.storage.topic": "offset-topic"}'
    • 使用YAML配置:

      $ bin/pulsar-admin source localrun \ --source-config-file debezium-mysql-source-config.yaml
  4. 订阅inventory.products表对应的变更 topic:

    $ bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0

    注意 topic 名由database.server.name+ 库名 + 表名构成,即dbserver1.inventory.products。

  5. 在 docker 中启动 MySQL 客户端:

    $ docker run -it --rm \ --name mysqlterm \ --link mysql \ --rm mysql:5.7 sh \ -c 'exec mysql -h"$MYSQL_PORT_3306_TCP_ADDR" -P"$MYSQL_PORT_3306_TCP_PORT" -uroot -p"$MYSQL_ENV_MYSQL_ROOT_PASSWORD"'
  6. 客户端弹出后,执行变更语句:

    mysql> use inventory; mysql> show tables; mysql> SELECT * FROM products; mysql> UPDATE products SET name='1111111111' WHERE id=101; mysql> UPDATE products SET name='1111111111' WHERE id=107;

    此时在订阅 topic 的终端窗口中即可看到products表的变更数据被完整保留在sub-products主题中。

3.4 PostgreSQL 实战示例

配置文件

JSON 形态:

{ "database.hostname": "localhost", "database.port": "5432", "database.user": "postgres", "database.password": "postgres", "database.dbname": "postgres", "database.server.name": "dbserver1", "schema.whitelist": "inventory", "pulsar.service.url": "pulsar://127.0.0.1:6650" }

YAML 形态(对应 debezium-postgres-source-config.yaml):

tenant: "public" namespace: "default" name: "debezium-postgres-source" topicName: "debezium-postgres-topic" archive: "connectors/pulsar-io-debezium-postgres-@pulsar:version@.nar" parallelism: 1 configs: ## config for pg, docker image: debezium/example-postgress:0.8 database.hostname: "localhost" database.port: "5432" database.user: "postgres" database.password: "postgres" database.dbname: "postgres" database.server.name: "dbserver1" schema.whitelist: "inventory" ## PULSAR_SERVICE_URL_CONFIG pulsar.service.url: "pulsar://127.0.0.1:6650"
启动步骤
  1. 启动 PostgreSQL 容器:

    $ docker pull debezium/example-postgres:0.8 $ docker run -d -it --rm --name pulsar-postgresql -p 5432:5432 debezium/example-postgres:0.8
  2. 启动 Pulsar standalone:bin/pulsar standalone。

  3. 启动连接器(JSON 或 YAML 方式):

    $ bin/pulsar-admin source localrun \ --archive connectors/pulsar-io-debezium-postgres-@pulsar:version@.nar \ --name debezium-postgres-source \ --destination-topic-name debezium-postgres-topic \ --tenant public \ --namespace default \ --source-config '{"database.hostname": "localhost","database.port": "5432","database.user": "postgres","database.password": "postgres","database.dbname": "postgres","database.server.name": "dbserver1","schema.whitelist": "inventory","pulsar.service.url": "pulsar://127.0.0.1:6650"}'

    或:

    $ bin/pulsar-admin source localrun \ --source-config-file debezium-postgres-source-config.yaml
  4. 订阅变更 topic:

    $ bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0
  5. 进入 PostgreSQL 客户端:

    $ docker exec -it pulsar-postgresql /bin/bash
  6. 执行数据变更:

    psql -U postgres postgres postgres=# \c postgres; You are now connected to database "postgres" as user "postgres". postgres=# SET search_path TO inventory; SET postgres=# select * from products; id | name | description | weight -----+--------------------+---------------------------------------------------------+-------- 102 | car battery | 12V car battery | 8.1 103 | 12-pack drill bits | 12-pack of drill bits with sizes ranging from #40 to #3 | 0.8 104 | hammer | 12oz carpenter's hammer | 0.75 105 | hammer | 14oz carpenter's hammer | 0.875 106 | hammer | 16oz carpenter's hammer | 1 107 | rocks | box of assorted rocks | 5.3 108 | jacket | water resistent black wind breaker | 0.1 109 | spare tire | 24 inch spare tire | 22.2 101 | 1111111111 | Small 2-wheel scooter | 3.14 (9 rows) postgres=# UPDATE products SET name='1111111111' WHERE id=107; UPDATE 1

    订阅终端将收到形如下方的 Debezium 变更事件(JSON 结构完整保留 before/after 镜像与 source 元信息):

    ----- got message ----- {"schema":{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"}],"optional":false,"name":"dbserver1.inventory.products.Key"},"payload":{"id":107}}...{"schema":{"type":"struct","fields":[{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":false,"field":"name"},{"type":"string","optional":true,"field":"description"},{"type":"double","optional":true,"field":"weight"}],"optional":true,"name":"dbserver1.inventory.products.Value","field":"before"},{"type":"struct","fields":[{"type":"int32","optional":false,"field":"id"},{"type":"string","optional":false,"field":"name"},{"type":"string","optional":true,"field":"description"},{"type":"double","optional":true,"field":"weight"}],"optional":true,"name":"dbserver1.inventory.products.Value","field":"after"},{"type":"struct","fields":[{"type":"string","optional":true,"field":"version"},{"type":"string","optional":true,"field":"connector"},{"type":"string","optional":false,"field":"name"},{"type":"string","optional":false,"field":"db"},{"type":"int64","optional":true,"field":"ts_usec"},{"type":"int64","optional":true,"field":"txId"},{"type":"int64","optional":true,"field":"lsn"},{"type":"string","optional":true,"field":"schema"},{"type":"string","optional":true,"field":"table"},{"type":"boolean","optional":true,"default":false,"field":"snapshot"},{"type":"boolean","optional":true,"field":"last_snapshot_record"}],"optional":false,"name":"io.debezium.connector.postgresql.Source","field":"source"},{"type":"string","optional":false,"field":"op"},{"type":"int64","optional":true,"field":"ts_ms"}],"optional":false,"name":"dbserver1.inventory.products.Envelope"},"payload":{"before":{"id":107,"name":"rocks","description":"box of assorted rocks","weight":5.3},"after":{"id":107,"name":"1111111111","description":"box of assorted rocks","weight":5.3},"source":{"version":"0.9.2.Final","connector":"postgresql","name":"dbserver1","db":"postgres","ts_usec":1559208957661080,"txId":577,"lsn":23862872,"schema":"inventory","table":"products","snapshot":false,"last_snapshot_record":null},"op":"u","ts_ms":1559208957692}}

3.5 MongoDB 实战示例

配置文件

JSON 形态:

{ "mongodb.hosts": "rs0/mongodb:27017", "mongodb.name": "dbserver1", "mongodb.user": "debezium", "mongodb.password": "dbz", "mongodb.task.id": "1", "database.whitelist": "inventory", "pulsar.service.url": "pulsar://127.0.0.1:6650" }

YAML 形态(对应 debezium-mongodb-source-config.yaml):

tenant: "public" namespace: "default" name: "debezium-mongodb-source" topicName: "debezium-mongodb-topic" archive: "connectors/pulsar-io-debezium-mongodb-@pulsar:version@.nar" parallelism: 1 configs: ## config for pg, docker image: debezium/example-postgress:0.10 mongodb.hosts: "rs0/mongodb:27017", mongodb.name: "dbserver1", mongodb.user: "debezium", mongodb.password: "dbz", mongodb.task.id: "1", database.whitelist: "inventory", ## PULSAR_SERVICE_URL_CONFIG pulsar.service.url: "pulsar://127.0.0.1:6650"
启动步骤
  1. 启动 MongoDB 容器并初始化数据:

    $ docker pull debezium/example-mongodb:0.10 $ docker run -d -it --rm --name pulsar-mongodb -e MONGODB_USER=mongodb -e MONGODB_PASSWORD=mongodb -p 27017:27017 debezium/example-mongodb:0.10

    进入容器初始化示例集合:

    ./usr/local/bin/init-inventory.sh

    若本机无法访问容器网络,可编辑/etc/hosts添加规则127.0.0.1 <容器ID>(容器 ID 通过docker ps -a查看)。

  2. 启动 Pulsar standalone:bin/pulsar standalone。

  3. 启动连接器:

    $ bin/pulsar-admin source localrun \ --archive connectors/pulsar-io-debezium-mongodb-@pulsar:version@.nar \ --name debezium-mongodb-source \ --destination-topic-name debezium-mongodb-topic \ --tenant public \ --namespace default \ --source-config '{"mongodb.hosts": "rs0/mongodb:27017","mongodb.name": "dbserver1","mongodb.user": "debezium","mongodb.password": "dbz","mongodb.task.id": "1","database.whitelist": "inventory","pulsar.service.url": "pulsar://127.0.0.1:6650"}'

    或:

    $ bin/pulsar-admin source localrun \ --source-config-file debezium-mongodb-source-config.yaml
  4. 订阅 topic:

    $ bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0
  5. 进入 MongoDB 客户端并更新文档:

    $ docker exec -it pulsar-mongodb /bin/bash
    mongo -u debezium -p dbz --authenticationDatabase admin localhost:27017/inventory db.products.update({"_id":NumberLong(104)},{$set:{weight:1.25}})

    订阅终端会收到类似下方的变更事件(after字段以 JSON 字符串承载 MongoDB 文档快照,source记录副本集、集合与操作序号):

    ----- got message ----- {"schema":{"type":"struct","fields":[{"type":"string","optional":false,"field":"id"}],"optional":false,"name":"dbserver1.inventory.products.Key"},"payload":{"id":"104"}}, value = {"schema":{"type":"struct","fields":[{"type":"string","optional":true,"name":"io.debezium.data.Json","version":1,"field":"after"},{"type":"string","optional":true,"name":"io.debezium.data.Json","version":1,"field":"patch"},{"type":"struct","fields":[{"type":"string","optional":false,"field":"version"},{"type":"string","optional":false,"field":"connector"},{"type":"string","optional":false,"field":"name"},{"type":"int64","optional":false,"field":"ts_ms"},{"type":"string","optional":true,"name":"io.debezium.data.Enum","version":1,"parameters":{"allowed":"true,last,false"},"default":"false","field":"snapshot"},{"type":"string","optional":false,"field":"db"},{"type":"string","optional":false,"field":"rs"},{"type":"string","optional":false,"field":"collection"},{"type":"int32","optional":false,"field":"ord"},{"type":"int64","optional":true,"field":"h"}],"optional":false,"name":"io.debezium.connector.mongo.Source","field":"source"},{"type":"string","optional":true,"field":"op"},{"type":"int64","optional":true,"field":"ts_ms"}],"optional":false,"name":"dbserver1.inventory.products.Envelope"},"payload":{"after":"{\"_id\": {\"$numberLong\": \"104\"},\"name\": \"hammer\",\"description\": \"12oz carpenter's hammer\",\"weight\": 1.25,\"quantity\": 4}","patch":null,"source":{"version":"0.10.0.Final","connector":"mongodb","name":"dbserver1","ts_ms":1573541905000,"snapshot":"true","db":"inventory","rs":"rs0","collection":"products","ord":1,"h":4983083486544392763},"op":"r","ts_ms":1573541909761}}.

四、FAQ:Postgres 连接器创建快照时挂起

在使用 Debezium Postgres 连接器同步数据时,若连接器线程长时间处于WAITING (parking)状态,典型的线程栈如下:

java.lang.Thread.State: WAITING (parking) at sun.misc.Unsafe.park(Native Method) at java.util.concurrent.locks.LockSupport.park(LockSupport.java:175) at java.util.concurrent.locks.AbstractQueuedSynchronizer$ConditionObject.await(...) at java.util.concurrent.LinkedBlockingDeque.putLast(LinkedBlockingDeque.java:396) at java.util.concurrent.LinkedBlockingDeque.put(LinkedBlockingDeque.java:649) at io.debezium.connector.base.ChangeEventQueue.enqueue(ChangeEventQueue.java:132) at io.debezium.connector.postgresql.PostgresConnectorTask$Lambda$203/385424085.accept(Unknown Source) at io.debezium.connector.postgresql.RecordsSnapshotProducer.sendCurrentRecord(RecordsSnapshotProducer.java:402) at io.debezium.connector.postgresql.RecordsSnapshotProducer.readTable(RecordsSnapshotProducer.java:321) at io.debezium.connector.postgresql.RecordsSnapshotProducer.lambda$takeSnapshot$6(RecordsSnapshotProducer.java:226) ... at io.debezium.connector.postgresql.PostgresConnectorTask.start(PostgresConnectorTask.java:126) at io.debezium.connector.common.BaseSourceTask.start(BaseSourceTask.java:47) at org.apache.pulsar.io.kafka.connect.KafkaConnectSource.open(KafkaConnectSource.java:127) at org.apache.pulsar.io.debezium.DebeziumSource.open(DebeziumSource.java:100)

该问题的根因是快照阶段变更事件队列(ChangeEventQueue)被写满后阻塞了生产者线程。解决方案是在配置文件中补充以下配置项,调大事件队列容量:

max.queue.size=

(具体数值依据表数据量设置,参考 Pulsar issue 4075 提供的 Kafka Connect 兼容层之上。

五、总结与选型建议

  • Canal source connector(pulsar-io/canal)面向 MySQL binlog 场景,基于 Canal 的getWithoutAck/ack语义实现拉取与确认,输出 JSON(CanalStringSource)或字节(CanalByteSource),适合已有 Canal 基础设施、希望最小化引入新组件的团队。
  • Debezium source connector(pulsar-io/debezium)覆盖面更广,仓库中已包含 MySQL、PostgreSQL、MongoDB、MSSQL、Oracle 五个分支;其 schema 历史(PulsarDatabaseHistory)与偏移量均存储于 Pulsar topic,天然适配 Pulsar 作为 CDC 事件中枢的架构,也便于接入下游实时数仓与事件驱动应用。

两条链路产出的变更事件都遵循"数据库逻辑名.库名.表名"的 topic 命名规则,消费者可按public/default/<server.name>.<db>.<table>订阅对应表,并利用 Pulsar 的持久化、分区与多租户能力构建高可用的实时数据管道。进一步阅读:CDC 连接器总览、Debezium source connector、CDC Canal Connector。

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

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询