Apache Beam DataFrame读取GZIP压缩JSON文件报错,实现是否正确?
问题分析与解决方案
你的实现存在兼容性问题,直接给read_json传compression_type="gzip"会触发Beam DataFrame与Pandas压缩文件处理逻辑的冲突,这就是报错AttributeError: 'CompressedFile' object has no attribute 'writable'的原因——Beam提供的压缩文件对象是只读的,而Pandas底层逻辑会尝试检查文件是否可写,导致属性不存在的错误。
正确实现方式
不要直接用read_json的compression_type参数,而是改用Beam原生的ReadFromText读取压缩文件,再转换为DataFrame:
import apache_beam as beam import pandas as pd from apache_beam.dataframe.convert import to_dataframe with beam.Pipeline() as pipeline: # 用Beam原生IO读取gzip压缩的JSON行文件,原生支持压缩格式 raw_json_lines = pipeline | beam.io.ReadFromText(input_file, compression_type="gzip") # 将读取到的行数据转换为Beam DataFrame beam_df = to_dataframe(raw_json_lines) # 解析每行JSON内容(如果需要结构化DataFrame) parsed_dataframe = beam_df.apply(lambda line: pd.read_json(line, lines=True), axis=1) # 后续可继续处理parsed_dataframe
补充说明
Beam DataFrame的read_json是对Pandasread_json的封装,但在分布式场景下,Pandas的压缩文件处理逻辑无法适配Beam的文件IO对象。改用Beam原生的ReadFromText可以绕过这个问题,因为它原生支持gzip、bz2等多种压缩格式,且能正确处理分布式文件系统中的压缩文件。
内容的提问来源于stack exchange,提问作者goutham
相关产品推荐
相关产品推荐

