ARTICLE DETAIL

资讯详情

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

Apache Beam 中 ReadFromCsv 读取 CSV 文件:PipelineOptions 自定义参数与 ReadViaPandas 底层实现解析

Apache Beam 中 ReadFromCsv 读取 CSV 文件:PipelineOptions 自定义参数与 ReadViaPandas 底层实现解析 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的 Python SDK 内置了 CSV 文件读写能力ReadFromCsv变换Transform可以把一个或多个逗号分隔值CSV文件读取为一个PCollection。本篇文章以仓库中的代码讲解文档 learning/prompts/code-explanation/08_io_csv.md 为骨架深入剖析CsvOptions自定义 PipelineOptions 如何注入--file_path命令行参数、ReadFromCsv如何基于 pandas 实现并给出可复制、可运行的完整示例以及参数取值范围说明。读完本文你将掌握如何在 Beam Python 管道中自定义命令行参数类、如何用ReadFromCsv读取 CSV 数据并配合Map(logging.info)查看内容以及该变换底层真实的调用链与实现细节。一、原代码逐段拆解CsvOptions 与 ReadFromCsv 的组合1.1 完整代码回顾讲解文档给出的示例代码是一个典型的 Beam Python 管道先用自定义的 PipelineOptions 子类接收--file_path参数再通过ReadFromCsv读取该 CSV 文件最后用Map(logging.info)将每一行数据打印到日志class CsvOptions(PipelineOptions): classmethod def _add_argparse_args(cls, parser): parser.add_argument( --file_path, defaultgs://your-bucket/your-file.csv, helpCsv file path ) options CsvOptions() with beam.Pipeline(optionsoptions) as p: output (p | Read from Csv file ReadFromCsv(pathoptions.file_path) | Log Data Map(logging.info))该代码由三个相互协作的部分组成自定义选项类参数注入、管道构建读取与处理以及日志输出结果观察。下面逐一拆解。1.2 自定义 PipelineOptions 子类命令行参数的声明与解析class CsvOptions(PipelineOptions): classmethod def _add_argparse_args(cls, parser): parser.add_argument( --file_path, defaultgs://your-bucket/your-file.csv, helpCsv file path )PipelineOptions是 Apache Beam Python SDK 提供的选项基类其核心机制是子类通过类方法_add_argparse_args(cls, parser)向参数解析器argparse parser注册自定义命令行参数。在仓库源码 sdks/python/apache_beam/options/pipeline_options.py 中可以看到PipelineOptions.__init__构建解析器后会遍历类型的 MRO方法解析顺序对每个声明了_add_argparse_args的类依次调用该方法完成参数注册parser _BeamArgumentParser(allow_abbrevFalse) for cls in type(self).mro(): if cls PipelineOptions: break elif _add_argparse_args in cls.__dict__: cls._add_argparse_args(parser) # type: ignore也就是说--file_path一经注册就可以通过命令行传参例如python pipeline.py --file_path gs://my-bucket/data.csv也可以像示例中那样直接以属性方式读取options.file_path。_add_argparse_args是 Beam 约定俗成的自定义选项扩展点仓库内几乎所有内置选项类如 Runner 相关选项都采用同样的写法。对parser.add_argument的各个参数做一下补充说明--file_path命令行参数名注意 Python SDK 内部通过argparse解析时会自动把短横线转换为下划线因此访问属性时使用options.file_pathdefaultgs://your-bucket/your-file.csv默认值示例中是一个占位 GCS 路径实际使用时建议改为真实的本地路径或对象存储路径helpCsv file path参数说明文字在--help输出中展示也方便管道使用者理解该参数的用途。从源码结构看CsvOptions()直接实例化时参数解析使用parse_known_args见 pipeline_options.py 第 384 行因此未被识别的参数不会导致崩溃这也使得自定义选项可以与其他标准选项如--runner、--project共存于同一命令行。二、ReadFromCsv 变换基于 pandas 的 CSV 读取2.1 变换声明与核心参数ReadFromCsv位于 Beam Python SDK 内置的 TextIO 连接器中sdks/python/apache_beam/io/textio.py它在模块顶部被导出__all__中包含ReadFromCsv。其完整的函数签名与参数说明如下来自 textio.py 第 975-995 行def ReadFromCsv( path: str, *, splittable: bool True, filename_column: Optional[str] None, **kwargs):path (str)要读取的文件路径支持通配符glob如*和?可以一次读取多个分片文件splittable (bool)文件是否可以在行边界上动态切分即文件的每一行是否代表一条完整记录。如果单条记录跨越多行例如某个带引号的字段内部含有换行符应将其设为False否则可能产生残缺记录设为False也可能禁用 liquid sharding液态分片即按需动态拆分的并行机制filename_column (str)如果非None会在每条记录上新增一列内容为该记录来源文件的文件名**kwargs透传给pandas.read_csv的额外关键字参数详见下文。2.2 底层实现ReadViaPandas 与 read_csvReadFromCsv的函数体非常简单本质上是ReadViaPandas的语法糖封装from apache_beam.dataframe.io import ReadViaPandas return ReadFromCsv ReadViaPandas( csv, path, splittablesplittable, filename_columnfilename_column, **kwargs)ReadViaPandas定义于 sdks/python/apache_beam/dataframe/io.py 第 806 行其__init__中根据format csv把filename_column注入 kwargs然后动态调用read_csv(...)构造 readerif format csv: kwargs[filename_column] filename_column self._reader globals()read_%s % format其中read_csv同样位于 io.py 第 89 行它把参数包装后交给 pandas 的pd.read_csv并支持两个关键增强incrementalTrue增量读取允许大文件被分批解析而不是一次性载入内存splitter_TextFileSplitter(args, kwargs) if splittable else None当splittableTrue时按换行边界切分文件从而支持动态拆分dynamic splitting与并行处理。read_csv的 docstring 还明确警告如果文件较大且记录中不包含带引号的换行符可以传splittableTrue启用按换行符的动态切分若记录包含引号内的换行符却使用该选项可能导致记录残缺和数据损坏。expand阶段io.py 第 822-829 行把 DataFrame 转换成PCollection先通过convert.to_pcollection(df, include_indexesFalse)输出并且对于 dtype 为object的列会显式转换为pd.StringDtype()保证列类型一致。2.3 pandas 参数透传真正可用的 kwargsReadFromCsv的**kwargs会一路透传给pandas.read_csvReadFromCsv的 docstring 通过append_pandas_args装饰器自动追加了 pandas 的完整参数文档见 textio.py 第 941-971 行的装饰器实现。这意味着你可以在调用时使用 pandas 的常见读取参数例如delimiter/sep指定分隔符默认逗号例如读取制表符分隔的 TSV 文件时传delimiter\tencoding指定文件编码如encodinglatin1仓库测试 textio_test.py 的test_non_utf8_csv_read_write正是用encodinglatin1读取非 UTF-8 的 CSVheader指定表头所在行headerNone表示文件无表头dtype指定各列的数据类型避免类型推断偏差names自定义列名列表skiprows跳过起始若干行。注意filepath_or_buffer与iterator两个 pandas 参数被显式排除见装饰器调用处exclude[filepath_or_buffer, iterator]因为它们由 Beam 的文件系统抽象与增量读取机制接管。三、管道主体读取、处理与结果观察with beam.Pipeline(optionsoptions) as p: output (p | Read from Csv file ReadFromCsv(pathoptions.file_path) | Log Data Map(logging.info))3.1 with 语句管理管道生命周期with beam.Pipeline(optionsoptions) as p:是 Beam Python 的推荐写法进入上下文后创建管道退出时自动执行p.run()并等待结果完成即run()wait_until_finish()的等效行为。optionsoptions把第一节定义的CsvOptions实例作为管道配置传入从而让--file_path的值在管道内部可见。3.2 数据流从 PCollection 到日志p | Read from Csv file ReadFromCsv(pathoptions.file_path)读取 CSV输出一个元素类型为命名元组namedtuple的PCollection每个元素对应 CSV 中的一行记录字段名与 CSV 表头列一一对应| Log Data Map(logging.info)对每个元素调用logging.info把每行数据输出到日志通常是标准输出/作业日志方便开发期观察读取结果。标签如Read from Csv file在分布式执行时用于在作业图中定位具体步骤同时也可作为步骤重命名的依据。四、实战完整可运行的示例与参数说明4.1 可直接运行的完整管道结合以上分析下面给出一个更完整、可直接运行的示例将默认路径替换为你的真实文件即可import logging import apache_beam as beam from apache_beam.io import ReadFromCsv from apache_beam.options.pipeline_options import PipelineOptions class CsvOptions(PipelineOptions): classmethod def _add_argparse_args(cls, parser): parser.add_argument( --file_path, defaultgs://your-bucket/your-file.csv, helpCsv file path) parser.add_argument( --with_filename, actionstore_true, helpWhether to include the source filename column) options CsvOptions() csv_kwargs {} if options.with_filename: csv_kwargs[filename_column] source_filename with beam.Pipeline(optionsoptions) as p: rows (p | Read from Csv file ReadFromCsv( pathoptions.file_path, **csv_kwargs) | Log Data Map(logging.info))运行方式python pipeline.py \ --file_path ./data/input.csv \ --runner DirectRunner4.2 读写配套WriteToCsv 快速生成测试数据仓库中与ReadFromCsv成对的是WriteToCsv同样位于 sdks/python/apache_beam/io/textio.py 第 1006-1035 行它接收 schema 化的PCollection写出以指定前缀命名的分片文件默认命名规则为path-XXXXX-of-NNNNN并自动关闭 DataFrame 索引indexFalse。测试 textio_test.py 的test_csv_read_write展示了完整的写-读-断言回路records [beam.Row(astr, bix) for ix in range(3)] with TestPipeline() as p: p | beam.Create(records) | beam.io.WriteToCsv(os.path.join(dest, out)) with TestPipeline() as p: pcoll ( p | beam.io.ReadFromCsv(os.path.join(dest, out*)) | beam.Map(lambda t: beam.Row(**dict(zip(type(t)._fields, t))))) assert_that(pcoll, equal_to(records))这段测试还揭示了一个重要细节写入会生成带分片后缀的多个文件如out-00000-of-00001.csv因此读取时必须使用通配符out*才能把全部分片读回来——这正是ReadFromCsv支持 glob 路径的典型应用场景。4.3 filename_column 的实测行为test_csv_read_with_filenametextio_test.py 第 1769-1792 行验证了filename_columnsource_filename的行为读取后每条记录会额外携带source_filename字段其值是该记录来源的实际分片文件名。这在多文件合并处理的场景例如按来源追溯数据中非常实用。五、小结自定义选项 内置 CSV 变换的完整调用链回顾整个调用链示例代码的技术脉络可以归纳为参数声明CsvOptions(PipelineOptions)通过_add_argparse_args注册--file_pathPipelineOptions.__init__依据 MRO 遍历所有子类并注册参数pipeline_options.py管道配置beam.Pipeline(optionsoptions)接收选项实例options.file_path直接读取解析后的值数据读取ReadFromCsv(path...)→ReadViaPandas(csv, ...)→read_csv(...)→pd.read_csv增量、可切分输出 DataFrame 后再转换为元素为命名元组的PCollectiontextio.py、dataframe/io.py结果观察Map(logging.info)将每行记录输出到日志。由此可以看出ReadFromCsv并不是一个手写的 CSV 解析器而是将 Beam 的分布式文件读取、动态切分能力与 pandas 成熟的 CSV 解析能力深度整合Beam 负责在哪些文件、按什么边界并行读取pandas 负责如何把一行文本解析成结构化字段。理解这层设计之后你在实际项目中就可以放心地把sep、encoding、dtype、skiprows等 pandas 参数透传给ReadFromCsv同时通过splittable与filename_column获得 Beam 特有的分布式读取能力。如需进一步深入可以继续阅读变换定义与全部参数sdks/python/apache_beam/io/textio.py底层 DataFrame 读取实现sdks/python/apache_beam/dataframe/io.pyCSV 读写与编码、filename_column 的测试验证sdks/python/apache_beam/io/textio_test.pyPipelineOptions 参数注册机制sdks/python/apache_beam/options/pipeline_options.py本文所讲解代码的原始出处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 DoFn 附加参数完全指南Timestamp、Window、PaneInfo 与 PipelineOptions 注入机制详解Apache Beam Java DoFn 附加参数完全指南Timestamp、Window、PaneInfo 与 PipelineOptions 注入机制详大数据批处理流处理数据工程Apache Beam Java Kata 实战用 TextIO.read() 从文本文件读取 PCollectionApache Beam Java Kata 实战用 TextIO.read 从文本文件读取 PCollection 本篇技术指南以 Apache Beam 官大数据批处理流处理数据工程OpenMMO服务器状态持久化SIGTERM优雅关闭全流程解析OpenMMO服务器状态持久化SIGTERM优雅关闭全流程解析 OpenMMO 是一款用 Rust Svelte 构建的开放世界 MMORPG。当运维人员游戏开发AI Agent人工智能上一篇CANN/PTO-ISAMegaMoE调度融合算子示例下一篇CANNOpsTransformer Flash Attention示例创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表