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

使用Kafka连接器映射Topic键值到数据库列时配置缺失报错排查

问题根源

报错No fields found using key and value schemas for table是因为JDBC Sink Connector默认会尝试将Kafka消息key/value的schema字段与数据库表列名一一匹配,但你的消息key的schema字段是MY_ID,value的schema字段是MY_ID/MY_INT/MY_NUMBER,和表的MY_KEY/MY_VALUE列完全不匹配,Connector找不到对应字段,因此报错。

你的需求是把整个key内容存入MY_KEY列,整个value内容存入MY_VALUE列,而非字段映射,因此需要通过Single Message Transforms (SMT) 转换消息结构,或调整Converter配置直接传递原始内容。


解决方案一:处理Avro格式消息(匹配当前Converter配置)

添加SMT将整个key和value包装成包含MY_KEY和MY_VALUE字段的结构体,让Connector能匹配表列:

完整配置修改

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaConnector
metadata:
  name: my-connector  # 原名称带空格,建议改为连字符格式避免资源名称问题
  labels:
    strimzi.io/cluster: kafka
spec:
  class: io.confluent.connect.jdbc.JdbcSinkConnector
  tasksMax: 1
  config:
    connection.url: <你的数据库URL>
    connection.user: <用户名>
    connection.password: <密码>
    topics: MY_TOPIC
    key.converter: io.confluent.connect.avro.AvroConverter
    value.converter: io.confluent.connect.avro.AvroConverter
    key.converter.schema.registry.url: http://schemaregistry.kafka:8085
    value.converter.schema.registry.url: http://schemaregistry.kafka:8085
    auto.create: false
    auto.evolve: false
    insert.mode: update
    table.name.format: POM_BL_LOG  # 修正为实际表名,原配置写的MY_TABLE
    # SMT配置:将整个key/value包装为对应字段
    transforms: hoistKey,hoistValue
    transforms.hoistKey.type: org.apache.kafka.connect.transforms.HoistField$Key
    transforms.hoistKey.field: MY_KEY
    transforms.hoistValue.type: org.apache.kafka.connect.transforms.HoistField$Value
    transforms.hoistValue.field: MY_VALUE
    # 更新模式必须指定主键
    pk.fields: MY_KEY
    pk.mode: record_key

关键配置说明

  • table.name.format:必须修正为实际表名POM_BL_LOG,否则Connector找不到目标表。
  • SMT转换:HoistField会把整个key包裹到MY_KEY字段、整个value包裹到MY_VALUE字段,让Connector能匹配表的对应列。
  • pk.fields/pk.mode:insert.mode: update要求指定主键,这里MY_KEY作为主键,配置后Connector会以此字段为更新依据。

解决方案二:支持非Avro格式(字符串key、完整JSON value直接存储)

如果需要处理字符串或原始JSON格式的消息,调整Converter为StringConverter并关闭schema校验:

字符串/JSON格式消息配置

apiVersion: kafka.strimzi.io/v1beta2
kind: KafkaConnector
metadata:
  name: my-connector
  labels:
    strimzi.io/cluster: kafka
spec:
  class: io.confluent.connect.jdbc.JdbcSinkConnector
  tasksMax: 1
  config:
    connection.url: <你的数据库URL>
    connection.user: <用户名>
    connection.password: <密码>
    topics: MY_TOPIC
    # 使用字符串Converter直接传递原始内容
    key.converter: org.apache.kafka.connect.storage.StringConverter
    value.converter: org.apache.kafka.connect.storage.StringConverter
    # 关闭schema校验(字符串/原始JSON无需schema)
    key.converter.schemas.enable: false
    value.converter.schemas.enable: false
    auto.create: false
    auto.evolve: false
    insert.mode: update
    table.name.format: POM_BL_LOG
    # SMT配置:将key/value包装为对应字段
    transforms: hoistKey,hoistValue
    transforms.hoistKey.type: org.apache.kafka.connect.transforms.HoistField$Key
    transforms.hoistKey.field: MY_KEY
    transforms.hoistValue.type: org.apache.kafka.connect.transforms.HoistField$Value
    transforms.hoistValue.field: MY_VALUE
    # 指定更新主键
    pk.fields: MY_KEY
    pk.mode: record_key

额外注意事项

  • 确保Kafka消息key的内容长度不超过数据库表MY_KEY列的VARCHAR2(50)限制,否则会插入失败。
  • MY_VALUE是CLOB类型,JDBC Connector会自动处理大文本内容,无需额外配置。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 00:34:58