Kafka JDBC Sink插入PG表报Struct字段名未正确指定异常如何解决
问题根因
- 配置项拼写错误:
value.converter.schema.enable参数名错误,正确应为value.converter.schemas.enable(少了末尾的s),导致JsonConverter无法正确解析消息中携带的Schema结构,直接触发字段名解析异常 - 无效/错误配置项:
pk.mofr为无效配置,属于拼写错误,且当前pk.mode设为none无需相关主键配置,直接删除即可mode: bulk是JDBC Source连接器专属配置,Sink连接器无此参数,直接删除key.converter配置值末尾多了多余空格,会导致类加载失败,需删除尾部空格- 存在重复配置:同时配置了
schemas.enable和分转换器的schema开关,建议删除全局的schemas.enable避免优先级冲突
- 网络配置问题:Kafka Connect运行在Docker容器中,
connection.url中的localhost指向容器内部而非宿主机,无法访问宿主机上的PostgreSQL服务 - 库表前置条件缺失:需提前在PostgreSQL的
postgres库下创建TESTDBschema,否则自动建表会因schema不存在失败
修复步骤
1. 提前创建PostgreSQL Schema
连接到PostgreSQL的postgres库,执行以下SQL:
CREATE SCHEMA IF NOT EXISTS TESTDB;
2. 修正连接器配置
使用修正后的配置创建连接器:
curl -X POST http://localhost:8082/connectors -H "Content-Type: application/json" -d '{ "name": "jdbc_sink_postgres_022", "config": { "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector", "connection.url": "jdbc:postgresql://<宿主机IP/同网络PostgreSQL服务名>:5432/postgres", "connection.user": "postgres", "connection.password": "postgres", "topics": "kafka_subs", "auto.create": "true", "insert.mode": "insert", "table.name.format": "TESTDB.subs", "pk.mode": "none", "poll.interval.ms": 60000, "value.converter.schemas.enable": "true", "key.converter.schemas.enable": "true", "value.converter": "org.apache.kafka.connect.json.JsonConverter", "key.converter": "org.apache.kafka.connect.json.JsonConverter" } }'
注意:将
<宿主机IP/同网络PostgreSQL服务名>替换为实际可访问的PostgreSQL地址,如果PostgreSQL也在同一个Docker Compose集群中,直接填服务名即可。
3. 发送测试消息
你原有消息格式符合带Schema的JsonConverter要求,无需修改,直接发送即可正常写入PostgreSQL的TESTDB.subs表。
内容的提问来源于stack exchange,提问作者Suvendu Ghosh
相关产品推荐
相关产品推荐

