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

Gobblin Kafka JSON转Avro失败:非JSON数组错误求助

解决Gobblin读取Kafka JSON转Avro到HDFS的Schema错误

首先看你的错误栈,核心问题是JsonStringToJsonIntermediateConverter在转换Schema时,发现输入的Schema不是JSON数组,导致抛出IllegalStateException: This is not a JSON Array。结合你的配置和Kafka消息,我整理了几个关键问题和修复方案:

1. 配置中的明显错误

a. 重复的extract.namespace

你配置了两次extract.namespace,后面的gobblin.source.extractor.filebased会覆盖前面的Kafka相关命名空间,这可能导致Schema解析逻辑混乱,保留一个即可:

extract.namespace=org.apache.gobblin.extract.kafka

b. 不匹配的Writer输出格式

你用了AvroDataWriterBuilder但设置writer.output.format=text,这完全矛盾,Avro Writer需要对应Avro格式:

writer.output.format=avro

c. source.schema的格式问题(核心错误原因)

在Properties配置文件中,字符串里的双引号需要转义,否则配置解析器会截断你的Schema字符串,导致Gson无法解析成JSON数组。另外你还拼写错了字段名ubdated_at(应该是updated_at),修正后的Schema配置:

source.schema=[{\"columnName\":\"name\", \"dataType\":{\"type\": \"string\"}}, {\"columnName\":\"city\", \"dataType\":{\"type\": \"string\"}}, {\"columnName\":\"age\", \"dataType\":{\"type\": \"integer\"}}, {\"columnName\":\"updated_at\", \"dataType\":{\"type\": \"string\"}}]

d. 多余的Schema注入配置

gobblin.converter.schemaInjector.schema=SCHEMA这个配置是多余的,它会干扰JsonStringToJsonIntermediateConverter对source.schema的读取,直接移除即可。

2. 修正后的完整配置

job.name=GobblinKafkaQuickStart
job.group=GobblinKafka
job.description=Gobblin quick start job for Kafka
job.lock.enabled=false
kafka.brokers=localhost:9092
source.class=org.apache.gobblin.source.extractor.extract.kafka.KafkaSimpleSource
extract.namespace=org.apache.gobblin.extract.kafka
converter.classes=org.apache.gobblin.converter.json.JsonStringToJsonIntermediateConverter, org.apache.gobblin.converter.avro.JsonIntermediateToAvroConverter
source.schema=[{\"columnName\":\"name\", \"dataType\":{\"type\": \"string\"}}, {\"columnName\":\"city\", \"dataType\":{\"type\": \"string\"}}, {\"columnName\":\"age\", \"dataType\":{\"type\": \"integer\"}}, {\"columnName\":\"updated_at\", \"dataType\":{\"type\": \"string\"}}]
writer.builder.class=org.apache.gobblin.writer.AvroDataWriterBuilder
writer.destination.type=HDFS
writer.output.format=avro
data.publisher.type=org.apache.gobblin.publisher.BaseDataPublisher
mr.job.max.mappers=1
metrics.reporting.file.enabled=true
metrics.log.dir=${gobblin.cluster.work.dir}/metrics
metrics.reporting.file.suffix=txt
bootstrap.with.offset=earliest

3. 额外检查点

  • 确认Kafka主题中的消息确实是你提供的格式:{"age": 36, "city": "London", "name": "John", "updated_at": "2020-05-19"}(注意修正后的字段名)
  • 确认Gobblin运行环境能访问到Kafka集群和HDFS,权限配置正确
  • 如果你用的是较新版本的Gobblin,确保JsonStringToJsonIntermediateConverter和JsonIntermediateToAvroConverter的类路径正确(这些类在gobblin-converter-json和gobblin-converter-avro模块中)

按照这个配置重新运行独立模式,应该就能解决Schema解析的错误了。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.07 18:18:02