首页/新闻资讯/正文详情

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

发布时间:2026/9/28 2:55:28 来源:云帆数科 栏目:资讯中心
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),仅供参考

相关推荐

ReScript reanalyze 死代码分析架构深度解析:从四阶段纯管道到响应式增量流水线
ReScript reanalyze 死代码分析架构深度解析:从四阶段纯管道到响应式增量流水线

编译器编程语言开发工具 【免费下载链接】rescript-compiler ReScript is a robustly typed language that compiles to efficient and human-readable JavaScript. 项目地址: https://gitcode.com/gh_mirrors/re/rescript-compiler 点击查看 免费下载 本文以 anal… · 2026/9/28 2:55:28

Clappr 事件系统完全指南:从 Player 映射事件到容器与播放层监听
Clappr 事件系统完全指南:从 Player 映射事件到容器与播放层监听

前端音视频插件系统 【免费下载链接】clappr An extensible, plugin-oriented, HTML5-first media player for the web 项目地址: https://gitcode.com/gh_mirrors/cl/clappr 点击查看 免费下载 Clappr 通过一套基于事件的通信机制连接 Player、Core、Container 与… · 2026/9/28 2:55:28

ClawX 内置 computer-use Skill:随包分发官方 CUA 0.25.0 的托管式安装与安全发现机制
ClawX 内置 computer-use Skill:随包分发官方 CUA 0.25.0 的托管式安装与安全发现机制

人工智能AI 应用桌面应用交互助手 【免费下载链接】ClawX ClawX is a desktop app that provides a graphical interface for OpenClaw AI agents. It turns CLI-based AI orchestration into a desktop experience without using the terminal. China website is https://claw… · 2026/9/28 2:55:28

Spingboot启动预热的实现
Spingboot启动预热的实现

启动预热的适用场景启动预热适合以下情况:数据主要来自第三方接口,无法直接从本地数据库读取。第三方接口响应较慢,首次访问容易超时。一个页面需要调用多个第三方接口或逐项查询。数据读取频繁,但变化不频繁。希望服务启动后&… · 2026/9/28 3:40:12

Understanding Driving Risks using Large Language Models: Toward Elderly Driver Assessment
Understanding Driving Risks using Large Language Models: Toward Elderly Driver Assessment

文章主要内容总结 本文研究了多模态大语言模型(具体为ChatGPT-4o)利用静态行车记录仪图像进行类人交通场景解读的潜力,重点聚焦与老年司机评估相关的三项任务:交通密度评估、交叉口可见性评估和停车标志识别。这些任务需上下文推理而非简单目标检测。研究采用零样本、少样… · 2026/9/28 3:32:43

Leveraging Large Language Models for Classifying App Users‘ Feedback
Leveraging Large Language Models for Classifying App Users‘ Feedback

文章主要内容总结 本文聚焦于利用大型语言模型(LLMs)解决应用用户反馈分类的挑战,传统方法依赖有监督机器学习,但受限于标注数据集的规模和质量。研究通过三个核心实验评估了4种先进LLMs(GPT-3.5-Turbo、GPT-4o、Flan-T5、Llama3-70b)的性能: LLMs在用户反馈分类中的基… · 2026/9/28 3:32:43

Using Large Language Models for Legal Decision-Making in Austrian Value-Added Tax Law: An Experim...
Using Large Language Models for Legal Decision-Making in Austrian Value-Added Tax Law: An Experim...

文章主要内容总结 本文通过实验评估了大型语言模型(LLMs)在奥地利及欧盟增值税(VAT)法框架下辅助法律决策的能力。研究聚焦于两种提升LLM性能的方法——微调(fine-tuning)和检索增强生成(RAG),并在两类案例中进行验证:一是权威教科书案例,二是税务咨询公司的真实案… · 2026/9/28 3:32:43

学Java别走弯路,这5个方向最吃香
学Java别走弯路,这5个方向最吃香

学Java的人很多,但学明白的人不多。有人学了半年还在写控制台程序,有人一年就能独当一面。差别不在天赋,而在方向。Java生态太庞大了,什么都学等于什么都没学。选对方向,事半功倍。今天盘点当前最吃香的5个Java方向&am… · 2026/9/28 3:32:15

AlphaAgents: Large Language Model based Multi-Agents for Equity Portfolio Constructions
AlphaAgents: Large Language Model based Multi-Agents for Equity Portfolio Constructions

AlphaAgents相关总结与翻译 一、文章主要内容总结 (一)研究背景与问题 传统股票投资组合管理依赖人类分析师处理海量信息(如财务披露、财报、市场新闻等),存在信息处理效率低、易受认知偏差(如损失厌恶、过度自信)影响的问题,可能错失投资收益机会。尽管AI在数据处理… · 2026/9/28 3:32:08

MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现
MATLAB雷达信号脉冲压缩仿真:LFM线性调频、匹配滤波与距离分辨率实现

简介:这套Matlab仿真工具完整呈现雷达信号脉冲压缩过程,从线性调频(LFM)信号生成、目标回波仿真到匹配滤波压缩处理均有可运行代码支撑,面向电子信息工程、计算机、数学等专业学生,适用于课程设计、期末大作… · 2026/9/27 0:00:01

汕头网站建设制作厂家避坑指南:5大注意事项救急
汕头网站建设制作厂家避坑指南:5大注意事项救急

汕头网站建设制作厂家避坑指南:5大注意事项救急 改个需求建站公司拖一周,这种憋屈事我见得太多了。 很多汕头老板找本地建站团队,签合同前看着方案挺美,一上线就变脸。 今天不聊虚的,直接拆解找 汕头网站建设制作厂家 时的5个核心 注意事项… · 2026/9/27 0:00:01

多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习
多模态虚假新闻检测实战:BERT+ResNet双塔与对比学习

简介:基于PyTorch的多模态虚假新闻检测项目完整代码包,面向自然语言处理与计算机视觉交叉方向的开发者、科研人员及毕业设计选题者,解决社交媒体中文本与图像联合识别虚假新闻的问题。系统以BERT预训练模型提取文本语义特征,以Res… · 2026/9/27 0:00:01

制作网页比较方便的软件怎么选?一文搞懂避坑指南
制作网页比较方便的软件怎么选?一文搞懂避坑指南

制作网页比较方便的软件怎么选?一文搞懂避坑指南 很多老板一上来就问:做个网站多少钱?但我反问他:你的域名买了吗?服务器租了吗?他一脸懵。这就是典型的“域名服务器搞不懂”。别急,今天咱们不聊虚的,直接 一文搞懂 那些让你头秃的技术名词。… · 2026/9/28 0:00:06

婚恋网站实战案例:避开3个高价坑,省钱50%还能跑赢流量
婚恋网站实战案例:避开3个高价坑,省钱50%还能跑赢流量

婚恋网站实战案例:避开3个高价坑,省钱50%还能跑赢流量 找婚恋网站建站公司,最怕的就是被坑高价。很多同行跟我吐槽,报价单上写得模棱两可,功能栏里全是“高级定制”、“专属UI”,结果落地全是套壳。今天不聊虚的,直接甩几个我经手的 实战案例… · 2026/9/28 0:00:19

济南做网站多少钱:3个案例拆解,防黑源码下载全攻略
济南做网站多少钱:3个案例拆解,防黑源码下载全攻略

济南做网站多少钱:3个案例拆解,防黑源码下载全攻略 上周济南一个做建材的老板找我,脸都绿了。他的官网首页弹出了赌博广告,后台被植入了挖矿脚本。他慌得问我:“网站被黑挂马不知道怎么办?能不能直接找之前的外包公司要源码下载,看看哪里被动了手脚?… · 2026/9/28 0:00:25

了解更多?预约专属演示

我们的顾问将为您一对一讲解产品与方案

企业微信二维码