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

