如何配置Kafka Connect Sink实现MSSQL跨库增删改操作全同步
MSSQL Kafka Connect Sink 配置方案
你需要使用Confluent JDBC Sink Connector适配Debezium生成的CDC事件结构,即可实现增删改操作全量同步,完整配置如下:
{ "name": "sql-sink", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "topics": "TOP1", "connection.url": "jdbc:sqlserver://<目标库IP>:1433;databaseName=<目标库名称>;encrypt=false;trustServerCertificate=true", "connection.user": "<目标库账号>", "connection.password": "<目标库密码>", "dialect.name": "SqlServerDatabaseDialect", "key.converter": "io.confluent.connect.avro.AvroConverter", "key.converter.schema.registry.url": "http://kafka01.xxxx.dev:8081", "value.converter": "io.confluent.connect.avro.AvroConverter", "value.converter.schema.registry.url": "http://kafka01.xxxx.dev:8081", "transforms": "unwrap", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.drop.tombstones": "false", "transforms.unwrap.delete.handling.mode": "none", "insert.mode": "upsert", "pk.mode": "record_key", "delete.enabled": "true", "auto.create": "false", "auto.evolve": "false", "errors.log.enable": "true", "errors.log.include.messages": "true" } }
核心参数说明
transforms.unwrap相关配置:负责提取Debezium CDC事件里的after字段作为写入目标库的实际数据,同时保留删除事件的墓碑消息,供Sink识别删除操作insert.mode=upsert:遇到主键冲突时自动执行更新操作,对应源库的UPDATE事件delete.enabled=true:开启删除事件支持,监听到CDC的删除事件后会自动删除目标库对应主键的记录dialect.name=SqlServerDatabaseDialect:指定SQL Server专属方言,避免语法适配问题pk.mode=record_key:使用Kafka消息的key作为表主键,Debezium默认会将源表主键设为消息key,保证更新、删除的准确性
注意事项
- 源表和目标表的结构必须完全对齐,字段名、字段类型、主键定义不能有差异,否则会出现写入失败
- 生产环境建议提前创建好目标表,关闭
auto.create和auto.evolve配置,避免自动建表出现字段类型不符合预期的问题 - 如果是首次同步,建议先启动Source连接器完成初始全量快照,再启动Sink连接器,避免全量数据遗漏
内容的提问来源于stack exchange,提问作者Arash Mousavi
相关产品推荐
相关产品推荐

