Skip to main content

概述

Debezium MySQL Source Connector 从 MySQL 读取表结构、存量行和 binlog 中的行级及 DDL 变更,并将它们作为 Kafka Connect SourceRecord 发布到 Kafka。首次运行通常先执行快照,再从快照边界继续增量读取;每张表的变更默认进入按 Topic 前缀、数据库名和表名组织的 Topic。它适用于把 MySQL 数据库变更事件持续送入 Kafka,供下游实时计算、缓存同步、审计或数据湖处理。

前置条件

  • MySQL 必须启用二进制日志,并使用 binlog_format=ROWbinlog_row_image=FULL;日志保留时间应覆盖预期的故障恢复和同步延迟。
  • 连接用户需要对快照表具备 SELECT 权限,并具备复制和位置读取所需的 REPLICATION SLAVEREPLICATION CLIENT 权限;发现数据库和快照锁定可能还需要 SHOW DATABASESRELOAD,表锁回退还需要 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
  • 重要级别:低
  • 有效值 / 注意事项truefalse;依赖: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
  • 重要级别:中
  • 有效值 / 注意事项truefalse;依赖:snapshot.mode=configuration_based

snapshot.mode.configuration.based.snapshot.schema

在 configuration_based 模式下决定是否快照表结构。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse;依赖:snapshot.mode=configuration_based

snapshot.mode.configuration.based.start.stream

在 configuration_based 模式下决定是否开始增量读取。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse;依赖:snapshot.mode=configuration_based

snapshot.mode.configuration.based.snapshot.on.schema.error

在 configuration_based 模式下决定遇到结构历史错误时是否请求快照。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse;依赖:snapshot.mode=configuration_based

snapshot.mode.configuration.based.snapshot.on.data.error

在 configuration_based 模式下决定日志位置不可用时是否请求快照。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse;依赖: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
  • 重要级别:低
  • 有效值 / 注意事项truefalse;依赖: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
  • 重要级别:低
  • 有效值 / 注意事项truefalse

类型转换

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
  • 重要级别:低
  • 有效值 / 注意事项truefalse;依赖:connect.keep.alive.interval.ms

connect.keep.alive.interval.ms

binlog 连接保活检查间隔,单位毫秒。
  • 类型LONG
  • 默认值60000
  • 重要级别:低
  • 有效值 / 注意事项:正整数;依赖:connect.keep.alive

use.nongraceful.disconnect

是否采用非正常关闭方式断开 binlog 客户端连接。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse

快照

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_microsecondsconnect;不接受 adaptive

enable.time.adjuster

是否调整两位年份的时间值。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项truefalse

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
  • 重要级别:低
  • 有效值 / 注意事项truefalse;不支持主键变化;依赖:incremental.snapshot.chunk.size

快照

snapshot.locking.mode

控制快照读取期间使用的锁策略。
  • 类型STRING
  • 默认值minimal
  • 重要级别:低
  • 有效值 / 注意事项extendedminimalminimal_perconaminimal_percona_no_table_locksnonecustomminimal 在读取结构和元数据时持有全局读锁,随后依赖 REPEATABLE READ 和兼容的事务存储引擎读取行;不保证任意存储引擎或并发 DDL 下的一致性。none 不获取快照锁,使用边界为 snapshot.mode=schema_onlyschema_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.topicschema.history.internal.kafka.bootstrap.servers

schema.history.internal.skip.unparseable.ddl

遇到无法解析的 DDL 时是否跳过该 DDL;跳过可能导致 Schema 元数据不完整。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项truefalse;跳过 DDL 可能导致后续事件缺少正确的字段定义;依赖:schema.history.internal

schema.history.internal.store.only.captured.tables.ddl

是否只在内部 Schema History 中保存已捕获表的 DDL。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项truefalse;扩大表过滤范围前应确认所需历史 Schema 仍可获得;依赖:schema.history.internaltable.include.listtable.exclude.list

schema.history.internal.store.only.captured.databases.ddl

是否只在内部 Schema History 中保存已捕获数据库的 DDL。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:低
  • 有效值 / 注意事项truefalse;扩大数据库过滤范围前应确认所需历史 Schema 仍可获得;依赖:schema.history.internaldatabase.include.listdatabase.exclude.list

扩展

converters

自定义类型转换器的别名列表。
  • 类型STRING
  • 默认值:无
  • 重要级别:低
  • 有效值 / 注意事项:逗号分隔的别名;每个别名的 .type 指定 CustomConverter 实现类

post.processors

事件后处理器的别名列表。
  • 类型STRING
  • 默认值:无
  • 重要级别:低
  • 有效值 / 注意事项:逗号分隔的别名;每个别名的 .type 指定后处理器实现类

事件

tombstones.on.delete

删除事件后是否发送同键的 null 值 tombstone。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:中
  • 有效值 / 注意事项truefalse

心跳

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
  • 重要级别:中
  • 有效值 / 注意事项truefalse

include.schema.comments

是否在 Schema 中包含数据库注释;可能增加内存占用。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse;可能增加内存占用

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
  • 重要级别:中
  • 有效值 / 注意事项truefalse;需要 MySQL 设置 binlog_rows_query_log_events=ON;可能暴露被过滤的数据

过滤

table.ignore.builtin

是否忽略数据库内置表。
  • 类型BOOLEAN
  • 默认值true
  • 重要级别:低
  • 有效值 / 注意事项truefalse;依赖: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
  • 重要级别:中
  • 有效值 / 注意事项truefalse;依赖: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
  • 重要级别:低
  • 有效值 / 注意事项truefalse;已弃用,当前没有替代配置。

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
  • 重要级别:中
  • 有效值 / 注意事项truefalse

errors.log.include.messages

是否在错误日志中包含消息内容;开启后可能暴露敏感数据。
  • 类型BOOLEAN
  • 默认值false
  • 重要级别:中
  • 有效值 / 注意事项truefalse;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。

维护大表同步时控制内存和发布批次

适用业务场景:持续变更量较大或下游发布速度波动,需要在内存占用、批量效率和捕获延迟之间做可控取舍。 配置示例:在快速开始配置中添加或覆盖以下项。
关键说明:队列用于吸收短时发布波动,不是持久化检查点;增大队列会增加内存占用,发布持续变慢时仍会增加捕获延迟和 binlog 保留压力。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.listtable.exclude.list,或 gtid.source.includesgtid.source.excludes
  • 不同表 Topic 或不同 Kafka 分区之间没有全局顺序;事务元数据用于关联事务信息,不构成下游原子事务。
  • 发生故障恢复时,最后一次持久化 offset 之后已发布的事件可能重复;该 Connector 不提供已确认的端到端 exactly-once 保证。
  • 删除事件和 tombstone 是两个不同记录;主键变更可能以旧键删除和新键创建发送到不同分区,消费者不能假设跨键原子更新。
  • Schema History 恢复依赖完整历史;丢失或清理 binlog 后重新快照只能恢复当前行状态,不能恢复已丢失的每条中间更新或删除。

常见问题

Connector 启动时报 MySQL binlog 配置错误怎么办?

确认 MySQL 已启用 binary logging,并检查 binlog_format=ROWbinlog_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,用于下游按键清理。删除事件本身仍包含被删除行的变更语义,两者不能混为一个记录。

为什么重启后下游可能再次收到已经处理过的事件?

Source offset 的持久化与事件发布不是一个不可分割的端到端事务。若事件已经发布,但 Worker 尚未持久化对应 offset 就发生故障,恢复后可能从较早位置重新读取并再次发布这些事件。下游应根据业务键和 Source 元数据实现幂等处理,并确保 offset、Schema History 和所需 binlog 一起保留。