☰
使用 Apache Beam 进行 AI/ML 数据探索与数据预处理流水线开发
2026/10/10 13:52:04 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache 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 Runner

2.1 为什么选择 Pandas 风格的 DataFrame API

Pandas,让开发者能在 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以增量的方式分块读取文件;对于不含引号换行的大文件,可以传入splittable=True参数启用基于换行符的动态切分(dynamic splitting)以提升并行度,但注意包含引号换行的记录使用该选项可能造成数据损坏。此外该模块还提供read_json、read_fwf、read_gbq(BigQuery 读取)以及to_csv等读写操作,均支持文件通配模式与任意 Beam 兼容文件系统。

2.2 在 JupyterLab 中交互式探索数据

DataFrame API 可与 Beam Interactive Runner 组合使用。Interactive Runner 是 Beam Python 流水线的交互式执行器,其构造函数定义在 interactive_runner.py,默认以DirectRunner作为底层执行器,支持缓存上次运行计算过的 PCollection(force_compute=False时只计算缺失数据的最小流水线片段)、渲染流水线图(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 数据预处理流水线由以下五个步骤构成:

  1. 读写数据(Read and write):从文件系统、数据库或消息队列中读取与写出数据。Apache Beam 拥有丰富的 内置 IO 连接器,例如本地/云文件系统文本、CSV、Parquet、BigQuery、Kafka、Pub/Sub 等,可无缝对接现有存储与消息基础设施。
  2. 数据清洗(Data cleaning):在数据进入模型之前进行过滤与清洗,例如移除重复或无关数据、纠正数据集中的错误、过滤离群点、处理缺失值。
  3. 数据变换(Data transformations):让数据符合模型训练所期望的输入,例如归一化、独热编码(one-hot encode)、缩放(scale)或向量化(vectorize)。
  4. 数据富集(Data enrichment):结合外部数据源使数据更有意义、更易于模型解释,例如把城市名或地址转换为坐标集合。
  5. 数据校验与指标(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, coordinates=beam.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.Filter):beam.Filter保留谓词返回True的元素。示例中filter_missing_data过滤掉weight为None的记录,这是处理缺失数据的常见策略之一。清洗阶段常见的操作还包括去重(beam.Distinct)、按条件裁剪离群点、字段纠错等,均可通过Filter/FlatMap组合实现。

数据变换(beam.FlatMap):变换阶段采用FlatMap对每条记录做 Min-Max 归一化,将三个数值字段分别缩放到约[0, 1]区间:

  • age:age / 100
  • height:(height - 150) / 50
  • weight:(weight - 50) / 50

这里用yield row保留"一对多"的灵活性——FlatMap返回迭代器,既能做一对一映射,也能在需要时展开为多条输出。除了这种手工缩放,Beam 官方还提供了更专业的 ML 预处理方案MLTransform(见 website/www/site/content/en/documentation/ml/preprocess-data.md),它封装了来自 TensorFlow Transforms(TFT)的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 的价值在于:它为每条数据注入"全体数据集级别"的外部信息,而无需在每条记录内复制这些数据,非常适合地址转坐标、外键关联、词表映射等富集场景。

指标统计(Metrics):Metrics.counter('main', 'counter')创建一个命名计数器(命名空间main、名称counter),count_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 的实现细节
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

Apache Beam is a unified programming model for Batch and Streaming data processing.

项目地址:https://gitcode.com/gh_mirrors/beam4/beam
点击查看免费下载

相关推荐

上一篇:PNChart与CoreGraphics:底层绘制原理深度剖析
下一篇:新贡献者流程(实验版本)

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询