Kafka Connect JdbcSinkConnector能否映射记录字段到不同名PostgreSQL列
Kafka JdbcSinkConnector 自定义字段映射配置方案
该需求完全可实现,无需修改源数据或目标表结构,通过 Kafka Connect 内置的 ReplaceField 单消息转换(SMT)能力即可完成字段重命名映射,具体配置如下:
核心逻辑
JdbcSinkConnector 默认按消息字段名与目标表列名严格匹配写入,只需在数据进入连接器写入逻辑前,将源端字段field_x重命名为目标列名column_field_x即可实现映射。
完整配置示例
在原有连接器配置基础上新增 SMT 转换规则,调整后的核心配置片段如下:
{ "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": "1", "topics": "你的Kafka Topic名称", "connection.url": "jdbc:postgresql://<PG地址>:<端口>/<数据库名>", "connection.user": "PG数据库用户名", "connection.password": "PG数据库密码", "auto.create": "false", "auto.evolve": "false", // 白名单需填写重命名后的目标字段名 "fields.whitelist": "column_field_x", // 新增字段重命名转换配置 "transforms": "renameFieldX", "transforms.renameFieldX.type": "org.apache.kafka.connect.transforms.ReplaceField$Value", "transforms.renameFieldX.renames": "field_x:column_field_x" }
配置说明
transforms:定义当前转换规则的别名,多个转换规则可用逗号分隔transforms.renameFieldX.type:指定使用值维度的ReplaceField转换,若需重命名Key维度的字段,将$Value替换为$Key即可transforms.renameFieldX.renames:配置映射规则,格式为源字段名:目标字段名,多个字段映射可逗号分隔,例如field1:col1,field2:col2fields.whitelist需填写重命名后的字段名,因为白名单过滤逻辑在SMT转换之后执行
注意事项
- 若消息为嵌套结构,需先通过
FlattenSMT将嵌套字段打平后再执行重命名操作 - 需保证重命名后的字段名与PostgreSQL目标表列名大小写完全匹配,PostgreSQL默认列名小写,若目标列名为大写需用双引号包裹
内容的提问来源于stack exchange,提问作者Marcel kobain
相关产品推荐
相关产品推荐

