尧图网络 高端网站定制 · 原创设计
免费咨询热线
400-888-6620
免费获取方案
使用 Apache Beam 进行 AI/ML 数据探索与数据预处理流水线开发
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 为 AI/ML 项目提供了一套统一的数据处理能力涵盖数据探索Data exploration、数据预处理Data preprocessing、数据后处理Data postprocessing与数据校验Data validation四类典型任务。本文以 Apache Beam 官方文档 website/www/site/content/en/documentation/ml/data-processing.md 为核心讲解如何利用 Beam Python SDK 的DataFrame API与Interactive Runner在 JupyterLab 笔记本中完成交互式数据探索并系统拆解一条覆盖读取、清洗、变换、富集、指标统计与写入全流程的 ML 数据预处理流水线。读完本文你将能够复用探索阶段的 Pandas 风格代码直接构建生产级预处理管道并掌握Metrics计数器、side input富集等 Beam 核心原语在 AI/ML 场景下的实战用法。一、Beam 数据处理的四类任务与两大主题在 AI/ML 项目中Apache Beam 数据处理通常划分为以下四类任务类型说明Data exploration数据探索在项目启动或数据发生变化时了解数据的属性、分布与统计特征Data preprocessing数据预处理变换数据使其满足模型训练所需的输入格式Data postprocessing数据后处理推理完成后将模型输出变换为有意义的业务结果Data validation数据校验检查数据质量发现离群点计算标准差与类别分布从整体上看这些处理可归并为两大主题数据探索与ML 数据流水线后者同时使用预处理与校验。数据后处理与预处理在本质上是类似的仅在于流水线的顺序与类型不同因此官方文档不再单独展开本文同样聚焦前两者。二、初始数据探索DataFrame API Interactive Runner2.1 为什么选择 Pandas 风格的 DataFrame APIPandas让开发者能在 Beam 流水线内使用熟悉的 Pandas 接口。Beam DataFrame API 本质上是 Beam 流水线之上的一个领域特定语言DSL类似于 Beam SQL。它基于 pandas 实现构建pandas 的 DataFrame 方法会在数据集子集上并行执行与原生 pandas 最大的区别在于所有操作都由 Beam API延迟执行deferred以适配 Beam 的并行处理模型参见 与 pandas 的差异。这意味着你可以用标准的 Pandas 命令构建复杂的数据处理流水线而无需显式书写ParDo、CombinePerKey等底层 Beam 原语探索阶段编写的代码可以直接复用到数据预处理流水线中实现一套代码、两处使用在部分场景下DataFrame API 会延迟到向量化的 pandas 实现上执行从而提升流水线效率。从源码实现看DataFrame API 提供了一整套 IO 入口。以read_csv为例其定义位于 sdks/python/apache_beam/dataframe/io.py底层通过 pandas 的pd.read_csv以增量的方式分块读取文件对于不含引号换行的大文件可以传入splittableTrue参数启用基于换行符的动态切分dynamic splitting以提升并行度但注意包含引号换行的记录使用该选项可能造成数据损坏。此外该模块还提供read_json、read_fwf、read_gbqBigQuery 读取以及to_csv等读写操作均支持文件通配模式与任意 Beam 兼容文件系统。2.2 在 JupyterLab 中交互式探索数据DataFrame API 可与 Beam Interactive Runner 组合使用。Interactive Runner 是 Beam Python 流水线的交互式执行器其构造函数定义在 interactive_runner.py默认以DirectRunner作为底层执行器支持缓存上次运行计算过的 PCollectionforce_computeFalse时只计算缺失数据的最小流水线片段、渲染流水线图render_option等能力。在 JupyterLab 笔记本中你可以用ib.collect()或ib.show()将 PCollection 物化出来查看。ib.show()见 interactive_beam.py会临时构建仅包含必要变换的流水线片段运行后以数据表形式可视化支持n最大元素数与duration最大读取时长限制并可开启visualize_data获得数据深入分析与统计概览控件ib.collect()见 interactive_beam.py则将 PCollection 物化为内存中的 DataFrame支持n、duration、raw_records等参数且能识别DeferredDataFrame自动完成到 PCollection 的转换。官方文档给出的数据探索示例可在笔记本中直接运行如下import apache_beam as beam from apache_beam.runners.interactive.interactive_runner import InteractiveRunner import apache_beam.runners.interactive.interactive_beam as ib p beam.Pipeline(InteractiveRunner()) beam_df p | beam.dataframe.io.read_csv(input_path) # 查看列名与数据类型 beam_df.dtypes # 生成描述性统计 ib.collect(beam_df.describe()) # 查看缺失值 ib.collect(beam_df.isnull())这段代码体现了迭代式开发的核心工作流先构建流水线定义再针对中间结果逐一查看确认数据形态后继续下一步骤最终将成熟代码平滑迁移到批处理预处理管道中。2.3 端到端参考示例仓库中提供了完整的端到端示例笔记本 examples/notebooks/beam-ml/dataframe_api_preprocessing.ipynb演示了如何使用 DataFrame API 同时完成数据探索与数据预处理可作为 AI/ML 项目实践的直接参照。三、ML 数据流水线的五个标准步骤一个典型的 ML 数据预处理流水线由以下五个步骤构成读写数据Read and write从文件系统、数据库或消息队列中读取与写出数据。Apache Beam 拥有丰富的 内置 IO 连接器例如本地/云文件系统文本、CSV、Parquet、BigQuery、Kafka、Pub/Sub 等可无缝对接现有存储与消息基础设施。数据清洗Data cleaning在数据进入模型之前进行过滤与清洗例如移除重复或无关数据、纠正数据集中的错误、过滤离群点、处理缺失值。数据变换Data transformations让数据符合模型训练所期望的输入例如归一化、独热编码one-hot encode、缩放scale或向量化vectorize。数据富集Data enrichment结合外部数据源使数据更有意义、更易于模型解释例如把城市名或地址转换为坐标集合。数据校验与指标Data validation and metrics确保数据满足流水线内可校验的特定要求并输出数据指标例如类别分布统计。3.1 完整示例一条覆盖全部步骤的预处理流水线官方文档提供了一个实现以上全部步骤的示例流水线import apache_beam as beam from apache_beam.metrics import Metrics with beam.Pipeline() as pipeline: # 步骤 1入口创建数据 input_data ( pipeline | beam.Create([ {age: 25, height: 176, weight: 60, city: London}, {age: 61, height: 192, weight: 95, city: Brussels}, {age: 48, height: 163, weight: None, city: Berlin}])) # 步骤 2清洗数据——过滤缺失值 def filter_missing_data(row): return row[weight] is not None cleaned_data input_data | beam.Filter(filter_missing_data) # 步骤 3变换数据——Min-Max 缩放 def scale_min_max_data(row): row[age] (row[age]/100) row[height] (row[height]-150)/50 row[weight] (row[weight]-50)/50 yield row transformed_data cleaned_data | beam.FlatMap(scale_min_max_data) # 步骤 4富集数据——通过 side input 加载坐标表 side_input pipeline | beam.io.ReadFromText(coordinates.csv) def coordinates_lookup(row, coordinates): row[coordinates] coordinates.get(row[city], (0, 0)) del row[city] yield row enriched_data ( transformed_data | beam.FlatMap(coordinates_lookup, coordinatesbeam.pvalue.AsDict(side_input))) # 步骤 5指标——使用 Metrics 计数器统计行数 counter Metrics.counter(main, counter) def count_data(row): counter.inc() yield row output_data enriched_data | beam.FlatMap(count_data) # 步骤 1出口写出数据 output_data | beam.io.WriteToText(output.csv)3.2 各步骤的实现要点与源码支撑输入数据beam.Create示例用beam.Create构造了三条用户记录age、height、weight、city四个字段其中第三条记录的weight为None用于演示缺失值场景。实际项目中此处通常替换为各类 IO 读取如 beam.io.ReadFromText 或 DataFrame API 的read_csv。数据清洗beam.Filterbeam.Filter保留谓词返回True的元素。示例中filter_missing_data过滤掉weight为None的记录这是处理缺失数据的常见策略之一。清洗阶段常见的操作还包括去重beam.Distinct、按条件裁剪离群点、字段纠错等均可通过Filter/FlatMap组合实现。数据变换beam.FlatMap变换阶段采用FlatMap对每条记录做 Min-Max 归一化将三个数值字段分别缩放到约[0, 1]区间age:age / 100height:(height - 150) / 50weight:(weight - 50) / 50这里用yield row保留一对多的灵活性——FlatMap返回迭代器既能做一对一映射也能在需要时展开为多条输出。除了这种手工缩放Beam 官方还提供了更专业的 ML 预处理方案MLTransform见 website/www/site/content/en/documentation/ml/preprocess-data.md它封装了来自 TensorFlow TransformsTFT的ScaleTo01、ScaleToZScore、ScaleByMinMax、Bucketize、ComputeAndApplyVocabulary、TFIDF、NGrams等变换并能通过write_artifact_location/read_artifact_location在训练与推理之间复用预处理参数如缩放用的均值、方差保证训练与推理数据预处理的一致性。数据富集side input AsDict富集步骤演示了 Beam 的**旁路输入side input**机制。pipeline | beam.io.ReadFromText(coordinates.csv)读取坐标文件beam.pvalue.AsDict(side_input)将其作为只读字典旁路传入coordinates_lookup函数以城市名作为键查询坐标查不到的取默认值(0, 0)最后删除原始city字段并yield新行。side input 的价值在于它为每条数据注入全体数据集级别的外部信息而无需在每条记录内复制这些数据非常适合地址转坐标、外键关联、词表映射等富集场景。指标统计MetricsMetrics.counter(main, counter)创建一个命名计数器命名空间main、名称countercount_data中调用counter.inc()每行递增一次。Beam Metrics 的实现位于 sdks/python/apache_beam/metrics支持 Counter、Distribution、Gauge 三类指标它们会在流水线执行后被收集并上报到 runner如 Dataflow 监控面板可用于监控数据量、观察类别分布或校验流水线是否按预期处理了全部记录。除计数器外Metrics.distribution可以记录数值的分布最小值/最大值/均值/分位数非常适合在数据校验阶段统计特征字段的取值分布。写出数据beam.io.WriteToText最终结果通过WriteToText写出为 CSV 文件。生产场景可根据数据规模与下游需求替换为其他连接器例如写入 BigQuery、Parquet 或 Kafka。四、实践建议与限制说明探索与生产代码复用在笔记本中用 DataFrame API Interactive Runner 完成探索后将验证过的 DataFrame 代码直接嵌入批处理流水线或通过DataframeTransform封装可显著缩短从探索到上线的周期关于 DataFrame 与 PCollection 的相互转换to_dataframe/to_pcollection可参考 Beam DataFrames 概览。环境要求DataFrame API 需要 Beam Python SDK 2.26.0 及以上版本推荐通过pip install apache_beam[dataframe]安装在 Beam 2.34.0 之后可用分布式 runner 上应保证 worker 与驱动端安装相同版本的 pandas。Interactive Runner 属于实验性模块源码注释中明确标注experimental, no backwards-compatibility guarantees适合开发探索阶段使用。数据校验的进一步深化若需要系统化的数据校验如计算标准差、类别分布、检测离群点可以结合 Metrics 的 Distribution 指标或借助MLTransform的 TFT 变换族在流水线内完成标准化与词表等统计型变换从而把校验与预处理统一到同一条流水线中。适用范围本文的示例流水线基于 Beam 批处理语义若涉及流式数据处理如从 Kafka 持续消费事件进行在线特征计算可参考仓库 sdks/python/apache_beam/io/kafka 相关文档与示例但数据探索阶段的 DataFrame 操作以全局窗口批处理为主要适用场景。五、扩展阅读Beam DataFrames 概览DataFrame API 的安装、用法与 PCollection 互转与 pandas 的差异DataFrame API 与原生 pandas 的行为差异使用 MLTransform 预处理数据基于 TFT 的标准化、分桶、词表等 ML 专用变换与训练/推理工件复用内置 IO 连接器流水线可用的各类读写连接器examples/notebooks/beam-ml/dataframe_api_preprocessing.ipynb数据探索 数据预处理端到端示例笔记本Interactive Runner 源码 与 interactive_beam 模块交互式执行与物化 API 的实现细节赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam AI/ML 流水线实战指南MLTransform 数据预处理与 RunInference 大规模推理Apache Beam AI/ML 流水线实战指南MLTransform 数据预处理与 RunInference 大规模推理 Apache Beam 是一个用大数据批处理流处理数据工程微信支付集成实战基于wechat3 SDK的JSAPI支付开发指南微信支付集成实战基于wechat3 SDK的JSAPI支付开发指南 微信支付作为主流的移动支付方式已成为众多开发者的必备技能。本文将为你介绍如何使用wech如何利用tinygrad数据流水线实现高效数据加载和预处理从理论到实践如何利用tinygrad数据流水线实现高效数据加载和预处理从理论到实践 tinygrad是一个轻量级的深度学习框架它不仅提供了类似于PyTorch的张量操作人工智能深度学习大模型上一篇PNChart与CoreGraphics底层绘制原理深度剖析下一篇新贡献者流程实验版本创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
RELATED

相关推荐

SpringBoot+微信小程序:网络安全科普系统论文转工程实战解析

SpringBoot+微信小程序:网络安全科普系统论文转工程实战解析

简介:面向微信小程序网络安全科普系统的开发需求,这份docx设计文档适用于毕业设计、课程作业或实际科普平台建设的学习者与开发者。系统采用Java语言、MySQL数据库、微信小程序及SpringBoot框架,构建了包含科普知识查阅、案例分析、在线评价交…

📅 2026/10/10 13:47:10
Spring Boot核心配置解析:绑定、多环境与加密实践

Spring Boot核心配置解析:绑定、多环境与加密实践

很多Java开发者第一次用Spring Boot,体验到的第一个“幸福感”就是不用再手写一堆XML了,但紧接着,就会被application.yml里的缩进、绑定规则和多环境切换折腾几回。application.yml这个看似不起眼的文件,其实是整个Spring Boot项目…

📅 2026/10/10 13:47:10
信创文件传输系统有哪些?主流形态、选型要点与避坑指南

信创文件传输系统有哪些?主流形态、选型要点与避坑指南

1. 先搞清楚:信创文件传输系统和日常用的文件传输工具有什么不同?先说个我几年前的真实经历。当时帮一家制造企业做供应链系统改造,对方IT负责人对着市面上七八套文件传输方案来回比较,越比越乱。他问了我一句话:“我们…

📅 2026/10/10 13:47:10
MORE NEWS

更多资讯

📰

Token是什么?NLP、认证、区块链、编译器中的四种含义与边界

第一次被Token这个词搞懵,是某次和同事调试一个跨平台系统。前端同事说Token过期了需要重新登录,算法同事说这段文本Token切得太多,后端同事说Token里的权限信息没带全。三个人用的都是同一个英文单词,但谁也没接住谁的话。那时候…

📰

豆瓣评论情感分析实战:SnowNLP清洗、校准与TF-IDF词云

简介:本资源是一份面向Python初学者与数据挖掘入门者的实战教学材料,聚焦豆瓣电影《肖申克救赎》评论文本的情感分析与可视化任务,通过调用轻量级中文NLP库SnowNLP实现情感倾向判断与词云生成,适用于课程实验、小规模文本分析项目…

📰

专科生论文降AI率实战指南:检测原理与9个工具组合用法

专科生写论文最怕什么?怕AI率。我见过太多同学,用DeepSeek或者豆包生成一段内容,复制粘贴进Word,交上去之后被AIGC检测拦下来,直接退回重写。有些更严重的,被老师约谈,说不清楚就按学术不端处理…

📰

并查集与贪心:破解情侣牵手的最小交换次数

1. 初见765:贪心能过,但我被"为什么正确"问住了我刷并查集(union-find)专项题单时,最先遇到的都是一批"给一堆连接关系,问连通块有几个"的直白题目。直到碰见力扣765情侣牵手&#xff…

📰

重言式判别程序设计:从表达式树到真值表的完整实现

简介:本资源是一份面向计算机专业本科生的数据结构课程设计实践材料,聚焦重言式(逻辑恒等式)的程序化判别实现,解决布尔表达式真值恒定性验证这一典型算法与数据结构综合应用问题。压缩包共4个文件,含2份Wo…

📰

用Python从零复现KCF目标跟踪算法:原理、实现与避坑指南

简介:KCF用Python代码复现.rar是一份基于Python实现的KCF(核相关滤波)目标跟踪算法复现工程,适合计算机视觉入门者、算法复现爱好者以及需要快速搭建跟踪demo的研究人员。压缩包采用rar格式,内含KCFpy-master项目&…

TODAY

今日更新

THIS WEEK

本周精选

THIS MONTH

本月热门

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

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

📞 💬