ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Beam Python SDK 从 JSON 文件读取数据:ReadFromJson 变换实战详解

Apache Beam Python SDK 从 JSON 文件读取数据:ReadFromJson 变换实战详解 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载导读本文围绕 Apache Beam Python SDK 内置的ReadFromJson变换系统讲解如何通过 PipelineOptions 自定义命令行参数、从本地或云端如 GCS读取 JSON 文件并以linesFalse模式将整个文件解析为单个 JSON 对象。文中以learning/prompts/code-explanation/07_io_json.md中的示例代码为骨架结合sdks/python/apache_beam/io/textio.py与sdks/python/apache_beam/dataframe/io.py的源码实现和sdks/python/apache_beam/io/textio_test.py中的测试用例深入说明orient、lines、dtype等核心参数的底层行为帮助读者写出可运行、可上生产环境的 JSON 读取流水线。一、示例代码总览三行看懂 JSON 读取流水线关联文档learning/prompts/code-explanation/07_io_json.md给出的核心示例是一段典型的“参数定义 流水线构建”代码class JsonOptions(PipelineOptions): classmethod def _add_argparse_args(cls, parser): parser.add_argument( --file_path, defaultgs://your-bucket/your-file.json, helpJson file path ) options JsonOptions() with beam.Pipeline(optionsoptions) as p: output (p | Read from Json file ReadFromJson( pathoptions.file_path, linesFalse ) | Log Data Map(logging.info))整段代码分为三层职责JsonOptions自定义参数类继承PipelineOptions并重写_add_argparse_args向命令行解析器注册--file_path参数默认指向gs://your-bucket/your-file.json参数实例化options JsonOptions()在 Python 解释器中构造默认配置在命令行运行时则自动解析--file_path...传入的值流水线执行beam.Pipeline(optionsoptions)携带参数构建流水线用ReadFromJson读取 JSON再用Map(logging.info)将每条记录打印到日志。下面依次拆解这三层并结合源码说明其底层机制。二、JsonOptions用 PipelineOptions 声明式管理文件路径2.1 PipelineOptions 与 _add_argparse_args 的约定PipelineOptions是 Apache Beam Python SDK 提供的命令行参数基类。子类只需实现_add_argparse_args(cls, parser)类方法在其中调用parser.add_argument(...)Beam 的选项解析机制就会自动将该参数合并到流水线可识别的参数集合中。示例中注册的参数解析规则参数名--file_path命令行传入时写作--file_pathgs://my-bucket/data.json默认值gs://your-bucket/your-file.json保证未显式传参时流水线也能有一个可解析的路径帮助文本Json file path在--help输出中提示用途。2.2 两种参数传入方式options JsonOptions()的妙处在于它同时支持两种使用场景脚本内硬编码直接实例化使用类中定义的默认值命令行覆盖通过python my_pipeline.py --file_pathgs://my-bucket/2024.json运行时Beam 自动将命令行值解析进options.file_path属性。随后在流水线中通过options.file_path访问该值将“参数声明”与“业务逻辑”解耦是 Beam Python 流水线中管理文件路径、表名、连接串等配置的推荐做法。三、ReadFromJson 变换核心签名与参数语义ReadFromJson定义于 sdks/python/apache_beam/io/textio.py 的apache_beam.io模块内其完整签名为def ReadFromJson( path: str, *, orient: str records, lines: bool True, dtype: Union[bool, dict[str, Any]] False, **kwargs):3.1 关键参数逐一说明参数默认值语义path必填要读取的文件路径支持 glob 通配符如*、?可一次匹配多个分片文件orientrecordsJSON 元素在文件中的组织格式默认records表示文件内容是一组形如{field1: value1, field2: value2, ...}的 JSON 对象列表linesTrue是否把每一行视为一条独立记录。True表示逐行解析每行一个 JSON 对象False表示将整个文件作为一个合法的 JSON 对象或数组解析。注意Beam 的默认值与 pandas 不同pandas 默认FalsedtypeFalse类型推断策略True时自动推断列类型传入{列名: 类型}字典时按列指定类型False时完全不推断类型。默认值与 pandas 不同pandas 默认True**kwargs—其余参数透传给 pandas 的read_json3.2 linesFalse 的含义整文件单对象解析示例代码特意设置linesFalse。这意味着ReadFromJson不会按行切分文件而是把整个文件当作一个完整的 JSON 对象或数组来解析——典型的场景是文件中是一个大的 JSON 数组[ {id: 1, name: Alice, score: 95}, {id: 2, name: Bob, score: 87} ]而linesTrue对应的则是“JSON Lines / NDJSON”格式每行一条独立记录{id: 1, name: Alice, score: 95} {id: 2, name: Bob, score: 87}从源码注释textio.py#L1053-L1054可以看到lines的语义正是“每条记录占一行”与“整个文件是一个合法 JSON 对象或列表”两种解析方式的切换开关。生产环境中流式追加日志数据通常用linesTrue可增量读取而静态导出的完整快照文件则常用linesFalse。3.3 输出形态schema 化的 PCollectionReadFromJson的输出不是原始 JSON 字符串而是经过 pandas 解析后、再转换为 Beam schema 元素的PCollection。测试 textio_test.py#L1837-L1851 给出了清晰的验证写入beam.Row(astr, bix)的记录读出后用zip(type(t)._fields, t)重建beam.Row与原始记录完全相等。这意味着下游可以像操作结构化数据一样直接通过字段名访问或使用Map将元素转换为字典、NamedTuple 等。四、源码深处ReadFromJson 的 pandas 驱动实现4.1 变换本身是薄封装阅读 textio.py#L1061-L1063 可见ReadFromJson内部只是ReadViaPandas(json, path, orient..., lines..., dtype...)的一层包装真正的读取逻辑落在 sdks/python/apache_beam/dataframe/io.py 中def read_json(path, *args, **kwargs): if nrows in kwargs: raise NotImplementedError(nrows not yet supported) elif kwargs.get(lines, False): # Work around https://github.com/pandas-dev/pandas/issues/34548. kwargs dict(kwargs, nrows1 63) return _ReadFromPandas( pd.read_json, path, args, kwargs, incrementalkwargs.get(lines, False), splitter_DelimSplitter(b\n, _DEFAULT_BYTES_CHUNKSIZE) if kwargs.get(lines, False) else None, binaryFalse)从源码结构可以提炼出三个重要的实现事实不支持nrows传入nrows会直接抛出NotImplementedError若linesTrue内部会用nrows1 63绕开 pandas issue #34548 的分块读取缺陷linesTrue走增量路径按\n分隔符_DelimSplitter切分数据流实现逐行增量解析支持大文件的分片与流式处理linesFalse走整体路径不启用分隔符由 pandas 一次性解析整个文件内容。4.2 ReadViaPandas把 DataFrame 转成 PCollectionReadFromJson最终落到 io.py#L806-L829 的ReadViaPandas.expand先让 pandas 读出 DataFrame再把 dtype 为object的列转为pd.StringDtype()对象序列化策略最后调用convert.to_pcollection(df)将 DataFrame 转换为 Beam PCollection。这也是为什么dtype参数的设置会直接影响下游拿到的元素类型。4.3 没有 pandas 时的降级行为textio.py#L1110-L1116 表明当环境缺少 pandas 时ReadFromJson/WriteToJson等变换会被替换为一个直接抛出ImportError(Please install apache_beam[dataframe])的占位函数。因此使用 JSON 读写能力前需要安装对应依赖例如pip install apache_beam[dataframe]五、写入侧对称能力WriteToJson 简要对照关联文档聚焦读取但textio.py中与ReadFromJson成对存在的是WriteToJsontextio.py#L1067-L1108它的核心参数包括path输出文件前缀实际文件名按path-XXXXX-of-NNNNN规则生成num_shards分片数量默认None由系统自动选择orient输出 JSON 的组织格式默认recordslines默认None源码中会在None时自动取orient records即 records 格式下默认按行输出。WriteToJson同样通过WriteViaPandas(json, ...)驱动底层调用DataFrame.to_json。测试 textio_test.py#L1853-L1881 的test_numeric_strings_preserved验证了读写往返过程中数字与字符串类型不会发生隐式转换这对保证数据一致性非常重要。六、完整可运行示例读 JSON → 结构化处理综合以上分析把关联文档的示例扩展为一个完整的、可直接运行的流水线import logging import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions class JsonOptions(PipelineOptions): classmethod def _add_argparse_args(cls, parser): parser.add_argument( --file_path, defaultgs://your-bucket/your-file.json, helpJson file path, supports glob patterns like gs://bucket/data*.json ) parser.add_argument( --json_lines, actionstore_true, defaultFalse, helpTreat each line as a separate JSON record (JSON Lines format) ) def run(): options JsonOptions() with beam.Pipeline(optionsoptions) as p: rows ( p | Read from Json file beam.io.ReadFromJson( pathoptions.file_path, linesoptions.json_lines, ) ) # rows 是 schema 化的元素可按字段访问 _ ( rows | Log Data Map(lambda row: logging.info( name%s score%s, row.name, row.score)) ) if __name__ __main__: logging.getLogger().setLevel(logging.INFO) run()运行方式本地 DirectRunnerpython pipeline.py --file_path./data.json --json_lines或使用 Dataflow Runner 读取 GCSpython pipeline.py \ --runnerDataflowRunner \ --projectyour-project \ --regionus-central1 \ --temp_locationgs://your-bucket/tmp \ --file_pathgs://your-bucket/your-file.json七、实践要点与注意事项lines的选择决定性能路径linesTrue时 Beam 采用按行增量解析_DelimSplitter按\n切分适合大文件与流式场景linesFalse时整文件一次性交给 pandas适合中小型完整 JSON 快照。path支持 glob需要读取多个分片文件时可直接写gs://bucket/data-*.jsonBeam 会并行处理匹配到的所有文件。dtype控制类型推断希望保持字符串原样避免数值被推断转换时使用默认False需要按列指定类型时传入{col_name: int64}这类字典。先安装 pandas 依赖缺少apache_beam[dataframe]时ReadFromJson会抛出ImportError务必在运行环境含 Dataflow worker 容器中安装该扩展。写入侧注意分片命名WriteToJson输出是path-XXXXX-of-NNNNN形式回读时用out*通配符即可一次读回全部分片见测试 textio_test.py#L1843-L1848 的往返验证模式。八、延伸阅读关联文档原文learning/prompts/code-explanation/07_io_json.mdReadFromJson/WriteToJson定义与完整 docstringsdks/python/apache_beam/io/textio.pypandas 驱动层实现read_json/to_json、ReadViaPandas/WriteViaPandassdks/python/apache_beam/dataframe/io.pyJSON 读写往返与类型保持测试sdks/python/apache_beam/io/textio_test.py同类 IO 讲解CSV 读取learning/prompts/code-explanation/08_io_csv.md赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Java Kata 实战用 TextIO.read() 从文本文件读取 PCollectionApache Beam Java Kata 实战用 TextIO.read 从文本文件读取 PCollection 本篇技术指南以 Apache Beam 官大数据批处理流处理数据工程Logto KOOK 连接器版本演进全解从 0.2.0 首个版本到 0.4.9 的源码级复盘Logto KOOK 连接器版本演进全解从 0.2.0 首个版本到 0.4.9 的源码级复盘 本文基于仓库内的 KOOK 连接器变更日志 https://li大数据批处理流处理数据工程shadPS4 安装与配置完整指南Windows / Linux / macOS 三平台跑通 PS4 模拟器shadPS4 安装与配置完整指南Windows / Linux / macOS 三平台跑通 PS4 模拟器 shadPS4 是一个用 C 写的 Play大数据批处理流处理数据工程上一篇greuler边样式与权重标注构建专业带权图可视化的完整教程下一篇cheatsheets-ai人工智能和机器学习的终极速查表指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表