如何用Apache Beam Python读取非JSONL格式的多行JSON文件?
遇到多行JSON文件(每个JSON对象跨多行,或者整个文件是一个JSON数组)的时候,默认的ReadFromFile逐行读取确实会失效,因为每行不是完整的JSON结构。这里有两种实用的解决方案,完全适配Beam的管道模型,比直接用FileSystem.open()更方便可靠:
方案1:处理多个独立的多行JSON对象
如果你的文件是多个独立的JSON对象(每个占多行),比如:
{
"name": "Alice",
"age": 30
}
{
"name": "Bob",
"age": 25
}
可以用ReadFromFile的read_all=True参数把整个文件内容作为单个元素读取,再配合自定义DoFn用JSON解码器逐个解析对象:
import apache_beam as beam import json class ParseMultiLineJson(beam.DoFn): def process(self, file_content): decoder = json.JSONDecoder() pos = 0 content_len = len(file_content) while pos < content_len: try: # 解析当前位置的第一个JSON对象,返回对象和新的位置 obj, pos = decoder.raw_decode(file_content, pos) yield obj except json.JSONDecodeError: # 跳过无效字符(比如换行、空格),继续往后解析 pos += 1 with beam.Pipeline() as p: parsed_json = ( p | "读取整个文件" >> beam.io.ReadFromFile("path/to/your/multi-line.json", read_all=True) | "解析多行JSON" >> beam.ParDo(ParseMultiLineJson()) # 这里可以添加后续的处理步骤,比如写入数据库、转换格式等 )
这个方法的核心是json.JSONDecoder.raw_decode(),它能从字符串中精准定位并解析完整的JSON对象,完全兼容嵌套结构,比正则分割更可靠。
方案2:处理整个文件是一个JSON数组
如果你的文件是一个大的JSON数组(所有数据包裹在[]里),比如:
[
{"name": "Alice", "age": 30},
{"name": "Bob", "age": 25}
]
处理起来更简单,直接读取整个文件后解析成数组,再展开每个元素:
import apache_beam as beam import json class ParseJsonArray(beam.DoFn): def process(self, file_content): # 把整个文件内容解析成JSON数组 json_array = json.loads(file_content) # 逐个输出数组里的元素 for item in json_array: yield item with beam.Pipeline() as p: parsed_json = ( p | "读取整个文件" >> beam.io.ReadFromFile("path/to/your/json-array.json", read_all=True) | "解析JSON数组" >> beam.ParDo(ParseJsonArray()) # 后续处理步骤 )
为什么不用FileSystem.open()?
你提到的FileSystem.open()确实能读取整个文件,但它是Beam的底层API,在分布式环境下(比如运行在Dataflow、Flink上),直接使用它需要自己处理文件分片、分布式存储访问等细节,而ReadFromFile已经封装了这些逻辑,能更好地适配Beam的并行处理模型,避免很多潜在问题。
注意事项
如果你的文件特别大(比如几十GB级别),read_all=True会把整个文件加载到内存,可能导致内存溢出。这种情况下,建议先把多行JSON预处理成JSONL(每行一个JSON对象)格式,再用常规的ReadFromFile配合json.loads处理会更高效。
内容的提问来源于stack exchange,提问作者Manuel RODRIGUEZ

