概述
Debezium SQL Server Source Connector 将 SQL Server 业务表的存量数据和后续行级变更发送到 Kafka,适用于下游数据同步、缓存更新和变化事件处理。首次接入时,它可以读取表的行快照,随后持续读取 SQL Server CDC 捕获的 INSERT、UPDATE 和 DELETE,而不是周期性扫描全表。 默认情况下,每张表对应一个数据 Topic,名称为topic.prefix.<database>.<schema>.<table>。消息键通常来自表的主键,消息值包含操作类型、变更前后数据及源端位点。表结构变化消息与内部结构历史分别保存,内部结构历史用于恢复,不是业务消费入口。
前置条件
- 数据库与表的 CDC:确认 SQL Server 的版本、版本类别及部署方式支持 CDC。由管理员先为数据库启用 CDC,再为需要同步的基础表启用 CDC;数据库级启用需要
sysadmin,表级启用需要db_owner。Connector 不会代为创建或刷新捕获实例。 - CDC 捕获作业与保留窗口:普通 SQL Server 部署需要运行 SQL Server Agent,并确保 CDC 捕获、清理作业正常。保留窗口应覆盖最长停机时间及恢复后的追赶时间;不同托管部署的捕获机制需按其自身要求准备,不能直接套用本地 Agent 的管理方式。
- 运行账号权限:账号应能读取快照源表及捕获列,并访问每个选定数据库的 CDC 元数据和变更数据。配置了 CDC gating role 时,非
sysadmin/db_owner账号还需加入该角色。CDC 管理权限不必赋予常驻运行账号;Agent 服务状态查询另需按服务器版本授予服务 DMV 可见权限,不能把VIEW SERVER STATE当作所有版本的统一要求。也可由 DBA 提供受控的状态查询入口;登录成功不代表这些读取权限已齐全。 - SQL Server TLS:快速开始要求服务器提供受信任的 TLS 证书,证书名称匹配连接主机名,且 Connect Worker 的 JVM 信任该证书链。使用私有 CA 时,应先将 CA 纳入 Worker 的信任配置。驱动默认启用连接加密且校验证书,不以关闭加密或跳过证书校验作为接入前提。
- 内部结构历史:为该 Connector 准备独立的 Kafka history Topic,保持单分区并持久保留完整历史,不让保留策略删除恢复所需记录。确保 history 客户端有读写权限,以及需要自动创建时的建 Topic 权限;这些客户端的认证不能假定完全继承 Worker 配置。
授权许可
使用 Apache License 2.0。快速开始
提前准备 Connect Cluster、Kafka、已启用 CDC 的 SQL Server 数据库和业务表,并确认网络连通、访问权限及 TLS 信任配置满足要求。创建与管理操作参见 AutoMQ 的管理 Connector。下面使用数据库账号认证,继承 Worker 已配置的消息 Converter,并通过受信任证书建立加密连接。1433、默认单 Task 和默认 initial 快照策略:首次无 Offset 时读取符合范围的 CDC 表存量,再持续同步变更。账号凭据应通过部署支持的安全配置机制提供,不写入公开文件或日志。示例依赖已准备好的 JVM 信任链;未配置证书信任的环境不能直接使用。
配置
连接与数据库范围
database.hostname
SQL Server 的可解析主机名或 IP 地址。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:必填;使用 TLS 时需与证书身份匹配。
database.port
SQL Server 连接端口。
- 类型:
int - 默认值:
1433 - 重要级别:高
- 有效值 / 注意事项:填写服务器实际监听端口。
database.instance
SQL Server 命名实例名称。
- 类型:
string - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:仅连接命名实例时填写,并核对其端口解析方式。
database.names
选择同一 SQL Server 连接上的数据库。
- 类型:
list - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:运行时需要非空、逗号分隔的数据库名;逐库准备 CDC 和权限。多库情况下按数据库分配 Task,不按表拆分数据库。
database.user
数据库连接账号。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:使用账号密码认证时填写;其他认证方式由所选 JDBC 模式决定,不能仅靠省略账号启用。
database.password
数据库连接密码。
- 类型:
password - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:按认证方式提供,不在日志或共享配置中暴露。
database.query.timeout.ms
数据库查询执行超时,单位毫秒。
- 类型:
int - 默认值:
600000 - 重要级别:低
- 有效值 / 注意事项:
0表示不限制查询执行时间;区别于连接超时。
connection.validation.timeout.ms
等待数据库连接校验完成的最长时间,单位毫秒。
- 类型:
long - 默认值:
60000 - 重要级别:低
- 有效值 / 注意事项:正整数;不是数据查询超时。
data.query.mode
控制增量 CDC 数据查询方式。
- 类型:
string - 默认值:
function - 重要级别:低
- 有效值 / 注意事项:
function调用 CDC 变更函数;direct直接查询变更表,需匹配其权限要求。
database.sqlserver.agent.status.query
查询 SQL Server Agent 是否正在运行。
- 类型:
string - 默认值:
SELECT CASE WHEN dss.[status]=4 THEN 1 ELSE 0 END AS isRunning FROM #db.sys.dm_server_services dss WHERE dss.[servicename] LIKE N'SQL Server Agent (%'; - 重要级别:中
- 有效值 / 注意事项:返回
isRunning状态,#db为数据库替换标记。此项在数据库子配置中读取,Connector 配置需使用database.前缀;字段后缀为sqlserver.agent.status.query。服务 DMV 权限依服务器版本而定;自定义查询需与 DBA 核对结果语义。查询失败不等同 Agent 已停止。
database.<JDBC-property-name>
向 Microsoft SQL Server JDBC 驱动透传连接属性,用于 TLS、认证及驱动连接行为。
- 类型:
string - 默认值:无固定默认值
- 重要级别:未声明
- 有效值 / 注意事项:去掉
database.或driver.前缀后交给驱动;同名属性冲突时driver.值优先。已单独列出的 Connector 连接参数仍按其自身语义填写,不用透传属性绕过数据库选择。驱动可接受属性不代表全部认证方式或可选依赖都已具备。
encrypt、trustServerCertificate、trustStore、trustStorePassword、trustStoreType 和 hostNameInCertificate。驱动在未覆盖时默认 encrypt=True、trustServerCertificate=false;这是驱动默认行为,不是本透传族的 ConfigDef 默认值。加密不能替代身份校验。证书链不受信任或名称不匹配时,应修复证书与信任配置,不通过禁用加密或跳过校验消除错误。其他 JDBC 属性按该驱动的属性名称、取值及依赖条件设置,不将其完整属性清单视为 Connector 的静态配置清单。
driver.*
使用独立前缀向 JDBC 驱动透传属性。
- 类型:
string - 默认值:无固定默认值
- 重要级别:未声明
- 有效值 / 注意事项:
*替换为实际 JDBC 属性名,去掉driver.前缀后传入驱动;与database.中同名驱动属性冲突时优先。属性可用性、驱动默认值和可选认证依赖与前一透传族遵循同样边界。
Topic 与同步范围
topic.prefix
为源服务器及其事件 Topic 提供标识和命名空间。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:必填;使用字母、数字、连字符、点和下划线,并保持唯一。上线后修改会改变 Topic 与源分区身份,不能作为透明重命名。
topic.naming.strategy
决定数据、结构变化、事务和心跳 Topic 的命名策略。
- 类型:
class - 默认值:
io.debezium.schema.SchemaTopicNamingStrategy - 重要级别:中
- 有效值 / 注意事项:自定义实现需可从插件类路径加载;默认数据 Topic 包含数据库、schema 和表名。
table.include.list
限定需要捕获的业务表。
- 类型:
list - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:逗号分隔的完整名称正则,SQL Server 表过滤使用
schema.table;点需转义。不与table.exclude.list同时设置。匹配不代表自动启用表 CDC。
table.exclude.list
排除不需要捕获的表。
- 类型:
list - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:逗号分隔的
schema.table正则;与table.include.list互斥。
table.ignore.builtin
是否忽略内置表。
- 类型:
boolean - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:
true/false;不改变外部系统对 CDC 基础表的限制。
column.include.list
选择事件值中保留的列。
- 类型:
list - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:逗号分隔的
schema.table.column正则;与column.exclude.list互斥。消息键另行形成,列过滤不等于过滤键内字段。
column.exclude.list
从事件值中排除列。
- 类型:
list - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:逗号分隔的完整列名正则;与
column.include.list互斥,不能据此保证敏感字段不出现在消息键。
skip.messages.without.change
在纳入的列没有变化时跳过消息发布。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:
true/false;结合列过滤理解其作用,不是任意行条件过滤。
skipped.operations
选择流阶段跳过的操作类型。
- 类型:
list - 默认值:
t - 重要级别:低
- 有效值 / 注意事项:
c、u、d、t或none,分别表示创建、更新、删除、截断或不跳过。存在t选项不表示 SQL Server CDC 能产生 TRUNCATE 事件;跳过业务操作会使下游状态不完整。
初始快照
snapshot.mode
决定启动时如何建立存量基线及进入增量读取。
- 类型:
string - 默认值:
initial - 重要级别:低
- 有效值 / 注意事项:
initial在无 Offset 时快照,完成后进入流;initial_only只快照;always每次启动快照;no_data只建立结构后进入流;schema_only是已弃用的no_data别名。when_needed按缺少或不可用位点等条件决定快照,不保证所有 CDC 清理情形都自动恢复。recovery重建结构历史,仅在停机后没有相关结构变化等安全前提下考虑;不能恢复已清理的历史事件。configuration_based由下列开关控制;custom需自定义 Snapshotter。
snapshot.isolation.mode
决定快照的隔离级别和并发影响。
- 类型:
string - 默认值:
repeatable_read - 重要级别:低
- 有效值 / 注意事项:
repeatable_read、exclusive、snapshot、read_committed、read_uncommitted。弱隔离不保证一致快照;snapshot需数据库启用ALLOW_SNAPSHOT_ISOLATION。锁定模式与锁超时会影响锁行为;database.applicationIntent精确为ReadOnly时实现会采用snapshot隔离。
snapshot.lock.timeout.ms
快照开始时等待表锁的最长时间,单位毫秒。
- 类型:
long - 默认值:
10000 - 重要级别:中
- 有效值 / 注意事项:锁获取超时会终止该次快照;调整前核对源库并发负载。
snapshot.locking.mode
在需要锁定的快照隔离模式下控制结构读取阶段的锁。
- 类型:
string - 默认值:
exclusive(公开运行参数定义的默认值) - 重要级别:低
- 有效值 / 注意事项:
exclusive、none、custom;对repeatable_read/exclusive隔离模式生效。none仅在快照期间无结构变化等安全条件下考虑;不保证所有快照阶段都无锁。custom需配套锁实现。
snapshot.locking.mode.custom.name
选择自定义快照锁实现。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:
snapshot.locking.mode=custom时必填,实现需可加载并提供匹配的名称。
snapshot.delay.ms
启动后延迟开始快照的时间,单位毫秒。
- 类型:
long - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:非负整数;不是流读取间隔。
snapshot.fetch.size
快照时每次读取到内存中的记录数。
- 类型:
int - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:非负整数;未设置时使用实现与 JDBC 的读取行为,不以示例批次值作为默认值。
snapshot.max.threads
执行快照使用的最大线程数。
- 类型:
int - 默认值:
1 - 重要级别:中
- 有效值 / 注意事项:正整数;快照线程数不同于 Connect Task 数,增加线程会增加数据库负载。
snapshot.include.collection.list
限定快照读取的数据集合。
- 类型:
list - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:完整集合名正则列表,需在已有捕获范围内;不替代流阶段表过滤。
snapshot.select.statement.overrides
声明需要覆盖快照 SELECT 的表。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:逗号分隔的表标识;SQL Server 使用
<database>.<schema>.<table>,并提供相应的动态 SQL 属性。WHERE 仅影响快照,不是持续行级过滤。
snapshot.select.statement.overrides.<table-id>
指定已声明表的快照 SELECT 语句。
- 类型:
string - 默认值:无
- 重要级别:未声明
- 有效值 / 注意事项:将
<table-id>替换为<database>.<schema>.<table>,且在snapshot.select.statement.overrides中列出;返回列应满足表结构和事件转换要求。
snapshot.tables.order.by.row.count
按行数安排初始快照表顺序。
- 类型:
string - 默认值:
disabled - 重要级别:中
- 有效值 / 注意事项:
disabled、ascending、descending;不改变同步范围。
snapshot.query.mode
选择快照查询生成方式。
- 类型:
string - 默认值:
select_all - 重要级别:低
- 有效值 / 注意事项:
select_all或custom;自定义方式需要查询提供者。
snapshot.query.mode.custom.name
选择自定义快照查询实现。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:
snapshot.query.mode=custom时提供匹配的实现名称,类及依赖应可加载。
snapshot.mode.custom.name
选择自定义 Snapshotter。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:
snapshot.mode=custom时必填,与实现的name()相匹配。
snapshot.mode.configuration.based.snapshot.data
配置型快照策略是否读取表数据。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅
snapshot.mode=configuration_based时生效。
snapshot.mode.configuration.based.snapshot.schema
配置型快照策略是否读取表结构。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅
snapshot.mode=configuration_based时生效,应与数据读取及历史恢复要求协调。
snapshot.mode.configuration.based.start.stream
配置型快照策略是否在快照后进入流读取。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅
snapshot.mode=configuration_based时生效。
snapshot.mode.configuration.based.snapshot.on.data.error
配置型快照策略遇到数据位点相关错误时是否重新快照。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅
snapshot.mode=configuration_based时生效;重新快照不能补回已清理的完整变化历史。
snapshot.mode.configuration.based.snapshot.on.schema.error
配置型快照策略遇到结构历史错误时是否重新快照。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:仅
snapshot.mode=configuration_based时生效;核对结构变化与保留位点是否兼容。
增量快照与信号
incremental.snapshot.chunk.size
增量快照每个分块的行数。
- 类型:
int - 默认值:
1024 - 重要级别:中
- 有效值 / 注意事项:定义接受非负整数;实际分块需有可推进的大小,并结合行大小和数据库负载评估。
incremental.snapshot.allow.schema.changes
控制增量快照期间的结构变化处理。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:此开关不能替代 CDC 捕获实例维护;不能据此认为 SQL Server 任意并发 DDL 都安全,尤其不应在增量快照期间更改主键。
incremental.snapshot.option.recompile
在增量快照 SELECT 中加入 OPTION(RECOMPILE)。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:可避免参数嗅探影响,但增加源数据库编译及 CPU 压力。
incremental.snapshot.watermarking.strategy
控制增量快照窗口水位信号的写入方式。
- 类型:
string - 默认值:
INSERT_INSERT(对应配置值insert_insert) - 重要级别:低
- 有效值 / 注意事项:
insert_insert写入开启和关闭信号;insert_delete写入开启信号,关闭时删除对应信号。需要信号数据集合及相应写入权限。
signal.data.collection
源端信号表的完整名称。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:SQL Server 使用
<database>.<schema>.<table>;源信号流程需准备规定结构、CDC 与相应访问权限,未配置时不启用该源端信号入口。
signal.enabled.channels
启用的信号通道。
- 类型:
list - 默认值:
source - 重要级别:中
- 有效值 / 注意事项:按已提供的通道选择,例如
source、kafka、file、jmx;所选通道需要配套资源及参数。
signal.poll.interval.ms
检查已注册信号通道的间隔,单位毫秒。
- 类型:
long - 默认值:
5000 - 重要级别:中
- 有效值 / 注意事项:正整数;与业务 CDC 轮询间隔不同。
signal.file
文件信号通道读取的文件。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:启用
file通道时必填,Task 所在 Worker 应可访问。
signal.kafka.bootstrap.servers
Kafka 信号通道的集群地址。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:启用
kafka通道时必填,应指向 Connect 使用的 Kafka 集群。
signal.kafka.topic
Kafka 信号通道读取的 Topic。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:启用
kafka通道时必填,并准备消费权限。
signal.kafka.groupId
Kafka 信号消费者的组标识。
- 类型:
string - 默认值:
kafka-signal - 重要级别:低
- 有效值 / 注意事项:属性名中的大写
I应保持不变;组隔离须与信号投递方式协调。
signal.kafka.poll.timeout.ms
等待 Kafka 信号的单次轮询超时,单位毫秒。
- 类型:
int - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:非负整数;仅 Kafka 信号通道生效。
signal.consumer.*
为 Kafka 信号通道透传消费者属性。
- 类型:
string - 默认值:无固定默认值
- 重要级别:未声明
- 有效值 / 注意事项:去掉前缀后交给 Kafka 客户端,用于信号消费者认证等;不等于 Worker 的全部消费者设置。
读取批次与内存
max.batch.size
每次向 Connect 返回的最大记录批次。
- 类型:
int - 默认值:
2048 - 重要级别:中
- 有效值 / 注意事项:正整数,且小于
max.queue.size。
max.queue.size
等待发送的变更事件队列容量。
- 类型:
int - 默认值:
8192 - 重要级别:中
- 有效值 / 注意事项:必须大于
max.batch.size;队列满时产生背压,不是无限缓冲。
max.queue.size.in.bytes
等待发送队列的字节容量限制。
- 类型:
long - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:非负整数;
0不启用此字节上限,事件数量上限仍生效。
max.iteration.transactions
控制每轮 CDC 查询纳入的事务数量,约束多表流读取的内存需求。
- 类型:
int - 默认值:
500 - 重要级别:中
- 有效值 / 注意事项:非负整数;不等于行数或 Kafka 消息批次大小,需结合事务大小评估。
query.fetch.size
JDBC 查询读取到内存的记录批次大小。
- 类型:
int - 默认值:
10000(SQL Server 公开运行参数定义的默认值) - 重要级别:中
- 有效值 / 注意事项:
0委托给 JDBC 默认 fetch 行为,区别于streaming.fetch.size的表读取上限。
streaming.fetch.size
流阶段每张表单次读取的最大行数。
- 类型:
int - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:
0表示不限制;非零值影响每表分批读取,不能当作事务边界。
poll.interval.ms
无新变更事件时等待下一次轮询的时间,单位毫秒。
- 类型:
long - 默认值:
500 - 重要级别:中
- 有效值 / 注意事项:正整数;更频繁轮询会增加查询开销。
streaming.delay.ms
快照完成后延迟进入流阶段的时间,单位毫秒。
- 类型:
long - 默认值:
0 - 重要级别:低
- 有效值 / 注意事项:非负整数;延迟期间仍需保留所需 CDC 数据。
executor.shutdown.timeout.ms
等待 Task 执行器关闭的最长时间,单位毫秒。
- 类型:
long - 默认值:
4000 - 重要级别:中
- 有效值 / 注意事项:正整数;不代替 Offset 与历史持久化。
消息键、类型与字段
message.key.columns
为匹配表指定消息键列。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:分号分隔的
<table-regex>:<column-1>,<column-2>;SQL Server 表匹配使用schema.table。未覆盖时使用表主键;自定义键不会在数据库中创建唯一约束,须自行保证稳定且唯一。
tombstones.on.delete
控制 DELETE 后是否发送同键、空值的 tombstone。
- 类型:
boolean - 默认值:
true - 重要级别:中
- 有效值 / 注意事项:关闭后仍有删除事件,但没有后续 tombstone;消费者应区分两者。
decimal.handling.mode
控制 DECIMAL / NUMERIC 的事件表示。
- 类型:
string - 默认值:
precise - 重要级别:中
- 有效值 / 注意事项:
precise使用 Connect Decimal 精确表示;string为字符串;double可能损失精度。money / smallmoney 也通过 decimal 转换路径处理。
binary.handling.mode
控制二进制列的事件表示。
- 类型:
string - 默认值:
bytes - 重要级别:低
- 有效值 / 注意事项:
bytes、base64、base64-url-safe、hex;与下游反序列化约定一致。
time.precision.mode
控制日期、时间及时间戳表示的精度。
- 类型:
string - 默认值:
adaptive - 重要级别:中
- 有效值 / 注意事项:
adaptive、adaptive_time_microseconds、isostring、microseconds、nanoseconds、connect。connect使用 Connect 毫秒精度;其他模式与实际列类型及转换路径结合,不能认为每种 SQL Server 类型都产生相同形式。
schema.name.adjustment.mode
调整消息 Schema 名称以适配 Converter。
- 类型:
string - 默认值:
none - 重要级别:低
- 有效值 / 注意事项:
none、avro、avro_unicode;后两者分别替换非法字符或进行 Unicode 转义。
field.name.adjustment.mode
调整消息字段名称以适配 Converter。
- 类型:
string - 默认值:
none - 重要级别:低
- 有效值 / 注意事项:
none、avro、avro_unicode;名称调整可能影响下游字段映射。
column.propagate.source.type
为匹配列的字段 Schema 附加源类型与长度信息。
- 类型:
list - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:完整列名正则列表,不改变列值表示。
datatype.propagate.source.type
按源数据类型名称选择需要传播原始类型信息的字段。
- 类型:
list - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:数据库原生类型名称的正则列表;区别于按列名匹配。
column.truncate.to.<N>.chars
将匹配字符串列截断到指定字符数。
- 类型:动态属性值为
string(列正则列表) - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:
N是属性名中的非负整数,值为逗号分隔的完整列名正则。静态模式项column.truncate.to.(d+).chars的类型元数据为int,但不能把模式名称当作实际配置键;具体数字属性的值不是字符数。
column.mask.with.<N>.chars
以指定数量的星号掩码替代匹配列内容。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:
N是非负整数,值为完整列名正则列表;静态模式名称column.mask.with.(d+).chars不直接作为配置键。掩码不能替代消息键字段的安全评估。
column.mask.hash.<algorithm>.with.salt.<salt>
以带盐哈希替代匹配列内容。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:对应模式项
column.mask.hash.([^.]+).with.salt.(.+);将算法和盐替换为实际名称和值,属性值为列正则列表。算法须由 JVM 支持,盐应妥善管理。
column.mask.hash.v2.<algorithm>.with.salt.<salt>
使用 v2 带盐哈希映射匹配列。
- 类型:
string - 默认值:无
- 重要级别:未声明
- 有效值 / 注意事项:值为完整列名正则列表,算法须由 JVM 支持;与原哈希方式不是可无条件互换的下游标识。
unavailable.value.placeholder
标记源端未提供的原始列值。
- 类型:
string - 默认值:
__debezium_unavailable_value - 重要级别:中
- 有效值 / 注意事项:区别于真实业务值及
null,下游不应将占位符当作真实内容。
include.schema.changes
是否向结构变化 Topic 发布结构消息。
- 类型:
boolean - 默认值:
true - 重要级别:中
- 有效值 / 注意事项:默认 Topic 名为
topic.prefix;与内部结构历史独立,不会自动刷新 SQL Server CDC 捕获实例。
include.schema.comments
是否将表与列注释纳入结构元数据。
- 类型:
boolean - 默认值:
false - 重要级别:中
- 有效值 / 注意事项:启用会增加结构元数据的内存开销。
provide.transaction.metadata
是否生成事务元数据及事件计数。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:不提供跨数据库或跨分区全局顺序,也不等同端到端 exactly-once。
transaction.metadata.factory
生成事务上下文与事务消息结构的工厂类。
- 类型:
class - 默认值:
io.debezium.pipeline.txmetadata.DefaultTransactionMetadataFactory - 重要级别:低
- 有效值 / 注意事项:自定义类及其依赖应在插件类路径中可加载。
sourceinfo.struct.maker
生成事件 source 结构的实现类。
- 类型:
class - 默认值:
io.debezium.connector.sqlserver.SqlServerSourceInfoStructMaker - 重要级别:低
- 有效值 / 注意事项:自定义实现可能改变下游依赖的源位点字段,应保持消费协议兼容。
extended.headers.enabled
是否附加 Debezium 上下文消息头。
- 类型:
boolean - 默认值:
true - 重要级别:低
- 有效值 / 注意事项:消息头携带来源上下文,不替代消息值中的业务字段。
结构历史与恢复
schema.history.internal
保存并恢复数据库结构历史的实现。
- 类型:
class - 默认值:
io.debezium.storage.kafka.history.KafkaSchemaHistory - 重要级别:低
- 有效值 / 注意事项:实现类及依赖需可加载,并配置该实现的存储属性。
schema.history.internal.kafka.bootstrap.servers
Kafka 结构历史客户端使用的集群地址。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:默认 Kafka history 实现下必填,使用 Connect 所在 Kafka 集群的
host:port列表。
schema.history.internal.kafka.topic
保存内部结构历史的 Topic。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:默认 Kafka history 实现下必填;单分区、完整历史持久保留,不与其他独立 Connector 共用。该要求不限制业务数据 Topic 分区数。
schema.history.internal.kafka.create.timeout.ms
创建 history Topic 的等待超时,单位毫秒。
- 类型:
long - 默认值:
30000 - 重要级别:低
- 有效值 / 注意事项:正整数;仅 Kafka history 实现生效。
schema.history.internal.kafka.query.timeout.ms
查询 Kafka 集群信息的超时,单位毫秒。
- 类型:
long - 默认值:
3000 - 重要级别:低
- 有效值 / 注意事项:正整数;仅 Kafka history 实现生效。
schema.history.internal.kafka.recovery.attempts
历史恢复时连续未读到记录的轮询次数。
- 类型:
int - 默认值:
100 - 重要级别:低
- 有效值 / 注意事项:与恢复轮询间隔共同决定无记录等待时长,不是数据库连接重试次数。
schema.history.internal.kafka.recovery.poll.interval.ms
恢复 history 数据时的轮询间隔,单位毫秒。
- 类型:
int - 默认值:
100 - 重要级别:低
- 有效值 / 注意事项:非负整数;仅 Kafka history 实现生效。
schema.history.internal.file.filename
文件结构历史存储路径。
- 类型:
string - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:选用
io.debezium.storage.file.history.FileSchemaHistory时必填;文件需持久存在,并考虑 Task 调度到其他 Worker 后的可访问性。
schema.history.internal.skip.unparseable.ddl
控制历史处理遇到不可解析 DDL 时是否跳过。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:跳过可能丢失结构元数据;此通用 history 选项不代表 SQL Server 流中所有 DDL 都可自动处理。
schema.history.internal.store.only.captured.databases.ddl
仅保存捕获数据库范围内的结构历史。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:缩小历史范围可能影响后续扩大同步范围时的恢复准备。
schema.history.internal.store.only.captured.tables.ddl
仅保存捕获表范围内的结构历史。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:后续纳入新表前需核对历史完整性,不能认为修改过滤就完成存量回填。
schema.history.internal.producer.*
为 Kafka history 写入客户端透传属性。
- 类型:
string - 默认值:无固定默认值
- 重要级别:未声明
- 有效值 / 注意事项:去掉前缀后传给 Kafka 客户端,可提供安全连接属性;不等同业务消息生产者配置。
schema.history.internal.consumer.*
为 Kafka history 恢复读取客户端透传属性。
- 类型:
string - 默认值:无固定默认值
- 重要级别:未声明
- 有效值 / 注意事项:去掉前缀后传给 Kafka 客户端;与 history 生产者认证协调,不能假定继承 Worker 的全部认证。
心跳与错误处理
heartbeat.interval.ms
发送心跳消息的间隔,单位毫秒。
- 类型:
int - 默认值:
0 - 重要级别:中
- 有效值 / 注意事项:非负整数;
0禁用心跳。不能仅靠心跳保证恢复位点仍在 CDC 保留窗口内。
heartbeat.topics.prefix
心跳 Topic 的命名前缀。
- 类型:
string - 默认值:
__debezium-heartbeat - 重要级别:低
- 有效值 / 注意事项:仅启用心跳时产生相应消息,不是业务表 Topic 前缀。
heartbeat.action.query
每次心跳执行的数据库查询。
- 类型:
string - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:查询需具备执行权限,并评估其写入与源库负载影响。
errors.max.retries
Connector 可重试错误的重试预算。
- 类型:
int - 默认值:
-1 - 重要级别:低
- 有效值 / 注意事项:
-1不限次数,0不重试,正数为上限;并非所有初始化、转换或 Kafka 写入错误都由此处理。权限和配置错误不能靠重试自愈。
retriable.restart.connector.wait.ms
可重试异常后等待重启的时间,单位毫秒。
- 类型:
long - 默认值:
10000 - 重要级别:低
- 有效值 / 注意事项:正整数;结合 CDC 保留时间评估长期故障风险。
event.processing.failure.handling.mode
处理损坏或无法处理的事件时的行为。
- 类型:
string - 默认值:
fail - 重要级别:中
- 有效值 / 注意事项:
fail停止;warn记录后跳过;ignore/skip跳过。跳过会丢失相应事件,不等同无损容错。
event.converting.failure.handling.mode
列值转换失败时的行为。
- 类型:
string - 默认值:
warn - 重要级别:中
- 有效值 / 注意事项:
fail抛出异常;warn将失败列置空并记录警告;skip将失败列置空并以 debug 记录。不能把此处的skip理解为完整事件必然被丢弃。
扩展处理与通知
converters
Debezium 自定义类型转换器的别名列表。
- 类型:
string - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:为每个别名提供
<converter-alias>.type;不同于 Kafka Connect 的key.converter/value.converter。
<converter-alias>.type
指定 Debezium 自定义转换器类。
- 类型:
class - 默认值:无
- 重要级别:未声明
- 有效值 / 注意事项:别名需在
converters中列出,实现及依赖应位于插件类路径。
<converter-alias>.<option>
为自定义转换器提供实现专属选项。
- 类型:
string - 默认值:无固定默认值
- 重要级别:未声明
- 有效值 / 注意事项:去掉别名前缀后传给所选实现,按该实现核对属性。
post.processors
Debezium 后处理器的别名列表。
- 类型:
string - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:每个别名需配置实现类,不等同 Connect SMT 列表。
<post-processor-alias>.type
指定后处理器实现类。
- 类型:
class - 默认值:无
- 重要级别:未声明
- 有效值 / 注意事项:别名需在
post.processors中列出,实现类应可加载。
<post-processor-alias>.<option>
为后处理器提供实现专属选项。
- 类型:
string - 默认值:无固定默认值
- 重要级别:未声明
- 有效值 / 注意事项:仅所选后处理器生效,按其接口与依赖配置。
notification.enabled.channels
启用的通知通道列表。
- 类型:
list - 默认值:无(
null) - 重要级别:中
- 有效值 / 注意事项:按可用通道配置;其中
sink指通知写入通道,不是 Sink Connector 方向。
notification.sink.topic.name
通知写入通道使用的 Topic。
- 类型:
string - 默认值:无(
null) - 重要级别:高
- 有效值 / 注意事项:启用
sink通知通道时必填,并准备相应写入权限。
custom.metric.tags
为 Connector 的 MBean 对象名称附加标签。
- 类型:
list - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:逗号分隔的
<key>=<value>列表;避免凭据与高基数值,并协调已有采集标签。
数据血缘
openlineage.integration.enabled
是否通过 OpenLineage 发布数据血缘元数据。
- 类型:
boolean - 默认值:
false - 重要级别:低
- 有效值 / 注意事项:启用时需要配套配置文件及可用的 OpenLineage 服务设置。
openlineage.integration.config.file.path
OpenLineage 客户端配置文件路径。
- 类型:
string - 默认值:
./openlineage.yml - 重要级别:低
- 有效值 / 注意事项:启用血缘集成时,Task 所在 Worker 需可读取;相对路径取决于运行目录。
openlineage.integration.job.description
血缘任务的说明。
- 类型:
string - 默认值:
Debezium change data capture job - 重要级别:低
- 有效值 / 注意事项:仅血缘集成启用时生效。
openlineage.integration.job.namespace
血缘任务命名空间。
- 类型:
string - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:结合血缘系统的任务标识约定填写。
openlineage.integration.job.owners
血缘任务负责人信息。
- 类型:
list - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:逗号分隔的
<key>=<value>列表,避免不必要的个人敏感信息。
openlineage.integration.job.tags
血缘任务标签。
- 类型:
list - 默认值:无(
null) - 重要级别:低
- 有效值 / 注意事项:逗号分隔的
<key>=<value>列表,不影响 CDC 捕获范围。
Connect 框架
connector.class
加载本 Connector 的实现类。
- 类型:
string - 默认值:无
- 重要级别:高
- 有效值 / 注意事项:必填,使用
io.debezium.connector.sqlserver.SqlServerConnector。
tasks.max
Connect 可创建的 Task 数量上限。
- 类型:
int - 默认值:
1 - 重要级别:高
- 有效值 / 注意事项:至少为
1。使用database.names时按数据库轮询分配,实际 Task 数为数据库数与此上限的较小值;同一数据库不会拆给多个 Task。单数据库调大此值不会增加表级并行度;此上限不保证跨库全局顺序。
key.converter
将 Connect 消息键序列化为 Kafka 消息键的 Converter。
- 类型:
class - 默认值:
null(Connector 级不覆盖) - 重要级别:低
- 有效值 / 注意事项:未设置时继承 Worker 配置;设置时需可加载的 Converter 类,并与下游反序列化协议一致。
value.converter
将 Connect 消息值序列化为 Kafka 消息值的 Converter。
- 类型:
class - 默认值:
null(Connector 级不覆盖) - 重要级别:低
- 有效值 / 注意事项:未设置时继承 Worker 配置;保留或调整事件结构前,应核对下游对 Schema、DELETE 与 tombstone 的处理。
errors.tolerance
控制 Connect 转换及相关受框架管理步骤的错误容忍行为。
- 类型:
string - 默认值:
none - 重要级别:中
- 有效值 / 注意事项:
none失败即停止;all可跳过问题记录,带来数据缺失风险。不替代 Connector 的数据库错误重试及事件处理策略。
最佳实践
首次接入时明确业务表范围
适用业务场景:第一次同步已有业务数据库,希望先建立指定业务表的存量基线,再持续接收它们的变化,避免把不相关的 CDC 表一并纳入。 配置示例:在快速开始配置中添加下列配置,若已存在table.exclude.list,应移除该冲突项。将占位符替换为需要同步的完整表名正则;点用反斜杠转义,例如文字形式的 dbo\.orders。
initial 已提供首次快照加持续增量,不必重复设置该默认项。快照会影响数据库负载并可能获取锁,应安排合适接入时间,核对快照行和后续增删改均进入对应 Topic。无稳定键的表不应直接假定可按行幂等更新。后续扩大过滤范围不是首次接入的简单重跑:已有 Offset 时修改列表不会自动补齐新增表存量,需要另外规划回填。
后续计划维护继续复用这一配置,不另建 Connector,也不清空 Offset 或更换前缀。停机前确认 Offset 已提交,备份并保留 Offset 存储和结构历史,确保 CDC 保留窗口覆盖停机与追赶;恢复后核对 Task 状态、追赶进度及新变更是否送达。普通至少一次运行仍可能重放已发送但未提交位点的事件,下游需按稳定业务键及源 LSN、事件序号设计幂等,不能只靠同事务共享的 commit_lsn 去重。
监控
监控内容
关注 Kafka Connect 服务健康、Connector / Task 状态、吞吐与端到端延迟、Offset 提交进度及失败、错误与重试,以及 Worker JVM 的堆内存、GC 和线程状态。持续观察暂停后能否恢复及故障后的积压是否收敛;只有部署启用了相应错误处理和 DLQ 时,才关注 DLQ 活动,不把 DLQ 当作数据库 CDC 读取故障的通用兜底。导入 Grafana 大盘
下载 AutoMQ Connect Cluster Dashboard,确认 Grafana 数据源能够查询采集到的 Connect / Worker 指标,且集群、Connector、Task 等标签与采集配置一致,再在 Grafana 导入 JSON 并选择对应数据源。限制条件
- CDC 捕获面向已启用的基础表,不能将视图或任意数据库对象作为同等捕获目标;不支持依赖 SQL Server CDC 获得 TRUNCATE 行事件。
- 单数据库只分配一个 Task,多库可按数据库分配 Task;不能通过增大
tasks.max将一个数据库按表并行,也不保证跨数据库、跨 Task 或跨 Kafka 分区的全局顺序。 - 表 include / exclude 互斥,列 include / exclude 也互斥;快照 SELECT 中的 WHERE 不能作为持续流的行级过滤。
- 表 DDL 不会自动刷新 CDC 捕获实例。需要 DBA 协同创建新实例并安全切换,Connector 尚未读完旧实例时不能提前删除它;不能认为所有 DDL 都会透明同步。
- Offset 与结构历史是不同的恢复资料,任一丢失都不能保证无缝续传;结构历史的单分区要求不适用于业务数据 Topic。
- CDC 清理后的历史变化不能靠重新快照完整补回;重新快照只能建立当前行状态,而且中断后的快照重跑可能重复发布行。
- 常规非事务性 Source 运行需按至少一次交付设计,不承诺端到端 exactly-once,下游幂等和事件去重仍由应用承担。
常见问题
连接时报证书不可信或主机名不匹配怎么办?
确认连接主机名与 SQL Server 证书匹配,服务器提供完整证书链,且 Worker JVM 信任其 CA。私有 CA 应纳入 Worker 的信任配置,并确认实际运行 Task 的 Worker 使用该配置。驱动默认加密且验证证书;修复证书与信任链,不用关闭加密或跳过校验掩盖问题。可以读取存量,但没有新的增量消息怎么办?
检查是否选用了initial_only,以及目标数据库与表的 CDC 是否启用、SQL Server Agent 和捕获作业是否运行、表是否通过过滤、运行账号是否能读取 CDC 数据。没有最大 LSN 也可能是尚未发生被捕获的变化;Agent 状态查询权限不足只说明无法完成该诊断,不直接证明 Agent 停止。先确认源端捕获表出现新记录,再排查 Task 状态及写入错误。
增大 Task 上限后为什么单库吞吐没有增加?
数据库是 Task 分配边界,一个数据库不会拆成多个 Task。核对瓶颈是否在源端捕获、查询、事件队列或 Kafka 写入,再结合负载和内存测量调整相应批次设置;不要用增加 Task 数替代表内并行,也不要将多个独立 Connector 共用同一 history Topic。重启后收到重复行或重复变更怎么办?
检查重启是否发生在初始快照未完成时、是否使用always,以及 Offset 是否成功持久化。保留原身份、数据库选择与结构历史,避免重置位点造成重新导入。故障窗口中的重复是至少一次运行需要处理的情况;下游按稳定键应用变化,并结合源位点与事件序号识别事件,正确处理 DELETE、主键变化和 tombstone。