大数据批处理流处理数据工程【免费下载链接】beamApache 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, defaultgs://your-bucket/your-file.json, helpJson file path ) def run(): This pipeline reads from Json file defined by the --file_path argument. options JsonOptions() with beam.Pipeline(optionsoptions) as p: output p | Read from Json file ReadFromJson( pathoptions.file_path, linesFalse ) if __name__ __main__: logging.getLogger().setLevel(logging.INFO) run()运行方式支持本地 Direct Runner 与远程 runnerpython 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(optionsoptions) as p以with语句创建管道。退出with块时管道会自动运行这是 Beam 2.x 中推荐的上下文管理用法。ReadFromJson(path..., linesFalse)内置的 JSON 读取变换负责把 JSON 文件内容转成 PCollection。这里特意指定linesFalse表示整个文件作为一个整体解析详见下文参数详解。二、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, orientorient, lineslines, dtypedtype, **kwargs)参数语义详解参数默认值说明path必填要读取的文件路径支持 glob 通配符如*、?因此可以一次读取一个目录下的多个文件如data/out*orientrecordsJSON 元素的组织格式。records表示文件内容是一个 JSON 对象列表形如[{field1: value1, field2: value2, ...}, ...]还支持 pandasread_json的其它取值如columns、index、split、table、values等linesTrue是否将每一行视为一条独立记录。True对应 JSON LinesNDJSON格式False表示整个文件是一个合法的 JSON 对象或数组dtypeFalse数据类型处理策略True时自动推断类型传入{列名: 类型}字典时按指定类型读取False时完全不推断类型。注意默认值与 pandas 不同pandas 默认会推断**kwargs—透传给pandas.read_json的其它参数如nrows之外的 pandas 参数。注意nrows在 Beam 的封装中暂不支持提示原文档示例中显式传了linesFalse而源码默认值是linesTrue。二者都是合法用法取决于你的源文件格式——如果文件是标准的 JSON 数组每行不是独立的 JSON应使用linesFalse如果是 JSON Lines 流式文件保留默认linesTrue即可。底层调用链从源码可以看到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当linesTrue时使用_DelimSplitter(b\n, _DEFAULT_BYTES_CHUNKSIZE)以换行符作为记录边界做增量读取incremental 模式这也正是 JSON Lines 文件能够被高效、分布式读取的底层机制当linesFalse时按整文件读取。ReadViaPandas.expand把读取结果转成 DataFrame 后对object类型的列统一转换为pd.StringDtype()objects_as_stringsTrue再通过convert.to_pcollection将 DataFrame 转换为 Beam PCollection。这意味着ReadFromJson输出的每个元素实际上是一行记录的命名元组row tuple字段名与 JSON 对象的键一一对应你可以直接通过属性访问。三、数据格式适配orient 与 lines 的搭配场景一标准 JSON 数组文件文件内容形如[ {name: alice, age: 30}, {name: bob, age: 25} ]应使用orientrecords默认配合linesFalsepcoll p | ReadFromJson(pathgs://bucket/array.json, linesFalse)场景二JSON LinesNDJSON文件文件每行一条记录适合流式追加和分布式分片读取{name: alice, age: 30} {name: bob, age: 25}使用默认的linesTrue即可。此时 Beam 的_DelimSplitter会按换行符切分记录天然支持多 worker 并行读取不同分片。场景三按 glob 读取多个文件path参数支持 glob 模式一次读取一批文件pcoll p | ReadFromJson(pathdata/2024/*.json, linesTrue)该能力在 textio_test.py 的读写往返测试中得到了验证测试先用WriteToJson写出分片文件形如out-00000-of-00003再通过ReadFromJson(os.path.join(dest, out*))读取并断言与原始beam.Row完全一致。四、类型保持数值与字符串不被误转换JSON 读取最容易踩的坑是数值与数字字符串被混淆。仓库测试 test_numeric_strings_preserved 专门验证了这一点写出as_string0、as_float_string0.0、as_int0、as_float0.0四种类型的数据后再读回并逐字段断言类型完全一致。这得益于ReadViaPandas.expand中对object类型列统一转换为pd.StringDtype()的处理——字符串列不会被 pandas 自动解析成数值。如果你需要更精确的类型控制可以直接使用dtype参数指定列类型pcoll p | ReadFromJson( pathdata.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, helpJson file path) parser.add_argument(--output, helpOutput path prefix) def run(): options JsonOptions() with beam.Pipeline(optionsoptions) as p: records ( p | Read from Json ReadFromJson( pathoptions.file_path, linesFalse) ) 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-NNNNNorientrecords默认输出为 JSON 对象数组linesNone时自动取orient records的结果——即 records 格式下默认输出为 JSON Linesnum_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 数据的标准入口它以pathorientlinesdtype四个核心参数覆盖了绝大多数 JSON 文件形态底层通过 pandasread_json与ReadViaPandas封装获得类型推断、换行增量读取和 DataFrame 到 PCollection 的自动转换。配合WriteToJson可以实现完整的 JSON 读写闭环而textio_test.py中的往返测试则保证了数值/字符串类型在读写过程中不被破坏。将本文示例中的JsonOptions与ReadFromJson组合使用即可快速搭建生产级的 JSON 数据接入管道。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐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),仅供参考