如何在SingleStore Kafka Connect中配置仅同步指定列数据
SingleStore Sink Connector 字段过滤配置方案
问题原因
SingleStore 官方提供的 Kafka Sink Connector 本身不支持 fields.whitelist 配置项,该配置是 Confluent JDBC Sink Connector 专属参数,所以在你当前的配置中不会生效。
实现方案
使用 Kafka Connect 原生提供的 ReplaceField 转换器实现字段白名单过滤,具体修改如下:
- 调整 transforms 配置,新增字段过滤的转换器节点
- 删除原有的无效配置
fields.whitelist - 按需补充需要保留的特殊字段(如删除标记字段、Debezium 附加字段)
修改后的完整配置
{ "name": "my-connector", "config": { "connector.class":"com.singlestore.kafka.SingleStoreSinkConnector", "tasks.max":"1", // 新增字段过滤转换器到列表 "transforms": "dropPrefix,unwrap,filterFields", "transforms.dropPrefix.type": "org.apache.kafka.connect.transforms.RegexRouter", "transforms.dropPrefix.regex": "dbserver1.inventory.(.*)", "transforms.dropPrefix.replacement": "$1", "errors.tolerance": "all", "errors.log.enable": "true", "errors.log.include.messages": "true", "transforms.unwrap.type": "io.debezium.transforms.ExtractNewRecordState", "transforms.unwrap.delete.handling.mode":"rewrite", "transforms.unwrap.add.fields": "ts_ms", // 新增字段过滤规则 "transforms.filterFields.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", // 不需要保留ts_ms和__deleted的话,直接写id,city即可 "transforms.filterFields.whitelist": "id,city,__deleted,ts_ms", "topics":"dbserver1.inventory.addresses", "connection.ddlEndpoint" : "memsql:3306", "connection.database" : "test", "connection.user" : "root", "connection.password": "password", "insert.mode": "upsert", "tableKey.primary.keyName" : "id", "auto.create": "true", "auto.evolve": "true", "singlestore.metadata.allow": true, "singlestore.metadata.table": "kafka_connect_transaction_metadata" } }
注意事项
__deleted是delete.handling.mode=rewrite模式下生成的删除标记字段,如果过滤掉该字段会导致CDC的删除事件无法同步到SingleStore- 字段名需要和Debezium输出的CDC数据中的字段名完全匹配,大小写敏感
- 转换器顺序不可调整,必须先执行unwrap解包Debezium嵌套结构,再执行字段过滤才能生效
内容的提问来源于stack exchange,提问作者yekmolsoheil
相关产品推荐
相关产品推荐

