You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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有更新数据,会自动执行更新操作而非插入重复记录。

额外检查点

  1. 确保Kafka Connect服务已正确加载JsonConverter(默认已包含在Connect分发包中)。
  2. 若需要清理错误数据或重置处理位置,可执行以下命令重置连接器偏移量:
curl -X POST "http://localhost:8082/connectors/zabbix-sink-connector/offsets" -H "Content-Type: application/json" -d '{"offset": 0}'

内容的提问来源于stack exchange,提问作者Dmitry

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.09 18:10:49