使用Apache Beam Python读取GCS多行JSON文件报错求助
解决Apache Beam读取GCS多行JSON数组的问题
问题原因
beam.io.ReadFromText默认按行读取文件,将每一行作为单独的PCollection元素传入后续处理。你的JSON是多行数组结构,第一行仅为[,单独传给json.loads时会触发JSON解码错误——因为单个[不是合法的JSON值。而单行JSON是完整的数组字符串,能被正常解析。
解决方案
方法1:合并所有行后解析
通过CombineGlobally将所有行拼接成完整的JSON字符串,再进行解析:
import json import apache_beam as beam from apache_beam.transforms.combiners import CombineGlobally with beam.Pipeline() as p: processed = ( p | "读取GCS文件" >> beam.io.ReadFromText(bucket) | "合并所有行" >> CombineGlobally(lambda lines: ''.join(lines)) | "解析完整JSON" >> beam.Map(json.loads) | "展开数组元素" >> beam.FlatMap(lambda arr: arr) # 如需单独处理数组内每个对象 | "打印结果" >> beam.Map(print) )
方法2:读取整个文件为单个字符串(单文件场景)
使用ReadAllFromText,通过min_bundle_size参数强制读取整个文件作为单个元素:
import json import apache_beam as beam with beam.Pipeline() as p: processed = ( p | "创建文件路径" >> beam.Create([bucket]) | "读取完整文件" >> beam.io.ReadAllFromText(min_bundle_size=1024*1024) # 设置足够大的bundle size | "解析JSON" >> beam.Map(json.loads) | "展开数组元素" >> beam.FlatMap(lambda arr: arr) | "打印结果" >> beam.Map(print) )
额外注意
你的原JSON中最后一个键值对后有多余逗号("entity": "Comunità del Collio",),虽然部分工具能兼容,但建议删除该逗号,避免潜在的解析异常。
内容的提问来源于stack exchange,提问作者Aurora
相关产品推荐
相关产品推荐

