尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
Apache Beam Python DataFrame API 示例管道实战:从 Wordcount 到航班延误分析
批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载导读Apache Beam Python SDK 在 2.30.0 版本引入了apache_beam.examples.dataframe示例模块集中展示了如何用 pandas 风格的 DataFrame API 编写批量与流式数据处理管道。本文以仓库中 examples/dataframe 目录 下的三个官方示例为主线深入讲解 DataFrame API 的两种典型用法——嵌入传统 Beam 管道通过to_dataframe/to_pcollection与 Beam Schema 协同与纯 DataFrame 端到端管道直接使用 DataFrame IO并结合源码说明其底层执行原理。读完本文你将能够独立运行这些示例、理解 GroupBy/merge/GroupBy.apply 等复杂聚合在分布式环境下的行为并学会如何在自有管道中复用这套 API。1. 前置条件与环境准备运行这些示例需要满足以下条件来自 README.mdApache Beam Python SDK 版本必须安装apache-beam2.30.0因为apache_beam.examples.dataframe模块在该版本才被加入其中航班延误示例更进一步标注为2.31.0 新增。pandas 版本DataFrame API 底层构建在 pandas 实现之上需要安装兼容的 pandas 版本详见官方 DataFrames 文档的前置要求部分。从 Beam 2.34.0 起最简单的安装方式是通过dataframe扩展pip install apache_beam[dataframe]注意在分布式 Runner 上执行 DataFrame 管道时Worker 节点上也应安装相同版本的 pandas否则可能因版本差异导致行为不一致。Beam DataFrames 的核心设计是pandas 的 DataFrame 方法被并行地调用在数据子集上而所有操作由 Beam API **延迟deferred**执行以适配 Beam 的并行处理模型。你可以把它看作 Beam 管道的一门领域专用语言DSL——与 Beam SQL 类似但接口是开发者熟悉的 pandas DataFrame API。2. Wordcount 管道将 DataFrame API 嵌入传统 Beam 管道Wordcount 是数据分析系统的Hello World仓库用 DataFrame API 实现了它代码位于 wordcount.py。2.1 核心实现解读该示例的精华在于演示了 DataFrame API 如何与更大的 Beam 管道集成先用传统 Beam 变换读取与切分文本再把 PCollection 转成 DataFrame 做聚合最后再转回 PCollection 继续用 Beam 变换处理。完整逻辑如下import apache_beam as beam from apache_beam.dataframe.convert import to_dataframe from apache_beam.dataframe.convert import to_pcollection from apache_beam.io import ReadFromText with beam.Pipeline(optionsPipelineOptions(pipeline_args)) as p: # 1) 用传统 Beam 变换读取文本并切词 lines p | Read ReadFromText(known_args.input) words ( lines | Split beam.FlatMap( lambda line: re.findall(r[\w], line)).with_output_types(str) # 2) 映射为 Row 对象生成适合转 DataFrame 的 schema | ToRows beam.Map(lambda word: beam.Row(wordword))) # 3) PCollection - DataFramepandas 风格聚合 df to_dataframe(words) df[count] 1 counted df.groupby(word).sum() counted.to_csv(known_args.output) # 4) DataFrame - PCollection继续用 Beam 变换处理 counted_pc to_pcollection(counted, include_indexesTrue) _ ( counted_pc | beam.Filter(lambda row: row.count 50) | beam.Map(lambda row: f{row.word}: {row.count}) | beam.Map(print))关键点逐一拆解beam.Row(wordword)创建 schema-aware PCollectionto_dataframe要求输入 PCollection 带 schema即 schema-aware PCollectionbeam.Row会根据字段自动推断并绑定 schema。代码中with_output_types(str)的作用是让FlatMap的输出类型被明确标注从而保证后续 schema 推断正确。to_dataframe与to_pcollection二者均位于 apache_beam/dataframe/convert.py实现 PCollection 与延迟 DataFrame之间的双向转换。to_pcollection(counted, include_indexesTrue)将聚合结果此时word是 groupby 后的索引连同索引一并写回 PCollection使每行成为带word、count字段的 Row。df[count] 1新增列这是 pandas 风格操作等价于为每个单词分配计数 1随后groupby(word).sum()完成词频统计。混合编程的价值聚合这种向量化操作交给 pandas 高效实现而词频过滤count 50与打印等元素级逻辑仍用 Beam 变换表达二者各取所长。2.2 运行与预期输出在本地运行默认使用 Direct Runnerpython -m apache_beam.examples.dataframe.wordcount \ --input gs://dataflow-samples/shakespeare/kinglear.txt \ --output counts输出会写入counts-XXXXX-of-YYYYY形式的分片文件XXXXX为分片序号YYYYY为总分片数内容形如KING: 243 LEAR: 236 DRAMATIS: 1 PERSONAE: 1 king: 65 of: 447 Britain: 2 OF: 15 FRANCE: 10 DUKE: 3 ...命令行参数说明定义在wordcount.py的 argparse 部分--input默认指向gs://dataflow-samples/shakespeare/kinglear.txt可省略--output为必填指定输出文件前缀其余参数如--runner、--project、--temp_location等可通过命令行透传给PipelineOptions。2.3 测试如何验证正确性wordcount_test.py 使用TestPipeline构造临时输入文本运行wordcount.run()后读取*.result*分片用正则(\S),([0-9])解析出word,count行与本地re.findall统计的期望词频逐项比对。它同时验证了两件事CSV 输出的表头格式为word,count以及聚合结果与顺序无关比对前先排序。3. 出租车示例管道纯 DataFrame 的端到端管道与 wordcount 不同taxiride.py 中的两条管道不使用任何 Beam 原语没有 ParDo、GroupByKey 等而是完全依靠 DataFrame API 及其自带的 DataFrame IO如read_csv、to_csv构建端到端管道。处理的数据是知名的 NYC Taxi 开放数据集。3.1 示例数据集说明仓库在公共 GCS 桶gs://apache-beam-samples中预置了如下数据快照详见 README.md路径说明gs://apache-beam-samples/nyc_taxi/2017/yellow_tripdata_2017-*.csv2017 年逐月出租车行程 CSV2018、2019 年有类似目录gs://apache-beam-samples/nyc_taxi/misc/sample.csv2019 年初的 100 万条记录抽样约 85 MiB适合本地处理gs://apache-beam-samples/nyc_taxi/misc/taxi_zone_lookup.csvZone ID 查询表供borough_enrich管道使用仓库内 data 目录 还存放了两份 2018 年全量数据的期望输出taxiride_2018_aggregation_truth.csv与taxiride_2018_enrich_truth.csv供集成测试做逐帧比对pd.testing.assert_frame_equal。3.2location_id_agg按落客点分组聚合管道核心逻辑极其简洁from apache_beam.dataframe.io import read_csv with pipeline as p: rides p | read_csv(input_path) # 统计每个 DOLocationID 的下客乘客总数 agg rides.groupby(DOLocationID).passenger_count.sum() agg.to_csv(output_path)read_csv(input_path)直接作为 PTransform 使用pandas 会从 CSV 首行推断列名因此passenger_count与DOLocationID直接可用。read_*支持文件模式glob以及所有 Beam 兼容文件系统。groupby(...).passenger_count.sum()一次分组聚合即完成按落客点 ID 汇总乘客数底层对应 Beam 的 group-by-key 与向量化 sum。全程唯一的传统 Beam 类型只有Pipeline实例本身。运行命令python -m apache_beam.examples.dataframe.taxiride \ --pipeline location_id_agg \ --input gs://apache-beam-samples/nyc_taxi/misc/sample.csv \ --output aggregation.csv输出写入aggregation.csv-XXXXX-of-YYYYY分片内容形如DOLocationID,passenger_count 1,3852 3,130 4,7725 5,24 6,37 7,7429 8,24 9,180 10,938 ...提示taxiride.py的 argparse 中--pipeline默认值为location_id_agg合法取值为location_id_agg与borough_enrich二选一--input默认即为上面的 sample.csv 路径。3.3borough_enrich与 Zone 查询表 join 后聚合这条管道在聚合基础上增加了DataFrame join把 Zone 查询表按LocationID设为索引再与行程数据的DOLocationID做左连接从而获得每个落客点所属的 Borough行政区最后按 Borough 汇总乘客数ZONE_LOOKUP_PATH gs://apache-beam-samples/nyc_taxi/misc/taxi_zone_lookup.csv def run_enrich_pipeline(pipeline, input_path, output_path, zone_lookup_pathZONE_LOOKUP_PATH): with pipeline as p: rides p | Read taxi rides read_csv(input_path) zones p | Read zone lookup read_csv(zone_lookup_path) # 先让 zones 以 LocationID 为索引再左连接 rides rides.merge( zones.set_index(LocationID).Borough, right_indexTrue, left_onDOLocationID, howleft) # 按 Borough 汇总乘客数 agg rides.groupby(Borough).passenger_count.sum() agg.to_csv(output_path)代码注释中特别说明了一个重要的实现取舍另一种更直观的写法rides.merge(zones[[LocationID,Borough]], howleft, left_onDOLocationID, right_onLocationID)虽然在 pandas 中更常见但不保留索引因而在分布式执行时无法并行所以示例选择了先set_index再按索引合并的方式。这是 DataFrame API 与本地 pandas 在并行语义上的典型差异。运行命令python -m apache_beam.examples.dataframe.taxiride \ --pipeline borough_enrich \ --input gs://apache-beam-samples/nyc_taxi/misc/sample.csv \ --output enrich.csv输出写入enrich.csv-XXXXX-of-YYYYY内容形如Borough,passenger_count Bronx,13645 Brooklyn,70654 EWR,3852 Manhattan,1417124 Queens,81138 Staten Island,531 Unknown,285273.4 测试与已知注意事项taxiride_test.py 把 10 行样本数据复制到100 个 CSV 文件模拟多文件读取再用本地 pandas 计算期望结果与管道输出逐行比对验证了read_csv对文件模式与多分片的正确处理。taxiride_it_test.py 则对 2018 年全年数据gs://apache-beam-samples/nyc_taxi/2018/*.csv做端到端验证与 truth 文件比对其中test_enrich特意将 worker 机型提升为e2-highmem-2注释说明标准机型在该管道上会 OOM内存不足——这提示 join 类 DataFrame 操作对 worker 内存有更高要求。测试代码中的 TODO 注释如 issue #20926指出示例当前输出的是浮点型总和而非整型这是已知的精度取舍不影响结果正确性。4. 航班延误管道窗口 GroupBy.apply 的复杂聚合flight_delays.py 是 2.31.0 加入的第三个示例它把 DataFrame API 用在带时间窗口的流式语义场景中先用传统 Beam 管道从 BigQuery 读取航班准点数据套用 24 小时滚动窗口并定义 Beam Schema再转成 DataFrame 执行复杂的GroupBy.apply聚合最后用to_csv写出。DataFrame 计算尊重之前施加的 24 小时窗口结果按天分文件输出。4.1 数据读取与窗口设置query f SELECT FlightDate AS date, IATA_CODE_Reporting_Airline AS airline, Origin AS departure_airport, Dest AS arrival_airport, DepDelay AS departure_delay, ArrDelay AS arrival_delay FROM apache-beam-testing.airline_ontime_data.flights WHERE FlightDate {start_date} AND FlightDate {end_date} AND DepDelay IS NOT NULL AND ArrDelay IS NOT NULL tbl ( p | read table beam.io.ReadFromBigQuery( queryquery, use_standard_sqlTrue) | assign timestamp beam.Map(lambda x: window.TimestampedValue(x, to_unixtime(x[date]))) # 用 beam.Select 确保数据带 schemalambda 中的类型转换保证类型推断正确 | set schema beam.Select( datelambda x: str(x[date]), airlinelambda x: str(x[airline]), departure_airportlambda x: str(x[departure_airport]), arrival_airportlambda x: str(x[arrival_airport]), departure_delaylambda x: float(x[departure_delay]), arrival_delaylambda x: float(x[arrival_delay]))) daily tbl | daily windows beam.WindowInto( beam.window.FixedWindows(60 * 60 * 24))注意这里的两个关键工程细节beam.Select显式重建 schema从 BigQuery 读出的行先被打上时间戳再用beam.Select配合 lambda 把每列强制转换为明确类型str/float保证to_dataframe能得到正确的 schema 与 dtype 推断。FixedWindows(60*60*24)即 24 小时窗口整个 DataFrame 计算都发生在该窗口上下文内因此后续聚合天然按天隔离。4.2 复杂聚合GroupBy.apply窗口化后的数据转为 DataFrame按航空公司分组并对每组应用自定义函数get_mean_delay_at_top_airports——统计该航司在最繁忙的 10 个机场上的平均延误def get_mean_delay_at_top_airports(airline_df): arr airline_df.rename(columns{ arrival_airport: airport }).airport.value_counts() dep airline_df.rename(columns{ departure_airport: airport }).airport.value_counts() total arr dep # 保留全部含重复以保证结果确定性 # pandas 1.4.0 可能输出 NaN这里显式丢弃 top_airports total.nlargest(10, keepall).dropna() at_top_airports airline_df[arrival_airport].isin( top_airports.index.values) return airline_df[at_top_airports].mean(numeric_onlyTrue) df to_dataframe(daily) result df.groupby(airline).apply(get_mean_delay_at_top_airports) result.to_csv(output)groupby(airline).apply(fn)对每个航司分组调用自定义函数是 DataFrame API 支持的最灵活的聚合方式适合无法用sum/mean等内置算子表达的复合逻辑。函数内value_counts统计起降频次并相加nlargest(10, keepall)取前 10 繁忙机场保留并列项以保证结果确定dropna()显式处理 pandas 1.4.0 可能引入的 NaN。date参数由自定义类型转换函数input_date校验只接受2002-01-01到2012-12-31范围内的日期数据集覆盖范围格式为%Y-%m-%d。4.3 运行与预期输出由于要从 BigQuery 读数据必须提供 GCPproject与temp_locationpython -m apache_beam.examples.dataframe.flight_delays \ --start_date 2012-12-24 \ --end_date 2012-12-25 \ --output gs://bucket/dir/delays.csv \ --project gcp-project \ --temp_location gs://bucket/dir输出文件按窗口分区命名形如gs://bucket/dir/delays.csv-2012-12-24T00:00:00-2012-12-25T00:00:00-XXXXX-of-YYYYY即输出前缀-窗口起始-窗口结束-分片序号-总分片数内容为每天各航司的平均延误airline,departure_delay,arrival_delay EV,10.01901901901902,4.431431431431432 HA,-1.0829015544041452,0.010362694300518135 UA,19.142555438225976,11.07180570221753 VX,62.755102040816325,62.61224489795919 WN,12.074298711144806,6.717968157695224 ...4.4 测试验证flight_delays_it_test.py 以2012-12-23至2012-12-25三天为输入运行管道然后按日期逐个匹配输出分片f{self.output_path}-{date}*与硬编码在测试类中的EXPECTED字典含 15 家航司每天的departure_delay、arrival_delay期望值用pd.testing.assert_frame_equal精确比对。这从侧面印证了按天分区输出、窗口内 DataFrame 聚合这两个核心行为的确定性。5. 深入DataFrame API 的底层机制与更多用法5.1 模块结构一览DataFrame API 的核心实现位于 apache_beam/dataframe 目录与示例直接相关的模块包括convert.pyto_dataframe/to_pcollection双向转换io.pyread_csv、to_csv等 DataFrame IO支持文件模式与 Beam 文件系统transforms.pyDataframeTransformPTransformframes.pyDeferredDataFrame等延迟帧类型expressions.py表达式图与求值模型。5.2 不写管道样板DataframeTransform如果不想手动编排to_dataframe/to_pcollection可以改用DataframeTransform见 transforms.py。它类似于 Beam SQL 中的SqlTransform接受一个接收并返回 DataFrame的函数作为 PTransform 直接应用在 schema-aware PCollection 上内部自动完成批量化、转换与求值from apache_beam.dataframe.transforms import DataframeTransform with beam.Pipeline() as p: ... | beam.Select(DOLocationIDlambda line: int(..), passenger_countlambda line: int(..)) | DataframeTransform(lambda df: df.groupby(DOLocationID).sum()) | beam.Map(lambda row: f{row.DOLocationID},{row.passenger_count}) ...从源码的expand方法可以看到其完整调用链输入 PCollection 通过convert.to_dataframe转为延迟帧 → 以位置参数tuple或关键字参数dict方式调用用户函数 → 结果经convert.to_pcollectionyield_elements默认schemas即逐行展开为 schema-aware PCollection转回。它同样支持多输入多输出output (pc1, pc2) | DataframeTransform(lambda df1, df2: ...) output {a: pc, ...} | DataframeTransform(lambda a, ...: ...) pc1, pc2 {a: pc} | DataframeTransform(lambda a: expr1, expr2) {...} {a: pc} | DataframeTransform(lambda a: {...})当输入 PCollection 的元素本身就是 pandas DataFrame/Series 时还需通过proxy参数提供一个同 dtype 的空实例用于类型推断。DataframeTransform特别适合把既能跑在本地 pandas、又能跑在 Beam 上的独立函数直接复用。5.3 与 pandas 的差异须知Beam DataFrames 目标是兼容 pandas 原生实现但二者存在语义差异详见仓库文档 differences-from-pandas.md操作延迟执行Beam 管道构建阶段只记录操作真正计算发生在管道运行期因此无法像本地 pandas 一样即时查看中间结果并行分区语义groupby、join、排序等操作受分区partitioning约束部分写法如前面提到的不保留索引的 merge无法并行化确定性要求为支持分布式重试与验证应避免依赖全局状态或非确定性操作示例代码中多处注释都围绕这一点做处理。6. 小结与进一步实践三个示例构成了 DataFrame API 的完整学习路径示例文件核心技巧新增版本Wordcountwordcount.pyBeam Schema to_dataframe/to_pcollection混合编程2.30.0出租车聚合taxiride.py纯 DataFrame 端到端管道、read_csv/to_csv2.30.0出租车 enrich同上DataFrame merge 索引策略2.30.0航班延误flight_delays.py24h 窗口 GroupBy.apply复杂聚合2.31.0继续深入可从三处入手其一阅读 convert.py 与 transforms.py 的测试文件理解转换边界与DataframeTransform的参数语义其二参照 taxiride_test.py 的多文件模式用法把read_csv的 glob 能力用到自己的分片数据上其三参考 flight_delays_it_test.py 的分窗口断言方式为自定义窗口聚合编写可验证的测试。对于交互式探索仓库还提供了 dataframes.ipynb 笔记本可以在 Colab 中边运行边理解 DataFrame API 的延迟执行语义。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐流媒体下载工具 N_m3u8DL-RE 使用指南从 M3U8 到 MP4 的完整配置与避坑手册流媒体下载工具 N_m3u8DL RE 使用指南从 M3U8 到 MP4 的完整配置与避坑手册 N_m3u8DL RE 是一款跨平台的流媒体下载器面向 MPApache Beam DataFrame API 实战Wordcount、纽约出租车与航班延误示例流水线深度解析Apache Beam DataFrame API 实战Wordcount、纽约出租车与航班延误示例流水线深度解析 导读 Apache Beam 的 Pyth大数据批处理流处理数据工程Apache Beam YAML 示例目录实战指南从 Wordcount 到 Kafka、Iceberg 与 ML 管道Apache Beam YAML 示例目录实战指南从 Wordcount 到 Kafka、Iceberg 与 ML 管道 本文基于 Apache Beam 仓大数据批处理流处理数据工程上一篇为什么源师兄Python IDE的编辑体验如此顺滑CodeMirror 6与Python结构高亮深度解析下一篇罗技鼠标宏真的会被封号吗logitech-pubg混淆设计与Ban风险分析创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED

相关推荐

CodeQL C 查询包 0.7.0 版本解析:新增函数级访问控制检测与 Zip Slip 查询更名

CodeQL C 查询包 0.7.0 版本解析:新增函数级访问控制检测与 Zip Slip 查询更名

静态分析SAST应用安全漏洞扫描代码质量 【免费下载链接】codeql CodeQL: the libraries and queries that power security researchers around the world, as well as code scanning in GitHub Advanced Security 项目地址: https://gitcode.com/gh_mirrors/co/code…

📅 2026/10/10 2:44:18
js-lingui Metro Transformer 深度指南:在 React Native 与 Expo 中直接编译 .po 目录文件

js-lingui Metro Transformer 深度指南:在 React Native 与 Expo 中直接编译 .po 目录文件

开发工具前端 【免费下载链接】js-lingui 🌍 📖 A readable, automated, and optimized (2 kb) internationalization for JavaScript 项目地址: https://gitcode.com/gh_mirrors/js/js-lingui 点击查看 免费下载 本篇指南围绕 Lingui 的 li…

📅 2026/10/10 2:44:18
板级适配 · SystemInit:给芯片上电后的“自检程序“动手

板级适配 · SystemInit:给芯片上电后的“自检程序“动手

阅读指引:这篇讲的是芯片"刚通电、还没进 main"那一刻发生的事。下面术语不少(CMSIS、HAL、FPU、VTOR…),但每个术语我都会先给一句直白的解释,跟着读就行,别被缩写吓到。核心其实就一句&#xf…

📅 2026/10/10 2:44:18
MORE NEWS

更多资讯

📰

老游戏低配优化指南:CPU单核与显存管理实战

1. 为什么十几年后还有人折腾这款老游戏每次看到有人问“这游戏都这么多年了,还有必要优化吗”,我都想回一句:你去试试在现在的机器上直接跑原版,看看那个帧数曲线有多酸爽。这款游戏当年是出了名的吃CPU,双核时代它能…

📰

Windows网页长截图实战:用Playwright实现全页高清截图

有时候你会遇到这样一种尴尬:一篇特别重要的网页文章或者产品介绍页,从头到尾几十屏,你想把它完整保存下来,发给别人或者归档留档。普通截图工具截了上半截就没下半截,滚动截屏在微信/手机里倒是能用,但到了…

📰

Beads 测试指南:从 Bazel 门禁到测试设计的完整实践手册

AI 应用Agent 记忆CLIMCP 服务项目管理人工智能 【免费下载链接】beads Beads - A memory upgrade for your coding agent 项目地址: https://gitcode.com/GitHub_Trending/beads1/beads 点击查看 免费下载 本篇技术指南围绕 Beads 开源仓库(engdocs/TE…

📰

问卷星逆向实战:参数复现与会话模拟两种路线全解析

“问卷星逆向”这个话题,常年挂在自动化测试、数据采集、业务流程验证这几类需求下面。你可能是想把自己搭的问卷系统跟问卷星上的公开问卷做数据打通,也可能是想给一套答题系统做接口自动化回归,还可能是需要一个受控的数据采集程序去处理已…

📰

cmux多路复用实战:单端口多协议分发与连接管理

1. 从“cmux”这个名字说起:它到底想解决什么问题第一次看到“cmux”这个词,很多人会下意识地把它拆成“c”和“mux”两部分。mux是multiplexer的缩写,也就是多路复用器,在通信和系统编程里是个老面孔了。而前缀“c”可以有很多种…

📰

CPA侧信道分析实战:泄漏模型选择、波形对齐与攻击验证指南

先说个有点丢人的事。我头一回在自己搭的功耗采集平台上跑Correlation Power Analysis(CPA),用的是 AES-128 第一个 S 盒输出的汉明重量做泄漏模型,跑完相关性矩阵后,最大相关系数对应的密钥字节居然和真实密钥差了整整…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

读完文章,想聊聊您的网站?

告诉我们您的行业与需求,资深顾问一对一梳理方案与报价,全程免费。

📞 💬