news 2026/9/29 20:49:35

使用 Apache Beam Python SDK 读取 JSON 文件:ReadFromJson 变换实战指南

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
使用 Apache Beam Python SDK 读取 JSON 文件:ReadFromJson 变换实战指南
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

导读

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'等)
linesTrue是否将每一行视为一条独立记录。True对应 JSON Lines(NDJSON)格式;False表示整个文件是一个合法的 JSON 对象或数组
dtypeFalse数据类型处理策略: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的实际执行路径是:

  1. ReadFromJson构造ReadViaPandas('json', path, ...)变换;
  2. ReadViaPandas(见 sdks/python/apache_beam/dataframe/io.py)在__init__中通过globals()['read_json']找到对应的读取器;
  3. 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时按整文件读取。
  4. 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*')读回并断言相等。

七、注意事项与限制

  1. 依赖要求:ReadFromJson依赖 pandas 的read_json,因此需要安装apache_beam[dataframe]附加依赖。源码中对此有兜底处理——如果 pandas 不可用,textio.py 会将这些变换替换为no_pandas函数并抛出ImportError('Please install apache_beam[dataframe]')。
  2. nrows限制:read_json封装(见 dataframe/io.py)对nrows参数明确抛出NotImplementedError,即暂不支持只读取前 N 行。
  3. 运行环境:示例默认未指定 runner,会在本地 Direct Runner 上执行;如需部署到 Dataflow、Spark 或 Flink 等分布式 runner,通过--runner与对应--project/--region等参数指定即可。
  4. 文件可读性: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.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/29 20:48:48

配电柜温湿度监控为何必须选RJ45以太网传感器

1. 为什么配电柜监控非要选RJ45以太网温湿度传感器——从“能用”到“必须用”的现场逻辑去年夏天,我接手一个老电厂的智能升级项目,目标很朴素:把32台高压配电柜的环境数据实时传回中控室。最初方案是用USB转串口RS485温湿度模块&#xff0c…

作者头像 李华
网站建设 2026/9/29 20:48:46

以太网变送器双协议批量配置实战指南

1. 为什么“双协议批量配置”不是锦上添花,而是大规模环境监测项目的生死线我接手过三个超500个点位的工业级环境监测项目,最深的体会是:设备部署完成≠系统可用。真正卡住交付进度、拖垮运维成本的,从来不是传感器精度或外壳防护…

作者头像 李华
网站建设 2026/9/29 20:47:33

ESP32上WebAssembly实战:.wasm为何不能当应用跑?

前几天有个朋友找我聊 ESP32 上的 WebAssembly 方案,开口就是一句:“我逻辑用 Rust 编成 .wasm 了,是不是可以直接烧进去当应用跑?”我愣了一下,然后意识到这不是他一个人的困惑。最近各种技术社群里,“把应…

作者头像 李华
网站建设 2026/9/29 20:47:27

工业相机选型实战:从分辨率、靶面到全局快门的完整指南

1. 选型先别看参数表,先搞清你的成像需求我收到过很多类似的私信:把某个相机的型号往我这儿一丢,问"这个能不能用来检测 PCB 焊点""这个适合做 OCR 吗"。说实话,这种问法很难回答,因为参数表是死的…

作者头像 李华
网站建设 2026/9/29 20:46:26

MQTT协议架构与工业落地:从发布订阅、QoS到Broker实战

先说一个真实场景:你走进一家智能工厂,几十台PLC、传感器、AGV小车在车间里跑,中控大屏上温度、振动、产量、设备状态实时跳动。这套数据采集和指令下发背后,用的通信协议十有八九就是MQTT。我当年第一次接触MQTT时,第…

作者头像 李华
网站建设 2026/9/29 20:46:18

800人园区网实战:从毕业论文到企业网规划与设计全流程

简介:这份毕业论文文档围绕旭日公司企业网络的规划与设计展开,面向网络工程、通信工程等专业的在校学生及需要撰写同类课题的从业者,可作为毕业设计选题参考与方案撰写范本。文档以企业网络建设为背景,系统梳理了从需求分析到落地…

作者头像 李华