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

Kafka Connect SpoolDir Json源连接器键提取与配置问题咨询

问题解决方案

一、无Schema连接器与Transform兼容问题

你遇到的“期望Struct但得到String”报错,核心原因是**value.converter使用了StringConverter,导致Kafka Connect将JSON内容以纯字符串形式传递,而ValueToKey和ExtractField Transform仅能处理结构化数据(如Map或Struct)**。

修改配置如下,替换字符串转换器为JSON转换器并禁用Schema:

name=SchemaLessJsonSpoolDir
tasks.max=1
connector.class=com.github.jcustenborder.kafka.connect.spooldir.SpoolDirSchemaLessJsonSourceConnector
input.path=/home/spooldirTest/json_input
input.file.pattern=^.*\.ndjson$
error.path=/home/spooldirTest/json_errors
finished.path=/home/spooldirTest/json_output
halt.on.error=false
topic=spooldir-schemaless-json-topic
# 替换为JSON转换器,禁用Schema以生成Map结构数据
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=false
transforms=extractKey,extractValue
transforms.extractKey.type=org.apache.kafka.connect.transforms.ValueToKey
transforms.extractKey.fields=key
transforms.extractValue.type=org.apache.kafka.connect.transforms.ExtractField$Value
transforms.extractValue.field=value

修改后,连接器会将每行NDJSON解析为Map类型数据,Transform就能正常提取key和value字段。

二、带Schema的连接器配置问题

1. 自动生成Schema无效(值全为null)的原因

自动生成Schema失败通常是因为JSON字段类型不一致(比如你的数据中value字段有时是数组、有时是字符串),连接器的类型推断无法兼容这种混合类型,导致解析失败,最终值为null。

解决思路:

  • 尽量统一JSON字段类型,确保同名字段的类型一致;
  • 调整连接器的类型推断配置,比如添加schema.generation.allow.missing.fields=true允许缺失字段,或schema.generation.type.inference.enabled=true强制启用类型推断。

2. 手动传入Schema的方法

不需要使用AVRO,用JSON Schema即可。具体配置步骤:

  1. 编写符合JSON Schema规范的Schema文件(比如schema.json),示例如下:
{
  "$schema": "http://json-schema.org/draft-07/schema#",
  "type": "object",
  "properties": {
    "key": {"type": "string"},
    "value": {"type": ["string", "array"], "items": {"type": "string"}}
  },
  "required": ["key"]
}
  1. 在连接器配置中指定Schema文件路径,同时关闭自动生成:
name=SchemaBasedSpoolDir
tasks.max=1
connector.class=com.github.jcustenborder.kafka.connect.spooldir.SpoolDirJsonSourceConnector
input.path=/home/spooldirTest/json_input
input.file.pattern=^.*\.ndjson$
error.path=/home/spooldirTest/json_errors
finished.path=/home/spooldirTest/json_output
halt.on.error=false
topic=spooldir-schema-json-topic
# 指定Schema文件路径
schema.path=/path/to/your/schema.json
# 关闭自动生成Schema
schema.generation.enabled=false
# 使用支持Schema的转换器
value.converter=org.apache.kafka.connect.json.JsonConverter
value.converter.schemas.enable=true

3. 关于Converter的选择

  • 如果使用Schema Registry,推荐用io.confluent.connect.avro.AvroConverter,配合Schema Registry存储和管理Schema;
  • 如果不需要Schema Registry,用org.apache.kafka.connect.json.JsonConverter并设置schemas.enable=true即可,它会将Schema嵌入到消息中。

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

相关产品推荐
方舟 Agent Plan

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

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