Skip to main content

概述

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,并通过受信任证书建立加密连接。
替换环境占位符,其中主机名应与证书名称匹配,Topic 前缀应在其他 Connector 中保持唯一,history Topic 不得与业务 Topic 或其他独立 Connector 共用。此示例使用默认端口 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 连接参数仍按其自身语义填写,不用透传属性绕过数据库选择。驱动可接受属性不代表全部认证方式或可选依赖都已具备。
TLS 属性包括 encrypttrustServerCertificatetrustStoretrustStorePasswordtrustStoreTypehostNameInCertificate。驱动在未覆盖时默认 encrypt=TruetrustServerCertificate=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
  • 重要级别:低
  • 有效值 / 注意事项cudtnone,分别表示创建、更新、删除、截断或不跳过。存在 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_readexclusivesnapshotread_committedread_uncommitted。弱隔离不保证一致快照;snapshot 需数据库启用 ALLOW_SNAPSHOT_ISOLATION。锁定模式与锁超时会影响锁行为;database.applicationIntent 精确为 ReadOnly 时实现会采用 snapshot 隔离。

snapshot.lock.timeout.ms

快照开始时等待表锁的最长时间,单位毫秒。
  • 类型long
  • 默认值10000
  • 重要级别:中
  • 有效值 / 注意事项:锁获取超时会终止该次快照;调整前核对源库并发负载。

snapshot.locking.mode

在需要锁定的快照隔离模式下控制结构读取阶段的锁。
  • 类型string
  • 默认值exclusive(公开运行参数定义的默认值)
  • 重要级别:低
  • 有效值 / 注意事项exclusivenonecustom;对 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
  • 重要级别:中
  • 有效值 / 注意事项disabledascendingdescending;不改变同步范围。

snapshot.query.mode

选择快照查询生成方式。
  • 类型string
  • 默认值select_all
  • 重要级别:低
  • 有效值 / 注意事项select_allcustom;自定义方式需要查询提供者。

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
  • 重要级别:中
  • 有效值 / 注意事项:按已提供的通道选择,例如 sourcekafkafilejmx;所选通道需要配套资源及参数。

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
  • 重要级别:低
  • 有效值 / 注意事项bytesbase64base64-url-safehex;与下游反序列化约定一致。

time.precision.mode

控制日期、时间及时间戳表示的精度。
  • 类型string
  • 默认值adaptive
  • 重要级别:中
  • 有效值 / 注意事项adaptiveadaptive_time_microsecondsisostringmicrosecondsnanosecondsconnectconnect 使用 Connect 毫秒精度;其他模式与实际列类型及转换路径结合,不能认为每种 SQL Server 类型都产生相同形式。

schema.name.adjustment.mode

调整消息 Schema 名称以适配 Converter。
  • 类型string
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项noneavroavro_unicode;后两者分别替换非法字符或进行 Unicode 转义。

field.name.adjustment.mode

调整消息字段名称以适配 Converter。
  • 类型string
  • 默认值none
  • 重要级别:低
  • 有效值 / 注意事项noneavroavro_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
关键说明:先由 DBA 为目标表启用 CDC,并核对账号读取权限,再创建 Connector。默认 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。

新增列后消息中仍看不到该列怎么办?

核对列过滤与事件 Converter,再检查 CDC 捕获实例是否仍使用旧列结构。基础表加列不会自动更新旧捕获表,需 DBA 按服务器支持的流程准备新实例,等待 Connector 安全切换,保留所需结构历史后再清理旧实例。不要为此直接删 history Topic 或清空 Offset。

停机过久导致 CDC 位点不可用怎么办?

先核对保留窗口、清理作业和当前可读位点,并保留 Offset 与 history 的备份。已被清理的完整变化历史无法从当前表重建;明确下游需要当前状态还是完整审计历史后,再规划重新建立基线或其他补数途径。不要假定某个自动快照模式能够覆盖所有失效路径,或把结构历史恢复当作日志恢复。