Apache Beam 机器学习数据预处理实战:从基础管线到 MLTransform 与 DataFrame API
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 为批流一体的数据处理提供了统一编程模型同样也是构建机器学习ML数据管线的理想框架。本篇技术指南以 learning/prompts/documentation-lookup-nolinks/51_ml_data_preprocessing.md 为核心骨架讲解如何在 Apache Beam 中完成数据探索、预处理、后处理与数据校验四个关键阶段并通过完整的可运行示例覆盖读取写入、清洗、转换、富集、度量与校验的典型预处理流程随后深入MLTransform类与 Beam DataFrames API帮助你掌握训练与推理阶段复用同一套预处理逻辑、以及用 pandas 风格代码交互式探索数据的实战能力。AI/ML 项目中数据处理的四个阶段在 Apache Beam 中机器学习项目的数据处理工作可以归纳为四个相互衔接的阶段数据探索Data exploration分析并理解数据集内的特征、模式与分布洞察不同变量之间的相互关系为后续特征工程提供依据。数据预处理Data preprocessing对原始数据执行清洗、转换与准备使其满足机器学习算法的输入要求。数据后处理Data postprocessing在模型推理输出之后对结果施加额外转换使其更易解释与阅读。数据校验Data validation评估数据的质量、一致性与正确性确保数据满足特定标准或准则适合目标分析或应用。Beam 的 PTransform 抽象天然支持将这些阶段全部实现为管线中的有向无环图DAG节点数据可以在同一套执行引擎如 Direct Runner、Flink、Spark、Dataflow上以批或流的方式流动。典型数据预处理管线的五个构成步骤一个典型的 Beam 数据预处理管线通常包含以下步骤读取与写入数据从各种数据源读取、向各种数据汇写出数据。数据清洗过滤与清理数据去除重复项、纠正错误、处理缺失值或过滤异常值。数据转换对数据进行缩放、编码或向量化为模型输入做准备。数据富集引入外部数据源增强数据集的丰富度与上下文信息。数据校验与指标校验数据质量并计算类别分布等统计指标。完整示例一个覆盖全流程的预处理管线原文档给出了一个实现上述全部步骤的 Python 示例下面将其完整展开并逐段说明每个环节的作用。import apache_beam as beam from apache_beam.metrics import Metrics with beam.Pipeline() as pipeline: # 创建数据 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}])) # 清洗数据过滤缺失值 def filter_missing_data(row): return row[weight] is not None cleaned_data input_data | beam.Filter(filter_missing_data) # 转换数据最小-最大缩放 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) # 富集数据从外部坐标文件按城市查表 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))) # 度量计数器统计处理过的数据条数 counter Metrics.counter(main, counter) def count_data(row): counter.inc() yield row output_data enriched_data | beam.FlatMap(count_data) # 写出数据 output_data | beam.io.WriteToText(output.csv)该管线依次执行的操作创建数据用beam.Create构造三条包含age、height、weight、city字段的示例记录其中第三条的weight为None用于演示缺失值场景。清洗数据beam.Filter(filter_missing_data)保留weight非空的记录即过滤掉 Berlin 这条缺失样本。转换数据beam.FlatMap(scale_min_max_data)对age、height、weight分别做手工的最小-最大缩放使数值落入近似 [0, 1] 区间。富集数据pipeline | beam.io.ReadFromText(coordinates.csv)作为旁路输入side input通过beam.pvalue.AsDict将城市映射为坐标字典在coordinates_lookup中按city查表得到coordinates字段并删除原city列查不到的城市回退到默认值(0, 0)。度量Metrics.counter(main, counter)定义命名计数器count_data每处理一条数据自增一次用于观察管线吞吐。写入数据beam.io.WriteToText(output.csv)将结果写出为文本文件。从源码看清洗与转换的语义上述示例中的beam.Filter与beam.FlatMap是 Apache Beam Python SDK 的核心转换。Filter接收返回布尔值的函数并仅保留返回True的元素FlatMap期望函数通过yield产出零个或多个输出元素因此既可用于一对一映射也可用于一对多展开——scale_min_max_data与coordinates_lookup中yield row的写法正对应这一语义。示例中input_data是dict元素的 PCollection字段通过row[key]访问属于典型的无 Schema 字典数据处理方式若使用带 Schema 的 PCollection则可通过beam.Row与列名属性访问进而与下文介绍的MLTransform的列式处理模型无缝衔接。MLTransform训练与推理一致的专用预处理变换除标准数据处理变换外Apache Beam 还提供了一组面向 ML 预处理的专用变换统一封装在MLTransform类中。该类简化了工作流并通过对训练与推理复用同一套步骤来保证数据一致性——这正是生产级 ML 系统中训练-服务偏斜train/serve skew的关键防线。MLTransform的实现位于 sdks/python/apache_beam/ml/transforms/base.py。从源码看它支持两种工作模式写模式write mode传入write_artifact_locationMLTransform将变换应用到数据并把产生的工件artifacts写入该存储目录。工件既包括施加在数据集上的变换本身也包括变换过程中生成的值例如ScaleTo01产生的min、max以及ScaleToZScore产生的mean、var。若该目录不存在会被自动创建写入会覆盖已有工件因此每个MLTransform实例应使用不同的目录。读模式read mode传入read_artifact_locationMLTransform从该目录读取工件并施加到数据上。此时无需再传入transforms列表因为变换已内嵌在工件中——推理阶段只需要拿到训练阶段产出的工件目录即可复现完全一致的预处理逻辑。源码中明确要求这两个参数只能指定其一base.py且必须指定至少一个读模式下若还传入transforms会直接抛出ValueErrorbase.py。MLTransform的典型用法import apache_beam as beam from apache_beam.ml.transforms import MLTransform from apache_beam.ml.transforms.tft import ScaleTo01, ComputeAndApplyVocabulary # 训练阶段写模式产出工件 with beam.Pipeline() as p: data p | beam.Create([{age: 25, height: 176}, {age: 61, height: 192}]) _ data | MLTransform( write_artifact_locationgs://my-bucket/ml-artifacts, transforms[ ScaleTo01(columns[age, height]), ComputeAndApplyVocabulary(columns[city]), ]) # 推理阶段读模式无需再传 transforms with beam.Pipeline() as p: data p | beam.Create([{age: 30, height: 180, city: Berlin}]) _ data | MLTransform( read_artifact_locationgs://my-bucket/ml-artifacts)transforms列表中的变换按声明顺序依次应用第 i 个变换的输入是第 (i-1) 个变换的输出源码注释明确指出多输入变换暂不支持见 base.py。此外MLTransform还提供with_transform()方法动态追加变换base.py以及with_exception_handling()方法开启异常处理模式将失败样本以 DLQ死信队列形式旁路输出避免单个坏样本拖垮整条预处理链base.py。基于 TFT 的专用预处理模块MLTransform可用于两类任务生成文本嵌入text embeddings以及实现 TensorFlow TransformsTFT库提供的专用处理模块。TFT 相关变换全部定义在 sdks/python/apache_beam/ml/transforms/tft.py 中主要包括变换类功能关键参数与产物ComputeAndApplyVocabulary计算并应用词表将输入值映射为词表索引top_k保留最高频 token 数、frequency_threshold频率阈值、num_oov_buckets越界哈希桶数为 0 时回退到default_value、split_string_by_delimiter按分隔符切分字符串、vocab_filename词表文件名实际保存时会追加列名后缀ScaleToZScorez-score 标准化使数据均值为 0、方差为 1额外产出工件列名_mean与列名_varScaleTo01最小-最大缩放到 [0, 1]额外产出工件列名_min与列名_maxScaleToGaussian缩放为近似标准正态分布均值为 0、方差为 1仅当列存在长尾时启用高斯变换否则等价于 z-score 标准化elementwise控制按元素还是整列计算ApplyBuckets按给定边界映射为桶索引边界须升序input boundaries[0]映射为 0input boundaries[-1]映射为len(boundaries)NaN 映射为len(boundaries)ApplyBucketsWithInterpolation在桶内插值并归一化到 [0, 1]NaN 映射为 0.5Bucketize按分位数自动分桶将数值转换为桶 IDnum_buckets桶数量、epsilon分位数误差容限须大于 0、elementwiseTFIDF对文本列应用 tf-idf 变换输出词表索引与 tfidf 权重两类工件vocab_size缺省时尝试用tft.get_num_buckets_for_transformed_feature推断失败则要求显式指定smoothTrue时对 tf-idf 分数做平滑ScaleByMinMax按最小-最大值缩放见 tft.pyNGrams/BagOfWords生成 n-gram 序列 / 词袋表示用于文本特征化HashStrings将字符串散列为桶适合大规模类别特征DeduplicateTensorPerRow按行对张量去重见 tft.py这些类都继承自TFTOperationtft.py并通过register_input_dtype声明可处理的输入类型如ScaleToZScore、ScaleTo01注册为 float 输入在内部调用tft.*底层函数完成实际计算。文本嵌入连接外部模型服务MLTransform的另一大用途是生成文本嵌入embeddings实现位于 sdks/python/apache_beam/ml/transforms/embeddings 目录内置多种模型接入方式TensorFlow HubTensorflowHubTextEmbeddings与TensorflowHubImageEmbeddingstensorflow_hub.py基于TFModelHandlerTensor在本地加载 TF Hub 模型。Hugging FaceSentenceTransformerEmbeddings使用 Sentence Transformers 本地模型InferenceAPIEmbeddings调用 Hugging Face Inference APIhuggingface.py。OpenAIOpenAITextEmbeddingsopen_ai.py通过远程模型处理器调用 OpenAI 文本嵌入接口。Vertex AIVertexAITextEmbeddings、VertexAIImageEmbeddings与VertexAIMultiModalEmbeddingsvertex_ai.py支持文本、图像与多模态输入。嵌入变换作为MLTransform的transforms列表成员使用与 TFT 变换一样遵循顺序执行与工件产出的统一机制。DataFrame APIpandas 风格的交互式探索与预处理若要以更贴近数据科学工作流的方式探索和预处理 ML 数据集可以借助 Apache Beam Python SDK 提供的 DataFrame API。该 API 构建在 pandas 之上允许使用标准 pandas 命令与数据交互例如import apache_beam as beam from apache_beam.dataframe.io import read_csv with beam.Pipeline() as p: df p | read_csv(dataset.csv) # 标准的 pandas 操作如列选择、分组、缺失值填充 result df[[age, height]].dropna() result | beam.io.WriteToText(preprocessed.csv)从源码看sdks/python/apache_beam/dataframe/io.py 中的read_csv等 I/O 入口将 pandas 的读取函数适配为 Beam 数据源如pd.read_csv对应beam.dataframe.io.read_csv使 DataFrame 可以融入 Beam 管线并与其他 PTransform 互操作。Beam DataFrames 配合 Apache Beam 交互式 Runnerinteractive runner与 JupyterLab 笔记本可以支持迭代式开发与管线图的实时可视化——交互式 Runner 位于 sdks/python/apache_beam/runners/interactive它通过缓存、后台作业与片段化执行让 DataFrame 代码在笔记本中逐步运行、即时反馈非常适合在写完整管线前先做数据探索和原型验证。实践建议训练与推理复用同一份工件把ScaleTo01、ScaleToZScore等变换放进write_artifact_location产出的工件推理服务只读工件目录从机制上杜绝训练/推理预处理不一致。区分手工缩放与 TFT 缩放示例中的手工(x-150)/50属于快速原型做法生产场景应优先使用ScaleTo01/ScaleToZScore/Bucketize它们会把min/max/mean/var、分位边界等统计量固化为工件随推理阶段自动复现。文本特征工程按需选型类别型特征用ComputeAndApplyVocabulary必要时配合top_k、frequency_threshold、num_oov_buckets长文本可用TFIDF、NGrams或BagOfWords语义向量则直接选用各嵌入提供方。异常处理保障鲁棒性开启with_exception_handling()将无法处理的坏样本旁路收集避免一条脏数据中断整条预处理链。迭代探索再固化先用 DataFrame API 交互式 Runner 在 JupyterLab 中探索数据分布与特征关系再将验证过的逻辑固化到MLTransform的变换列表中形成可重复、可版本化的生产预处理管线。以上内容均可在仓库对应源码中进一步验证MLTransform核心实现见 sdks/python/apache_beam/ml/transforms/base.pyTFT 变换见 sdks/python/apache_beam/ml/transforms/tft.py嵌入变换见 sdks/python/apache_beam/ml/transforms/embeddingsDataFrame I/O 见 sdks/python/apache_beam/dataframe/io.py交互式 Runner 见 sdks/python/apache_beam/runners/interactive。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐5分钟掌握WinCDEmuWindows终极免费虚拟光驱解决方案5分钟掌握WinCDEmuWindows终极免费虚拟光驱解决方案 还在为ISO镜像文件无法直接访问而烦恼吗想象一下你下载了一个游戏安装包或者系统镜像却因存储驱动开发Apache Beam AI/ML 流水线实战指南MLTransform 数据预处理与 RunInference 大规模推理Apache Beam AI/ML 流水线实战指南MLTransform 数据预处理与 RunInference 大规模推理 Apache Beam 是一个用大数据批处理流处理数据工程Apache Beam Beam ML 示例 Notebook 体系RunInference、ModelHandler 与 MLTransform 的机器学习实战Apache Beam Beam ML 示例 Notebook 体系RunInference、ModelHandler 与 MLTransform 的机器学习大数据批处理流处理数据工程上一篇Jqwik 项目常见问题解决方案下一篇5 步跑通 Upscayl 免费 AI 图像放大7 个内置模型 4 个关键参数完整指南创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考