- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
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 connector | org.apache.pulsar.io.canal.CanalStringSource |
| Debezium source connector | org.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()进入核心拉取循环:
connector.connect()建立连接,connector.subscribe()订阅目标;connector.getWithoutAck(batchSize)拉取一批 binlog 消息;- 通过 MessageUtils.messageConverter 将 protobuf 原消息转换为
FlatMessage(扁平化后的列结构,包含isKey、isNull、mysqlType、columnName、columnValue等字段); - 封装为
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 | 否 | false | true时通过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 | 是 | null | Debezium 实现的具体 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 | 是 | null | Kafka Connect 提供的记录 key 转换器 |
value.converter | 是 | null | Kafka Connect 提供的记录 value 转换器 |
database.history | 是 | null | 数据库历史类名 |
database.history.pulsar.topic | 是 | null | 连接器写入并恢复 DDL 语句的数据库历史 topic。注意:该 topic 仅供内部使用,消费者不应使用 |
database.history.pulsar.service.url | 是 | null | 历史 topic 使用的 Pulsar 集群服务地址 |
pulsar.service.url | 是 | null | Debezium 偏移量 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 | 是 | null | MongoDB 副本集主机端口逗号分隔列表(host或host:port形式) |
mongodb.name | 是 | null | 标识连接器及其监控的 MongoDB 副本集/共享集群的唯一名称,每个服务器至多由一个 Debezium 连接器监控 |
mongodb.user | 是 | null | 连接 MongoDB 的数据库用户名(仅在启用认证时需要) |
mongodb.password | 是 | null | 连接 MongoDB 的密码(仅在启用认证时需要) |
mongodb.task.id | 是 | null | MongoDB 连接器 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"启动步骤
启动带示例数据库的 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本地以 standalone 模式启动 Pulsar:
$ bin/pulsar standalone以 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
订阅
inventory.products表对应的变更 topic:$ bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0注意 topic 名由
database.server.name+ 库名 + 表名构成,即dbserver1.inventory.products。在 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"'客户端弹出后,执行变更语句:
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"启动步骤
启动 PostgreSQL 容器:
$ docker pull debezium/example-postgres:0.8 $ docker run -d -it --rm --name pulsar-postgresql -p 5432:5432 debezium/example-postgres:0.8启动 Pulsar standalone:
bin/pulsar standalone。启动连接器(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订阅变更 topic:
$ bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0进入 PostgreSQL 客户端:
$ docker exec -it pulsar-postgresql /bin/bash执行数据变更:
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"启动步骤
启动 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查看)。启动 Pulsar standalone:
bin/pulsar standalone。启动连接器:
$ 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订阅 topic:
$ bin/pulsar-client consume -s "sub-products" public/default/dbserver1.inventory.products -n 0进入 MongoDB 客户端并更新文档:
$ docker exec -it pulsar-mongodb /bin/bashmongo -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
相关推荐
Apache Pulsar CDC 连接器实战指南:Canal 与 Debezium Source Connector 捕获数据库变更入 Pulsar
Apache Pulsar CDC 连接器实战指南:Canal 与 Debezium Source Connector 捕获数据库变更入 Pulsar CDC(
消息队列后端流处理Apache Pulsar CDC Connector:用 Debezium 与 Canal 将数据库变更日志接入 Pulsar
Apache Pulsar CDC Connector:用 Debezium 与 Canal 将数据库变更日志接入 Pulsar Pulsar 的 CDC(Ch
消息队列后端流处理Apache Pulsar Debezium Source Connector 实战指南:MySQL/PostgreSQL/MongoDB 变更数据捕获(CDC)接入
Apache Pulsar Debezium Source Connector 实战指南:MySQL/PostgreSQL/MongoDB 变更数据捕获(CDC
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考