Kafka-Connect无法发送无Schema JSON至PostgreSQL的问题求助
解决Kafka Connect JdbcSinkConnector写入PostgreSQL的Schema错误
问题根源
你遇到的错误是因为JdbcSinkConnector默认要求数据带有结构化Schema,但Zabbix发送的是无Schema的原始JSON(以HashMap格式被Connect读取),且当前连接器配置未正确适配这种无Schema数据。即使全局设置了CONNECT_VALUE_CONVERTER_SCHEMAS_ENABLE=false,如果连接器自身未明确配置转换器参数,全局配置可能不会被继承。
正确配置方案
需要在连接器配置中明确指定JSON转换器并禁用Schema,同时调整主键模式以适配无Schema数据:
修改后的连接器创建命令
curl -X POST "http://localhost:8082/connectors" -H "Content-Type: application/json" -d '{ "name": "zabbix-sink-connector", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "tasks.max": 1, "topics": "zabbix-webhook", "connection.url": "jdbc:postgresql://<host>/etl?currentSchema=test", "connection.user": "<username>", "connection.password": "<password>", "auto.create": "true", "auto.evolve": "true", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "value.converter.schemas.enable": "false", "pk.mode": "record_value", "pk.fields": "itemid", "insert.mode": "upsert" } }'
关键配置说明
value.converter与value.converter.schemas.enable:强制连接器使用JSON转换器处理无Schema的JSON数据,覆盖全局配置可能存在的不一致。pk.mode与pk.fields:指定从数据中提取itemid作为主键,满足JdbcSink对数据唯一性的要求(避免默认的pk.mode=none带来的Schema依赖)。auto.evolve:自动适配PostgreSQL表结构,当Zabbix数据字段发生变化时无需手动修改表。insert.mode=upsert:如果同一itemid有更新数据,会自动执行更新操作而非插入重复记录。
额外检查点
- 确保Kafka Connect服务已正确加载
JsonConverter(默认已包含在Connect分发包中)。 - 若需要清理错误数据或重置处理位置,可执行以下命令重置连接器偏移量:
curl -X POST "http://localhost:8082/connectors/zabbix-sink-connector/offsets" -H "Content-Type: application/json" -d '{"offset": 0}'
内容的提问来源于stack exchange,提问作者Dmitry
相关产品推荐
相关产品推荐

