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

使用Kafka Connect JDBC Sink连接Kafka主题至PostgreSQL时遇错误求助

问题排查:Kafka JDBC Sink Connector 连接 PostgreSQL 失败

配置中的明显错误

先看你提供的连接器配置,存在几处拼写和逻辑问题:

  1. 拼写错误

    • transfores.Flatten.type 应为 transforms.Flatten.type(多了一个字母e)
    • transforms.Flatten.deliniter 应为 transforms.Flatten.delimiter(单词拼写错误)
  2. Transform 逻辑冗余

    • transforms.RenameFields.renames 设置为value:value,timestamp:timestamp,相当于没有重命名字段,完全可以移除这个Transform,减少不必要的处理环节。
  3. Upsert 模式与主键配置冲突

    • 你设置了insert.mode: "upsert",但同时配置pk.mode: "none"。Upsert(更新插入)需要指定主键来判断操作类型,两者冲突会导致连接器无法正常执行upsert逻辑。如果不需要更新操作,改成insert.mode: "insert";如果需要upsert,必须配置pk.mode(比如pk.mode: "record_value")并指定pk.fields为目标表的主键字段。

补充排查点(基于JDBC Sink常见报错场景)

结合这类问题的常规排查逻辑,还需要确认以下内容:

  • Schema Registry 连通性:确认Schema Registry服务(http://localhost:8081)正常运行,且连接器所在节点能访问该地址。使用Avro格式时,Schema Registry不可用会直接导致消息解析失败。
  • PostgreSQL 表结构匹配:确认temperature表的字段名称、数据类型,与Kafka消息经过Transform后的字段完全兼容。比如消息中的timestamp字段类型,是否与表中对应字段的TIMESTAMP类型匹配。
  • 数据库权限与连通性:确认postgres用户拥有jdbcsink数据库中temperature表的插入/更新权限,同时检查jdbc:postgresql://localhost:5432/jdbcsink地址正确、PostgreSQL服务正常运行。
  • 批量配置合理性:batch.size: "2"设置过小,会导致频繁的数据库请求,建议调整为合理数值(比如默认的1000),减少资源开销。

修正后的示例配置

{
    "name": "temperature_jdbcsink",
    "config": {
        "connector.class": "io.confluent.connect.jdbc.JdbcSinkConnector",
        "task.max": "1",
        "topics": "temperature",
        "key.converter": "io.confluent.connect.avro.AvroConverter",
        "value.converter": "io.confluent.connect.avro.AvroConverter",
        "key.converter.schema.registry.url": "http://localhost:8081",
        "value.converter.schema.registry.url": "http://localhost:8081",
        "transforms": "Flatten",
        "transforms.Flatten.type": "org.apache.kafka.connect.transforms.Flatten$value",
        "transforms.Flatten.delimiter": "_",
        "connection.url": "jdbc:postgresql://localhost:5432/jdbcsink",
        "connection.user": "postgres",
        "connection.password": "postgres",
        "insert.mode": "insert",
        "batch.size": "1000",
        "table.name.format": "temperature",
        "pk.mode": "none",
        "db.timezone": "Asia/Kolkata"
    }
}

(注:如果需要启用upsert模式,需添加主键配置,例如目标表主键为id,则补充pk.mode: "record_value"和pk.fields: "id")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.13 22:50:26