使用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
相关产品推荐
相关产品推荐

