概述
Debezium MySQL Source Connector 从 MySQL 读取表结构、存量行和 binlog 中的行级及 DDL 变更,并将它们作为 Kafka Connect SourceRecord 发布到 Kafka。首次运行通常先执行快照,再从快照边界继续增量读取;每张表的变更默认进入按 Topic 前缀、数据库名和表名组织的 Topic。它适用于把 MySQL 数据库变更事件持续送入 Kafka,供下游实时计算、缓存同步、审计或数据湖处理。前置条件
- MySQL 必须启用二进制日志,并使用
binlog_format=ROW和binlog_row_image=FULL;日志保留时间应覆盖预期的故障恢复和同步延迟。 - 连接用户需要对快照表具备
SELECT权限,并具备复制和位置读取所需的REPLICATION SLAVE、REPLICATION CLIENT权限;发现数据库和快照锁定可能还需要SHOW DATABASES、RELOAD,表锁回退还需要LOCK TABLES。托管 MySQL 可能使用不同的授权名称。 database.server.id必须在当前运行的 MySQL 客户端中唯一;如果使用 GTID 或只读增量快照,源服务器还必须保留所需的 GTID 历史。- Kafka Schema History 使用专用的单分区 Topic,保留完整 DDL 历史,并允许 Connector 使用对应的 Kafka 客户端配置读写和创建该 Topic。
授权许可
使用 Apache License 2.0。快速开始
准备 Connect Cluster、Kafka 和要捕获的 MySQL 数据库,确认 Connector 能访问数据库、Kafka 以及 Schema History Topic,并完成必要授权。Connector 的创建和管理方式请参阅管理 Connector。 下面的配置会对已有表执行初始快照,然后持续读取 binlog;请替换尖括号中的环境值。topic.prefix、Schema History Topic 和 database.server.id 应在多个 Connector 间分别保持唯一;密码应通过部署环境支持的安全配置方式提供。
配置
标识
topic.prefix
为该 MySQL 服务器/集群建立 Topic 命名空间,必须在不同 Connector 之间保持唯一。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:仅字母、数字、连字符、点和下划线;必填
数据库连接
database.hostname
MySQL 服务器的可解析主机名或 IP 地址。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:^[a-zA-Z0-9-_.]+$;必填
database.port
MySQL 服务器端口。
- 类型:
INT - 默认值:
3306 - 重要级别:高
- 有效值 / 注意事项:整数
database.user
连接 MySQL 使用的用户名。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填
database.password
连接 MySQL 使用的密码。
- 类型:
PASSWORD - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:按类型填写;未列出的值没有额外约束。
database.query.timeout.ms
等待数据库查询完成的超时时间,单位毫秒。
- 类型:
INT - 默认值:
600000 - 重要级别:低
- 有效值 / 注意事项:整数;0 表示不限制
database.initial.statements
每次建立 JDBC 连接时执行的 SQL;适合设置会话参数,不应执行数据修改。
- 类型:
STRING - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:以分号分隔 SQL;字面分号使用 ;;
database.server.id
Connector 作为 binlog 客户端加入 MySQL 集群时使用的唯一数字 ID。
- 类型:
LONG - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:正长整数;在运行中的数据库客户端之间唯一;必填
database.server.id.offset
并行快照连接使用的 server ID 偏移量。
- 类型:
LONG - 默认值:
10000 - 重要级别:高
- 有效值 / 注意事项:长整数;依赖:
snapshot.max.threads
SSL/TLS
database.ssl.keystore
客户端密钥库文件路径。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:依赖:
database.ssl.keystore.password
database.ssl.keystore.password
客户端密钥库的密码。
- 类型:
PASSWORD - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:依赖:
database.ssl.keystore
database.ssl.truststore
用于验证服务器证书的信任库文件路径。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:依赖:
database.ssl.truststore.password
database.ssl.truststore.password
验证信任库完整性所需的密码。
- 类型:
PASSWORD - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:依赖:
database.ssl.truststore
数据库连接
database.jdbc.driver
连接 MySQL 的 JDBC 驱动类。
- 类型:
CLASS - 默认值:
com.mysql.cj.jdbc.Driver - 重要级别:低
- 有效值 / 注意事项:有效类名
database.protocol
JDBC 驱动使用的连接协议。
- 类型:
STRING - 默认值:
jdbc:mysql - 重要级别:低
- 有效值 / 注意事项:依赖:
database.jdbc.driver
SSL/TLS
database.ssl.mode
控制数据库连接是否加密,以及是否校验证书和主机身份。
- 类型:
STRING - 默认值:
preferred - 重要级别:中
- 有效值 / 注意事项:disabled, preferred, required, verify_ca, verify_identity;依赖:
database.ssl.keystore,database.ssl.truststore
错误处理
event.processing.failure.handling.mode
控制处理损坏或无法解析事件时是停止、记录并跳过,还是直接跳过。
- 类型:
STRING - 默认值:
fail - 重要级别:中
- 有效值 / 注意事项:fail, warn, ignore, skip
吞吐与缓冲
max.batch.size
每次向 Kafka Connect 返回的 SourceRecord 最大数量。
- 类型:
INT - 默认值:
2048 - 重要级别:中
- 有效值 / 注意事项:正整数;依赖:
max.queue.size
max.queue.size
已读取但尚未转发的变更事件队列最大数量,必须大于 max.batch.size。
- 类型:
INT - 默认值:
8192 - 重要级别:中
- 有效值 / 注意事项:正数;必须大于
max.batch.size;依赖:max.batch.size
poll.interval.ms
没有新事件时等待下一次轮询的时间。
- 类型:
LONG - 默认值:
500 - 重要级别:中
- 有效值 / 注意事项:正整数
max.queue.size.in.bytes
变更事件队列的字节上限;与记录数上限共同限制缓冲区。
- 类型:
LONG - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:非负数;0 表示不限制字节数
事件
provide.transaction.metadata
是否生成事务元数据及事件计数。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:
true或false;依赖:transaction.metadata.factory
skipped.operations
指定增量同步阶段跳过的操作类型。
- 类型:
LIST - 默认值:
t - 重要级别:低
- 有效值 / 注意事项:r,c,u,d,t,none;none 不能与其他操作同时配置
快照
snapshot.delay.ms
开始快照前的等待时间,单位毫秒。
- 类型:
LONG - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:非负长整数
streaming.delay.ms
快照结束后开始增量读取前的等待时间,单位毫秒。
- 类型:
LONG - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:非负长整数
snapshot.include.collection.list
通过表名正则表达式限定快照读取范围。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:用于选择快照表的正则表达式列表;依赖:
snapshot.mode
snapshot.fetch.size
快照查询每批从数据库获取的行数。
- 类型:
INT - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:非负整数
snapshot.max.threads
快照使用的最大线程数;增加线程不会增加 Kafka Connect Task 数量。
- 类型:
INT - 默认值:
1 - 重要级别:中
- 有效值 / 注意事项:正整数;依赖:
database.server.id.offset
snapshot.mode.custom.name
自定义快照模式实现的名称。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:snapshot.mode=custom 时不能为空;依赖:
snapshot.mode
snapshot.mode.configuration.based.snapshot.data
在 configuration_based 模式下决定是否快照行数据。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;依赖:snapshot.mode=configuration_based
snapshot.mode.configuration.based.snapshot.schema
在 configuration_based 模式下决定是否快照表结构。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;依赖:snapshot.mode=configuration_based
snapshot.mode.configuration.based.start.stream
在 configuration_based 模式下决定是否开始增量读取。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;依赖:snapshot.mode=configuration_based
snapshot.mode.configuration.based.snapshot.on.schema.error
在 configuration_based 模式下决定遇到结构历史错误时是否请求快照。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;依赖:snapshot.mode=configuration_based
snapshot.mode.configuration.based.snapshot.on.data.error
在 configuration_based 模式下决定日志位置不可用时是否请求快照。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;依赖:snapshot.mode=configuration_based
错误处理
retriable.restart.connector.wait.ms
可重试错误触发重启后的等待时间,单位毫秒;不表示所有数据库错误都可重试。
- 类型:
LONG - 默认值:
10000 - 重要级别:低
- 有效值 / 注意事项:正长整数
吞吐与缓冲
query.fetch.size
增量读取过程中 JDBC 查询的批量取行大小。
- 类型:
INT - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:非负数;0 表示使用 JDBC 默认值
错误处理
errors.max.retries
可重试错误的最大重试次数。
- 类型:
INT - 默认值:
-1 - 重要级别:低
- 有效值 / 注意事项:-1 表示无限重试,0 表示不重试,正数表示重试次数
增量快照
incremental.snapshot.watermarking.strategy
增量快照窗口开启和关闭时写入信号表的水位标记策略。
- 类型:
STRING - 默认值:
INSERT_INSERT - 重要级别:低
- 有效值 / 注意事项:INSERT_INSERT, INSERT_DELETE;依赖:
signal.data.collection
数据库连接
connection.validation.timeout.ms
验证数据库连接的超时时间,单位毫秒。
- 类型:
LONG - 默认值:
60000 - 重要级别:低
- 有效值 / 注意事项:正长整数
运行时
executor.shutdown.timeout.ms
等待执行器关闭的超时时间,单位毫秒。
- 类型:
LONG - 默认值:
4000 - 重要级别:中
- 有效值 / 注意事项:正长整数
集成
openlineage.integration.enabled
是否启用 OpenLineage 集成。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:
true或false;依赖:openlineage.integration.config.file.path
openlineage.integration.config.file.path
OpenLineage 配置文件路径。
- 类型:
STRING - 默认值:
./openlineage.yml - 重要级别:低
- 有效值 / 注意事项:依赖:
openlineage.integration.enabled
openlineage.integration.job.namespace
OpenLineage 作业的命名空间。
- 类型:
STRING - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:依赖:
openlineage.integration.enabled
openlineage.integration.job.description
OpenLineage 作业的描述。
- 类型:
STRING - 默认值:
Debezium change data capture job - 重要级别:低
- 有效值 / 注意事项:依赖:
openlineage.integration.enabled
openlineage.integration.job.tags
OpenLineage 作业的标签。
- 类型:
LIST - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:key=value 形式的列表;依赖:
openlineage.integration.enabled
openlineage.integration.job.owners
OpenLineage 作业的所有者信息。
- 类型:
LIST - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:key=value 形式的列表;依赖:
openlineage.integration.enabled
事件
extended.headers.enabled
是否添加扩展事件头。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:
true或false
类型转换
decimal.handling.mode
控制 DECIMAL/NUMERIC 值在事件中的表示方式。
- 类型:
STRING - 默认值:
precise - 重要级别:中
- 有效值 / 注意事项:precise, string, double
快照
snapshot.lock.timeout.ms
快照获取锁的超时时间,单位毫秒。
- 类型:
LONG - 默认值:
10000 - 重要级别:中
- 有效值 / 注意事项:长整数;获取锁的超时时间;依赖:
snapshot.locking.mode
数据库连接
connect.timeout.ms
建立数据库连接的超时时间,单位毫秒。
- 类型:
INT - 默认值:
30000 - 重要级别:中
- 有效值 / 注意事项:正整数
connect.keep.alive
是否启用 binlog 连接保活线程。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:
true或false;依赖:connect.keep.alive.interval.ms
connect.keep.alive.interval.ms
binlog 连接保活检查间隔,单位毫秒。
- 类型:
LONG - 默认值:
60000 - 重要级别:低
- 有效值 / 注意事项:正整数;依赖:
connect.keep.alive
use.nongraceful.disconnect
是否采用非正常关闭方式断开 binlog 客户端连接。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false
快照
snapshot.mode
决定何时执行快照,以及快照后是否继续增量读取。
- 类型:
STRING - 默认值:
initial - 重要级别:低
- 有效值 / 注意事项:always, when_needed, initial, initial_only, never, configuration_based, custom, no_data, recovery, schema_only, schema_only_recovery;recovery 的安全前提为已有 Offset 且自上次 Connector 停机后没有结构变化;停机后有结构变化时不得使用,不以兼容 DDL 判断放宽。所需 binlog 仍须可用,重建结构历史不能补回缺失日志;when_needed 不处理 Schema History 错误,也不恢复已丢失的历史事件。依赖:
snapshot.mode.custom.name,snapshot.mode.configuration.based.snapshot.data,snapshot.mode.configuration.based.snapshot.schema,snapshot.mode.configuration.based.start.stream
snapshot.query.mode
选择快照查询实现。
- 类型:
STRING - 默认值:
select_all - 重要级别:低
- 有效值 / 注意事项:select_all, custom;依赖:
snapshot.query.mode.custom.name
snapshot.query.mode.custom.name
自定义快照查询实现的名称。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:snapshot.query.mode=custom 时不能为空;依赖:
snapshot.query.mode
类型转换
bigint.unsigned.handling.mode
控制 BIGINT UNSIGNED 超出有符号范围时的表示方式。
- 类型:
STRING - 默认值:
long - 重要级别:中
- 有效值 / 注意事项:long, precise
time.precision.mode
控制时间、日期和时间戳字段的精度表示。
- 类型:
STRING - 默认值:
adaptive_time_microseconds - 重要级别:中
- 有效值 / 注意事项:
adaptive_time_microseconds、connect;不接受adaptive
enable.time.adjuster
是否调整两位年份的时间值。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:
true或false
schema.name.adjustment.mode
调整 Schema 名称以适配序列化格式的命名规则。
- 类型:
STRING - 默认值:
none - 重要级别:低
- 有效值 / 注意事项:avro, avro_unicode, none
快照
min.row.count.to.stream.results
快照查询改用流式结果读取的表行数阈值。
- 类型:
INT - 默认值:
1000 - 重要级别:低
- 有效值 / 注意事项:非负数;0 表示对所有查询结果使用流式读取
增量快照
incremental.snapshot.chunk.size
增量快照每个分块读取的记录数量。
- 类型:
INT - 默认值:
1024 - 重要级别:中
- 有效值 / 注意事项:非负整数;依赖:
signal.data.collection
incremental.snapshot.allow.schema.changes
增量快照期间是否允许 Schema 变更。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:
true或false;不支持主键变化;依赖:incremental.snapshot.chunk.size
快照
snapshot.locking.mode
控制快照读取期间使用的锁策略。
- 类型:
STRING - 默认值:
minimal - 重要级别:低
- 有效值 / 注意事项:
extended、minimal、minimal_percona、minimal_percona_no_table_locks、none或custom。minimal在读取结构和元数据时持有全局读锁,随后依赖 REPEATABLE READ 和兼容的事务存储引擎读取行;不保证任意存储引擎或并发 DDL 下的一致性。none不获取快照锁,使用边界为snapshot.mode=schema_only或schema_only_recovery,且必须通过外部机制阻止整个快照期间的并发 DDL;枚举值校验通过不表示其他组合安全。custom还需指定snapshot.locking.mode.custom.name。
内部 Schema History
schema.history.internal
Schema History 实现类,默认使用 Kafka 存储实现。
- 类型:
CLASS - 默认值:
io.debezium.storage.kafka.history.KafkaSchemaHistory - 重要级别:低
- 有效值 / 注意事项:SchemaHistory 实现类;支持 Kafka 和文件型存储。依赖:
schema.history.internal.kafka.topic、schema.history.internal.kafka.bootstrap.servers。
schema.history.internal.skip.unparseable.ddl
遇到无法解析的 DDL 时是否跳过该 DDL;跳过可能导致 Schema 元数据不完整。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:
true或false;跳过 DDL 可能导致后续事件缺少正确的字段定义;依赖:schema.history.internal
schema.history.internal.store.only.captured.tables.ddl
是否只在内部 Schema History 中保存已捕获表的 DDL。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:
true或false;扩大表过滤范围前应确认所需历史 Schema 仍可获得;依赖:schema.history.internal、table.include.list、table.exclude.list
schema.history.internal.store.only.captured.databases.ddl
是否只在内部 Schema History 中保存已捕获数据库的 DDL。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:
true或false;扩大数据库过滤范围前应确认所需历史 Schema 仍可获得;依赖:schema.history.internal、database.include.list、database.exclude.list
扩展
converters
自定义类型转换器的别名列表。
- 类型:
STRING - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:逗号分隔的别名;每个别名的 .type 指定 CustomConverter 实现类
post.processors
事件后处理器的别名列表。
- 类型:
STRING - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:逗号分隔的别名;每个别名的 .type 指定后处理器实现类
事件
tombstones.on.delete
删除事件后是否发送同键的 null 值 tombstone。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:中
- 有效值 / 注意事项:
true或false
心跳
heartbeat.interval.ms
心跳事件的发送间隔;0 表示不发送心跳。
- 类型:
INT - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:非负数;0 disables;依赖:
heartbeat.topics.prefix
heartbeat.topics.prefix
心跳 Topic 的名称前缀。
- 类型:
STRING - 默认值:
__debezium-heartbeat - 重要级别:低
- 有效值 / 注意事项:依赖:
heartbeat.interval.ms
信号
signal.data.collection
接收增量快照或其他信号的 MySQL 表。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:完整的 database.table 名称;依赖:
signal.enabled.channels
signal.poll.interval.ms
检查新信号的间隔,单位毫秒。
- 类型:
LONG - 默认值:
5000 - 重要级别:中
- 有效值 / 注意事项:正整数;依赖:
signal.enabled.channels
signal.enabled.channels
启用的信号通道列表。
- 类型:
LIST - 默认值:
source - 重要级别:中
- 有效值 / 注意事项:通道名称;默认启用 source;依赖:
signal.data.collection
Topic 命名
topic.naming.strategy
生成数据、Schema 变更及辅助 Topic 名称的策略类。
- 类型:
CLASS - 默认值:
io.debezium.schema.SchemaTopicNamingStrategy - 重要级别:中
- 有效值 / 注意事项:TopicNamingStrategy 实现类
通知
notification.enabled.channels
启用的通知通道列表。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:通知通道名称;依赖:
notification.sink.topic.name
notification.sink.topic.name
sink 通知通道发布通知的 Kafka Topic。
- 类型:
STRING - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:notification.enabled.channels 包含 sink 时必填;依赖:
notification.enabled.channels
事件
transaction.metadata.factory
生成事务元数据的工厂实现类。
- 类型:
CLASS - 默认值:
io.debezium.pipeline.txmetadata.DefaultTransactionMetadataFactory - 重要级别:低
- 有效值 / 注意事项:TransactionMetadataFactory 实现类;依赖:
provide.transaction.metadata
监控
custom.metric.tags
添加到监控指标上的自定义标签。
- 类型:
LIST - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:key=value 形式的列表
过滤
column.include.list
只保留匹配列的值;值为逗号分隔的正则表达式。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:正则表达式列表;不能同时配置 column.exclude.list;依赖:
column.exclude.list
column.exclude.list
排除匹配列的值;不能与 column.include.list 同时配置。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:正则表达式列表;不能同时配置 column.include.list;依赖:
column.include.list
table.include.list
只捕获匹配的数据库表;值为逗号分隔的正则表达式。
- 类型:
LIST - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:正则表达式列表;不能同时配置 table.exclude.list;依赖:
table.exclude.list,database.include.list
table.exclude.list
排除匹配的数据库表;不能与 table.include.list 同时配置。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:正则表达式列表;不能同时配置 table.include.list;依赖:
table.include.list
键映射
message.key.columns
为指定表定义消息键使用的列。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:分号分隔的 table:key 表达式,须符合消息键表达式格式
快照
snapshot.select.statement.overrides
使用自定义 SELECT 查询的表名列表。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:逗号分隔的完整表名;每张表对应
snapshot.select.statement.overrides.<DB>.<TABLE>查询配置;依赖:snapshot.mode
字段脱敏
column.mask.hash.([^.]+).with.salt.(.+)
使用指定算法和盐值对匹配列进行哈希脱敏。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:动态配置名;指定列正则表达式、哈希算法和盐值
column.mask.with.(d+).chars
将匹配列的值替换为指定长度的掩码。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:动态配置名;指定整数掩码长度和列正则表达式列表
column.truncate.to.(d+).chars
将匹配列的值截断到指定字符数。
- 类型:
INT - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:动态配置名;指定整数截断长度
Schema 事件
include.schema.changes
是否发布 Schema 变更事件。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:中
- 有效值 / 注意事项:
true或false
include.schema.comments
是否在 Schema 中包含数据库注释;可能增加内存占用。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;可能增加内存占用
column.propagate.source.type
为匹配列附加源数据库类型信息。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:列名正则表达式列表
datatype.propagate.source.type
为匹配数据库类型的列附加源类型信息。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:数据库类型正则表达式列表
快照
snapshot.tables.order.by.row.count
按行数决定快照表的读取顺序。
- 类型:
STRING - 默认值:
disabled - 重要级别:中
- 有效值 / 注意事项:ascending, descending, disabled
心跳
heartbeat.action.query
发送心跳时执行的数据库查询。
- 类型:
STRING - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:仅在心跳间隔非零时执行;依赖:
heartbeat.interval.ms
事件
include.query
是否在事件中包含原始 SQL;需要 MySQL 记录行事件查询,且可能暴露被过滤的数据。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;需要 MySQL 设置 binlog_rows_query_log_events=ON;可能暴露被过滤的数据
过滤
table.ignore.builtin
是否忽略数据库内置表。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:
true或false;依赖:database.include.list
database.include.list
只捕获匹配的数据库;值为逗号分隔的正则表达式。
- 类型:
LIST - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:正则表达式列表;不能同时配置 database.exclude.list;依赖:
database.exclude.list,table.include.list
database.exclude.list
排除匹配的数据库;不能与 database.include.list 同时配置。
- 类型:
LIST - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:正则表达式列表;不能同时配置 database.include.list;依赖:
database.include.list
吞吐与缓冲
binlog.buffer.size
binlog 前瞻缓冲区大小,用于识别事务提交或回滚。
- 类型:
INT - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:非负数;0 表示不启用前瞻缓冲
错误处理
event.deserialization.failure.handling.mode
事件反序列化失败时的处理方式;此配置已弃用。
- 类型:
STRING - 默认值:
fail - 重要级别:中
- 有效值 / 注意事项:fail, warn, ignore, skip;设置此项会记录弃用警告;已弃用;替代项为
event.processing.failure.handling.mode
inconsistent.schema.handling.mode
处理缺少对应表结构的事件时采用的策略。
- 类型:
STRING - 默认值:
fail - 重要级别:中
- 有效值 / 注意事项:fail, warn, skip
GTID
gtid.source.filter.dml.events
是否将 GTID 来源过滤应用到行级变更事件。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:中
- 有效值 / 注意事项:
true或false;依赖:gtid.source.includes,gtid.source.excludes
gtid.source.includes
限定 GTID 来源 UUID 范围;不能与 gtid.source.excludes 同时配置。
- 类型:
LIST - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:GTID 来源 UUID 表达式;不能同时配置 gtid.source.excludes;依赖:
gtid.source.excludes
gtid.source.excludes
排除指定 GTID 来源 UUID;不能与 gtid.source.includes 同时配置。
- 类型:
STRING - 默认值:无
- 重要级别:中
- 有效值 / 注意事项:GTID 来源 UUID 表达式;设置 includes 时不接受此项;依赖:
gtid.source.includes
元数据
sourceinfo.struct.maker
生成事件 Source 元数据结构的实现类。
- 类型:
CLASS - 默认值:
io.debezium.connector.mysql.MySqlSourceInfoStructMaker - 重要级别:低
- 有效值 / 注意事项:SourceInfoStructMaker 实现类
Schema History 后端
schema.history.internal.name
Schema History 后端的逻辑名称。
- 类型:
STRING - 默认值:未声明默认值;运行时注入
<logical-name>-schemahistory - 重要级别:低
- 有效值 / 注意事项:后端逻辑名称;依赖:
schema.history.internal
schema.history.internal.kafka.topic
Kafka Schema History 使用的 Topic;必须使用专用的单分区 Topic 并保留完整 DDL 历史。
- 类型:
STRING - 默认值:未声明
- 重要级别:高
- 有效值 / 注意事项:KafkaSchemaHistory 必填;必填;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
schema.history.internal.kafka.bootstrap.servers
读取和写入 Schema History 的 Kafka 集群地址,通常应与 Connect 使用的集群相同。
- 类型:
STRING - 默认值:未声明
- 重要级别:高
- 有效值 / 注意事项:KafkaSchemaHistory 必填;通常使用同一 Kafka 集群;必填;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
schema.history.internal.kafka.recovery.poll.interval.ms
恢复 Schema History 时轮询 Kafka 的间隔,单位毫秒。
- 类型:
INT - 默认值:
100 - 重要级别:低
- 有效值 / 注意事项:非负整数;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
schema.history.internal.kafka.recovery.attempts
恢复 Schema History 时允许的连续空轮询次数。
- 类型:
INT - 默认值:
100 - 重要级别:低
- 有效值 / 注意事项:整数;空轮询的总等待时间为次数乘以轮询间隔;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
schema.history.internal.kafka.query.timeout.ms
查询 Schema History Kafka 集群信息的超时时间,单位毫秒。
- 类型:
LONG - 默认值:
3000 - 重要级别:低
- 有效值 / 注意事项:正整数;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
schema.history.internal.kafka.create.timeout.ms
创建 Schema History Topic 的超时时间,单位毫秒。
- 类型:
LONG - 默认值:
30000 - 重要级别:低
- 有效值 / 注意事项:正整数;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
schema.history.internal.consumer.*
传递给 Schema History Kafka 消费者的客户端配置。
- 类型:
MAP - 默认值:未声明
- 重要级别:低
- 有效值 / 注意事项:去掉前缀后传递给 Kafka 消费者的配置;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
schema.history.internal.producer.*
传递给 Schema History Kafka 生产者的客户端配置。
- 类型:
MAP - 默认值:未声明
- 重要级别:低
- 有效值 / 注意事项:去掉前缀后传递给 Kafka 生产者的配置;依赖:
schema.history.internal=io.debezium.storage.kafka.history.KafkaSchemaHistory
信号后端
signal.file
file 信号通道读取的文件路径。
- 类型:
STRING - 默认值:未声明默认值;运行时回退为
file-signals.txt - 重要级别:高
- 有效值 / 注意事项:仅在启用 file 信号通道时读取;省略时使用 file-signals.txt。;依赖:
signal.enabled.channels
signal.kafka.topic
kafka 信号通道读取的 Topic。
- 类型:
STRING - 默认值:未声明默认值;运行时回退为
<connector logical name>-signal - 重要级别:高
- 有效值 / 注意事项:仅在启用 kafka 信号通道时读取;省略时使用
<connector logical name>-signal。;依赖:signal.enabled.channels
signal.kafka.bootstrap.servers
kafka 信号通道使用的 Kafka 集群地址。
- 类型:
STRING - 默认值:未声明
- 重要级别:高
- 有效值 / 注意事项:启用 kafka 信号通道时必填;必填;依赖:
signal.enabled.channels
signal.kafka.poll.timeout.ms
kafka 信号通道每次轮询的超时时间,单位毫秒。
- 类型:
INT - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:非负整数;依赖:
signal.enabled.channels
signal.kafka.groupId
kafka 信号通道使用的消费者组 ID。
- 类型:
STRING - 默认值:
kafka-signal - 重要级别:低
- 有效值 / 注意事项:依赖:
signal.enabled.channels
signal.consumer.*
传递给 Kafka 信号消费者的客户端配置。
- 类型:
MAP - 默认值:未声明
- 重要级别:低
- 有效值 / 注意事项:去掉 signal.consumer. 前缀后传递给 Kafka 消费者的配置;依赖:
signal.enabled.channels
Kafka Connect 框架
name
Connector 的全局唯一名称。
- 类型:
STRING - 默认值:未声明
- 重要级别:高
- 有效值 / 注意事项:非空的 Connector 名称;必填
connector.class
Connector 实现类;MySQL Source 使用 io.debezium.connector.mysql.MySqlConnector。
- 类型:
STRING - 默认值:未声明
- 重要级别:高
- 有效值 / 注意事项:Connector 实现类;使用 io.debezium.connector.mysql.MySqlConnector;必填
tasks.max
该 Connector 可创建的最大 Task 数;MySQL CDC 实际只处理一个 Task。
- 类型:
INT - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:框架接受至少为 1 的整数,但 MySQL Connector 在生成 Task 配置时拒绝大于 1 的值并抛错;只能使用
1。
tasks.max.enforce
是否强制 Task 数量不超过 tasks.max;关闭它不会绕过 MySQL Connector 的单 Task 限制。
- 类型:
BOOLEAN - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:
true或false;已弃用,当前没有替代配置。
key.converter
SourceRecord 键的 Converter 类;未在 Connector 级别设置时使用 Worker 配置。
- 类型:
CLASS - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:可实例化的 Converter 子类
value.converter
SourceRecord 值的 Converter 类;未在 Connector 级别设置时使用 Worker 配置。
- 类型:
CLASS - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:可实例化的 Converter 子类
header.converter
SourceRecord 事件头的 Converter 类。
- 类型:
CLASS - 默认值:无
- 重要级别:低
- 有效值 / 注意事项:可实例化的 HeaderConverter 子类
transforms
按应用顺序列出单消息转换的别名。
- 类型:
LIST - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:不重复的转换别名
predicates
供单消息转换使用的条件判断别名。
- 类型:
LIST - 默认值:空列表
- 重要级别:低
- 有效值 / 注意事项:不重复的条件判断别名;依赖:
transforms
config.action.reload
外部配置提供者的值变化时是否重启 Connector。
- 类型:
STRING - 默认值:
restart - 重要级别:低
- 有效值 / 注意事项:none, restart
errors.retry.timeout
框架处理可重试错误的总时长,单位毫秒。
- 类型:
LONG - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:毫秒;-1 表示无限重试
errors.retry.delay.max.ms
框架连续重试之间的最大等待时间,单位毫秒。
- 类型:
LONG - 默认值:
60000 - 重要级别:中
- 有效值 / 注意事项:毫秒;依赖:
errors.retry.timeout
errors.tolerance
框架处理错误时是停止还是容忍并跳过;不覆盖任意 binlog 或启动错误。
- 类型:
STRING - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:none, all
errors.log.enable
是否记录框架处理错误及相关上下文。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false
errors.log.include.messages
是否在错误日志中包含消息内容;开启后可能暴露敏感数据。
- 类型:
BOOLEAN - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true或false;false 表示不向日志写入消息内容;依赖:errors.log.enable
快照
snapshot.locking.mode.custom.name
custom 锁定模式使用的 SnapshotLock 实现名称。
- 类型:
STRING - 默认值:未声明
- 重要级别:中
- 有效值 / 注意事项:
snapshot.locking.mode=custom时指定的 SnapshotLock 实现名称,应与实现的name()返回值一致。
snapshot.select.statement.overrides.<DB>.<TABLE>
指定表快照使用的自定义 SELECT 查询。
- 类型:
STRING - 默认值:未声明
- 重要级别:中
- 有效值 / 注意事项:对应表须列入
snapshot.select.statement.overrides;缺少该表的查询配置时会记录警告并跳过该覆盖查询。
Schema History 后端
schema.history.internal.file.filename
文件型 Schema History 后端保存历史的本地路径。
- 类型:
STRING - 默认值:未声明
- 重要级别:中
- 有效值 / 注意事项:选择 FileSchemaHistory 时必填;有效的本地文件路径。;必填;依赖:
schema.history.internal=io.debezium.storage.file.history.FileSchemaHistory
最佳实践
首次接入时限制快照范围并继续增量同步
适用业务场景:首次接入已有 MySQL 数据库,只需要建立指定表的初始基线,并在快照完成后继续接收这些表的新变化。 配置示例:在快速开始配置中添加以下项;若已设置表排除列表,先移除冲突项。snapshot.include.collection.list 控制本次快照读取哪些表,table.include.list 同时限制后续增量捕获范围。两个列表应保持一致,避免快照基线与后续变更范围不一致。 Schema History Topic 应专用、单分区并保留完整 DDL;不要按普通业务 Topic 清理,也不要随意删除已提交 offset。正常重启依赖原有逻辑身份、offset、Schema History 和仍可读取的 binlog。
已有存量数据但仅从 Connector 启动时捕获后续变更
适用业务场景:MySQL 中已有存量数据,但下游不需要接收这些已有行的快照,只需要从 Connector 启动时开始捕获后续 INSERT、UPDATE 和 DELETE 变更。 配置示例:在快速开始配置中添加或覆盖以下项;按需保留table.include.list 以限制后续捕获范围。
no_data 不发送已有行的数据快照,但仍会读取表结构和其他必要的 Schema 信息,随后继续读取 binlog 进行 CDC。该模式不会自动补齐 Connector 启动前的历史变更;如果需要为下游建立存量基线,应使用 initial 快照场景或通过其他数据初始化流程完成。仍需确保 Schema History 和所需 binlog 可用,并避免删除已提交 offset。
维护大表同步时控制内存和发布批次
适用业务场景:持续变更量较大或下游发布速度波动,需要在内存占用、批量效率和捕获延迟之间做可控取舍。 配置示例:在快速开始配置中添加或覆盖以下项。max.queue.size 必须大于 max.batch.size,应结合 Worker 堆和下游吞吐逐步调整。
监控
监控内容
监控 Worker、Connector 和 Task 的健康状态及重启次数,关注 Source 吞吐、端到端延迟、捕获延迟、offset 提交、错误与重试、队列积压以及 Worker JVM 堆和 GC。若启用了错误容忍或死信队列,再关注跳过记录和 DLQ 活动;错误日志中避免输出消息内容和凭据。导入 Grafana 大盘
从 下载 AutoMQ Connect Cluster Dashboard 获取 JSON,配置 Grafana 使用对应的指标数据源并确保标签与 Connect 集群一致,然后在 Grafana 中导入该 Dashboard。限制条件
- MySQL CDC 由单个 Kafka Connect Task 处理;
tasks.max>1会在生成 Task 配置时抛错,不能并行扩展 CDC 吞吐,快照线程也不会创建额外 Task。 - 同一层级的 include 和 exclude 配置不能同时使用,例如
table.include.list与table.exclude.list,或gtid.source.includes与gtid.source.excludes。 - 不同表 Topic 或不同 Kafka 分区之间没有全局顺序;事务元数据用于关联事务信息,不构成下游原子事务。
- 发生故障恢复时,最后一次持久化 offset 之后已发布的事件可能重复;该 Connector 不提供已确认的端到端 exactly-once 保证。
- 删除事件和 tombstone 是两个不同记录;主键变更可能以旧键删除和新键创建发送到不同分区,消费者不能假设跨键原子更新。
- Schema History 恢复依赖完整历史;丢失或清理 binlog 后重新快照只能恢复当前行状态,不能恢复已丢失的每条中间更新或删除。
常见问题
Connector 启动时报 MySQL binlog 配置错误怎么办?
确认 MySQL 已启用 binary logging,并检查binlog_format=ROW、binlog_row_image=FULL。同时检查用户复制权限、database.server.id 是否唯一,以及 Connector 到 MySQL 端口的访问。
为什么提高 tasks.max 后无法启动多个 Task?
这是该 Connector 的处理模型:MySQL binlog CDC 由一个 Task 处理。tasks.max>1 会在生成 Task 配置时被拒绝,而不是被忽略;需要从过滤范围、批次、队列和下游发布能力着手调优。
重启后提示找不到 Schema 或无法恢复位置怎么办?
先确认 Connector 的逻辑身份、offset Topic、Schema History Topic 和 Kafka 连接配置未改变,并确认历史 Topic 未被清理且仍为单分区。Schema History 丢失与 binlog 位置过期是不同问题:when_needed 不处理 Schema History 错误,重建当前结构也不能补回丢失的 binlog 或历史结构变化。重新快照只能建立当前行状态,不能恢复已过期日志中的全部中间事件。
为什么下游看到删除记录后又收到 null 值?
当tombstones.on.delete=true 时,删除事件之后会发送同键的 null 值 tombstone,用于下游按键清理。删除事件本身仍包含被删除行的变更语义,两者不能混为一个记录。