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即可。具体配置步骤:
- 编写符合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"] }
- 在连接器配置中指定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
相关产品推荐
相关产品推荐

