- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
JSON 是数据工程中最通用的交换格式之一。Apache Beam 为 Python SDK 提供了内置的ReadFromJson变换,让你无需手写解析逻辑即可将 JSON 文件(含 JSON Lines)读取为 PCollection,并借助 pandas 的read_json语义获得类型推断、orient 灵活适配等能力。本文以 Beam 仓库中的官方提示文档为基础,结合sdks/python/apache_beam/io/textio.py与sdks/python/apache_beam/dataframe/io.py的真实实现,讲解ReadFromJson的完整用法、关键参数语义、底层原理与常见实战场景,读完后你可以直接写出可运行、可上生产环境的 JSON 读取管道。
一、核心示例:完整可运行的 JSON 读取管道
以下代码来自仓库中的 官方提示文档,它演示了如何用 Apache Beam Python SDK 从 JSON 文件读取数据:
import logging import apache_beam as beam from apache_beam import Map from apache_beam.io.textio import ReadFromJson from apache_beam.options.pipeline_options import PipelineOptions class JsonOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_argument( '--file_path', default="gs://your-bucket/your-file.json", help='Json file path' ) def run(): """ This pipeline reads from Json file defined by the --file_path argument. """ options = JsonOptions() with beam.Pipeline(options=options) as p: output = p | "Read from Json file" >> ReadFromJson( path=options.file_path, lines=False ) if __name__ == "__main__": logging.getLogger().setLevel(logging.INFO) run()运行方式(支持本地 Direct Runner 与远程 runner):
python read_json_pipeline.py --file_path gs://your-bucket/your-file.json python read_json_pipeline.py --file_path ./local/data.json # 本地文件同样支持逐段拆解
JsonOptions类:它继承自PipelineOptions,这是 Beam Python SDK 提供的命令行参数解析基类。通过覆写_add_argparse_args,我们可以自定义--file_path参数,并设置默认值与帮助文本。这样管道就能通过命令行注入要读取的 JSON 文件路径,而无需硬编码。代码中的options.file_path即对应命令行传入的--file_path。beam.Pipeline(options=options) as p:以with语句创建管道。退出with块时管道会自动运行,这是 Beam 2.x 中推荐的上下文管理用法。ReadFromJson(path=..., lines=False):内置的 JSON 读取变换,负责把 JSON 文件内容转成 PCollection。这里特意指定lines=False,表示整个文件作为一个整体解析(详见下文参数详解)。
二、ReadFromJson 变换的源码级解析
ReadFromJson定义在 sdks/python/apache_beam/io/textio.py 中,它并不是从零实现的读取器,而是对 pandasread_json的封装:
@append_pandas_args(pandas.read_json, exclude=['path_or_buf']) def ReadFromJson( path: str, *, orient: str = 'records', lines: bool = True, dtype: Union[bool, dict[str, Any]] = False, **kwargs): """A PTransform for reading json values from files into a PCollection. Args: path (str): The file path to read from. The path can contain glob characters such as ``*`` and ``?``. orient (str): Format of the json elements in the file. Default to 'records', meaning the file is expected to contain a list of json objects like `{field1: value1, field2: value2, ...}`. lines (bool): Whether each line should be considered a separate record, as opposed to the entire file being a valid JSON object or list. Defaults to True (unlike Pandas). dtype (bool): If True, infer dtypes; if a dict of column to dtype, then use those; if False, then don’t infer dtypes at all. Defaults to False (unlike Pandas). **kwargs: Extra arguments passed to `pandas.read_json` (see below). """ from apache_beam.dataframe.io import ReadViaPandas return 'ReadFromJson' >> ReadViaPandas( 'json', path, orient=orient, lines=lines, dtype=dtype, **kwargs)参数语义详解
| 参数 | 默认值 | 说明 |
|---|---|---|
path | 必填 | 要读取的文件路径,支持 glob 通配符(如*、?),因此可以一次读取一个目录下的多个文件(如data/out*) |
orient | 'records' | JSON 元素的组织格式。'records'表示文件内容是一个 JSON 对象列表,形如[{field1: value1, field2: value2, ...}, ...];还支持 pandasread_json的其它取值(如'columns'、'index'、'split'、'table'、'values'等) |
lines | True | 是否将每一行视为一条独立记录。True对应 JSON Lines(NDJSON)格式;False表示整个文件是一个合法的 JSON 对象或数组 |
dtype | False | 数据类型处理策略:True时自动推断类型;传入{列名: 类型}字典时按指定类型读取;False时完全不推断类型。注意默认值与 pandas 不同(pandas 默认会推断) |
**kwargs | — | 透传给pandas.read_json的其它参数(如nrows之外的 pandas 参数)。注意nrows在 Beam 的封装中暂不支持 |
提示:原文档示例中显式传了
lines=False,而源码默认值是lines=True。二者都是合法用法,取决于你的源文件格式——如果文件是标准的 JSON 数组(每行不是独立的 JSON),应使用lines=False;如果是 JSON Lines 流式文件,保留默认lines=True即可。
底层调用链
从源码可以看到,ReadFromJson的实际执行路径是:
ReadFromJson构造ReadViaPandas('json', path, ...)变换;ReadViaPandas(见 sdks/python/apache_beam/dataframe/io.py)在__init__中通过globals()['read_json']找到对应的读取器;read_json(见 sdks/python/apache_beam/dataframe/io.py)最终调用pd.read_json,并将其包装进_ReadFromPandas:- 当
lines=True时,使用_DelimSplitter(b'\n', _DEFAULT_BYTES_CHUNKSIZE)以换行符作为记录边界做增量读取(incremental 模式),这也正是 JSON Lines 文件能够被高效、分布式读取的底层机制; - 当
lines=False时按整文件读取。
- 当
ReadViaPandas.expand把读取结果转成 DataFrame 后,对object类型的列统一转换为pd.StringDtype()(objects_as_strings=True),再通过convert.to_pcollection将 DataFrame 转换为 Beam PCollection。
这意味着ReadFromJson输出的每个元素实际上是一行记录的命名元组(row tuple),字段名与 JSON 对象的键一一对应,你可以直接通过属性访问。
三、数据格式适配:orient 与 lines 的搭配
场景一:标准 JSON 数组文件
文件内容形如:
[ {"name": "alice", "age": 30}, {"name": "bob", "age": 25} ]应使用orient='records'(默认)配合lines=False:
pcoll = p | ReadFromJson(path="gs://bucket/array.json", lines=False)场景二:JSON Lines(NDJSON)文件
文件每行一条记录,适合流式追加和分布式分片读取:
{"name": "alice", "age": 30} {"name": "bob", "age": 25}使用默认的lines=True即可。此时 Beam 的_DelimSplitter会按换行符切分记录,天然支持多 worker 并行读取不同分片。
场景三:按 glob 读取多个文件
path参数支持 glob 模式,一次读取一批文件:
pcoll = p | ReadFromJson(path="data/2024/*.json", lines=True)该能力在 textio_test.py 的读写往返测试中得到了验证:测试先用WriteToJson写出分片文件(形如out-00000-of-00003),再通过ReadFromJson(os.path.join(dest, 'out*'))读取并断言与原始beam.Row完全一致。
四、类型保持:数值与字符串不被误转换
JSON 读取最容易踩的坑是数值与数字字符串被混淆。仓库测试 test_numeric_strings_preserved 专门验证了这一点:写出as_string="0"、as_float_string="0.0"、as_int=0、as_float=0.0四种类型的数据后,再读回并逐字段断言类型完全一致。这得益于ReadViaPandas.expand中对object类型列统一转换为pd.StringDtype()的处理——字符串列不会被 pandas 自动解析成数值。
如果你需要更精确的类型控制,可以直接使用dtype参数指定列类型:
pcoll = p | ReadFromJson( path="data.json", dtype={"age": int, "name": str} )五、从 ReadFromJson 到完整管道:一个加工示例
读取 JSON 只是起点。结合 Beam 的常规变换,可以立刻对数据进行加工。下面是一个读取 JSON 数组并统计年龄平均值的完整示例:
import logging import apache_beam as beam from apache_beam.io.textio import ReadFromJson from apache_beam.options.pipeline_options import PipelineOptions class JsonOptions(PipelineOptions): @classmethod def _add_argparse_args(cls, parser): parser.add_argument('--file_path', help='Json file path') parser.add_argument('--output', help='Output path prefix') def run(): options = JsonOptions() with beam.Pipeline(options=options) as p: records = ( p | "Read from Json" >> ReadFromJson( path=options.file_path, lines=False) ) avg_age = ( records | beam.Map(lambda row: row.age) | beam.CombineGlobally(beam.combiners.MeanCombineFn()) ) _ = avg_age | "Write result" >> beam.io.WriteToText(options.output) if __name__ == "__main__": logging.getLogger().setLevel(logging.INFO) run()由于ReadFromJson输出的是带字段名的行对象,row.age可以直接按字段名访问,无需再手写json.loads解析。
六、配套写入:WriteToJson 实现读写闭环
与ReadFromJson配对的是WriteToJson(同样定义在 textio.py),它把 PCollection 写出为 JSON 文件。它的关键参数与读取端对称:
path:输出文件前缀,配合num_shards与file_naming决定最终分片文件名(默认形如path-XXXXX-of-NNNNN);orient='records':默认输出为 JSON 对象数组;lines:None时自动取orient == 'records'的结果——即 records 格式下默认输出为 JSON Lines;num_shards:分片数量,None表示由系统自动选择最优值。
读写配对的完整闭环在 JsonTest.test_json_read_write 中演示:beam.Create(records) | WriteToJson(dest/'out')写出,再用ReadFromJson(dest/'out*')读回并断言相等。
七、注意事项与限制
- 依赖要求:
ReadFromJson依赖 pandas 的read_json,因此需要安装apache_beam[dataframe]附加依赖。源码中对此有兜底处理——如果 pandas 不可用,textio.py 会将这些变换替换为no_pandas函数并抛出ImportError('Please install apache_beam[dataframe]')。 nrows限制:read_json封装(见 dataframe/io.py)对nrows参数明确抛出NotImplementedError,即暂不支持只读取前 N 行。- 运行环境:示例默认未指定 runner,会在本地 Direct Runner 上执行;如需部署到 Dataflow、Spark 或 Flink 等分布式 runner,通过
--runner与对应--project/--region等参数指定即可。 - 文件可读性:
path可以是本地路径、GCS 路径(gs://)等 Beam 文件系统支持的任意位置,glob 通配符同样适用。
总结
ReadFromJson是 Apache Beam Python SDK 中读取 JSON 数据的标准入口:它以path+orient+lines+dtype四个核心参数覆盖了绝大多数 JSON 文件形态,底层通过 pandasread_json与ReadViaPandas封装获得类型推断、换行增量读取和 DataFrame 到 PCollection 的自动转换。配合WriteToJson可以实现完整的 JSON 读写闭环,而textio_test.py中的往返测试则保证了数值/字符串类型在读写过程中不被破坏。将本文示例中的JsonOptions与ReadFromJson组合使用,即可快速搭建生产级的 JSON 数据接入管道。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Python SDK 从 JSON 文件读取数据:ReadFromJson 变换实战详解
Apache Beam Python SDK 从 JSON 文件读取数据:ReadFromJson 变换实战详解 导读 本文围绕 Apache Beam Pyt
大数据批处理流处理数据工程Apache Beam Python SDK 从 Kafka 读取数据:ReadFromKafka 变换实战指南
Apache Beam Python SDK 从 Kafka 读取数据:ReadFromKafka 变换实战指南 Apache Beam 提供了统一的批流处理编
大数据批处理流处理数据工程Apache Beam Python SDK 实战:使用 ReadFromAvro 与 PipelineOptions 读取 Avro 文件
Apache Beam Python SDK 实战:使用 ReadFromAvro 与 PipelineOptions 读取 Avro 文件 导读 本文以 Ap
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考