多彩编程 多彩编程MZPH · CODE BLOG
ARTICLE DETAIL

文章详情

深耕前端与后端开发技术的一线实战笔记与踩坑复盘。

Apache Beam YAML 示例目录全解:从 Wordcount 到内置变换的实战入门

Apache Beam YAML 示例目录全解:从 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 提供了一套基于 YAML 声明式配置的 Pipeline 描述方式让开发者无需编写 Python 代码即可定义完整的数据处理流水线。本篇文章以仓库中 sdks/python/apache_beam/yaml/examples 目录下的示例集为核心系统讲解 Beam YAML 的统一起跑命令、经典 Wordcount 示例的完整拆解以及内置的 Element-wise逐元素与 Aggregation聚合两类变换的配置语法。读完本文你将能够独立阅读、运行并仿写这些 YAML 示例掌握MapToFields、Filter、Explode、Combine等核心变换的字段级配置方法并理解这些示例如何通过仓库内的测试框架被自动验证。一、示例目录概览与统一运行方式examples目录是整个 Beam YAML 框架的官方示例集合它既是初学者的入门教材也是框架内置变换用法的活字典。目录结构如下sdks/python/apache_beam/yaml/examples/ ├── wordcount_minimal.yaml # 入门首选单词计数 ├── simple_filter.yaml # 读 CSV → 过滤 → 写 CSV ├── simple_filter_and_combine.yaml # 过滤 聚合组合 ├── regex_matches.yaml # 正则匹配 ├── transforms/ │ ├── elementwise/ # 逐元素变换 │ │ ├── map_to_fields_callable.yaml │ │ ├── filter_callable.yaml │ │ └── explode.yaml │ └── aggregation/ # 聚合变换 │ ├── combine_sum.yaml │ ├── combine_sum_minimal.yaml │ ├── combine_count_minimal.yaml │ ├── combine_mean_minimal.yaml │ ├── combine_max_minimal.yaml │ ├── combine_min_minimal.yaml │ ├── combine_multiple_aggregations.yaml │ ├── top_largest_per_key.yaml │ ├── top_smallest_per_key.yaml │ └── group_into_batches.yaml └── testing/ └── examples_test.py # 示例自动化验证测试这些示例全部通过同一个命令行入口运行。根据 README.md 的说明任何 YAML 示例都可以用如下命令执行python -m apache_beam.yaml.main --pipeline_spec_file/path/to/example.yaml其中apache_beam.yaml.main是 Beam YAML 框架的命令行入口模块--pipeline_spec_file指定要加载的 YAML 文件路径。运行前需要确保环境中已安装 Apache Beam Python SDK即当前仓库的sdks/python目录。需要注意的是wordcount_minimal.yaml、simple_filter.yaml等示例默认读取存储在 Google Cloud Storage 上的公开数据文件gs://前缀直接运行需要配置 Google Cloud 凭据Application Default Credentials。更简单的本地调试方式是把ReadFromText/ReadFromCsv中的path改为本地文件路径或参考下文将要介绍的测试框架自动注入本地输入。二、WordcountBeam YAML 的“Hello World”wordcount_minimal.yaml 被 README 明确推荐为入门起点。它读取一个文本文件按词切分、按词分组并统计每个词的出现次数——这是所有 Beam SDK 共有的经典示例完整展示了 Beam YAML 的多项核心能力。完整配置如下pipeline: type: chain transforms: # 读取文本文件输出一个字段 line 的行集合 - type: ReadFromText config: path: gs://dataflow-samples/shakespeare/kinglear.txt # 将每行的 line 字段切分为单词列表 - type: MapToFields config: language: python fields: words: callable: | import re def my_mapping(row): return re.findall(r[A-Za-z\], row.line.lower()) # 将每个单词列表展开为独立行 - type: Explode config: fields: words # 把字段重命名为 word - type: MapToFields config: fields: word: words # 按 word 分组并统计计数 - type: Combine config: language: python group_by: word combine: count: value: word fn: count # 输出结果 - type: LogForTesting2.1 流水线骨架type: chain整个 Pipeline 由type: chain声明表示各变换按列表顺序串联执行上一个变换的输出自动作为下一个变换的输入。这是 Beam YAML 中最常用、最易读的流水线组织方式simple_filter.yaml则展示了另一种方式——通过input字段显式指定上游变换的输出适合构造非线性的流水线拓扑。2.2 每个变换的职责拆解ReadFromText读取path指定的文本文件把每一行包装成包含单个字段line的 Row。MapToFields Python callable通过language: python声明使用 Python 表达式fields中words字段的callable是一段完整的 Python 函数。该函数接收一个row用正则[A-Za-z]将row.line转小写后切分为单词列表。Explodefields: words将每个单词列表炸开为多个独立 Row每个 Row 只含一个单词。这是列表字段扁平化的关键操作。MapToFields字段重命名word: words这种输出字段名: 输入字段名的简写形式完成字段重命名把words改名为word。Combinegroup_by: word指定分组键combine下声明聚合结果字段count其value: word表示对word字段做聚合fn: count指定聚合函数为计数。这是按 key 分组 多指标聚合的声明式写法。LogForTesting把结果 Row 逐条打印到日志便于本地观察输出也用于测试断言。2.3 预期输出示例自带的验收标准该 YAML 文件末尾以注释形式给出了预期输出这是 Beam YAML 示例约定俗成的自带测试断言# Expected: # Row(wordking, count311) # Row(wordlear, count253) # Row(worddramatis, count1) # Row(wordpersonae, count1) # Row(wordof, count483) # Row(wordbritain, count2)即使没有真实运行从配置结构也能推演出数据流kinglear.txt文本 → 小写分词 → 单词展开 → 按词分组计数 → 输出形如Row(wordking, count311)的结果。三、Element-wise 变换逐元素处理数据README 指出逐元素Element-wise示例主要演示内置的映射类变换MapToFields、Filter和Explode这三个变换分别对应投影/计算新字段、过滤行、展开列表字段三种常见操作。对应示例位于 transforms/elementwise 目录。3.1MapToFields字段投影与 Python 表达式map_to_fields_callable.yaml 演示了用MapToFields对每一行执行字段级计算。示例先用Create变换在内存中构造数据elements列表每行是带#前缀的字符串再对每个元素做字符串清理pipeline: type: chain transforms: - type: Create name: Gardening plants config: elements: - # Strawberry - # Carrot # ... 更多元素 - type: MapToFields name: Strip header config: language: python fields: element: callable: lambda row: row.element.strip(# \\n) - type: LogForTesting要点Create变换用于在 YAML 内直接内联构造测试数据是纯本地调试的利器其输出字段由elements中第一个元素的键自动推断。MapToFields的fields定义输出字段名elementcallable可以是单行 lambda如本示例或多行 Python 函数如 Wordcount 中的my_mapping。language: python声明表达式语言除 Python 外Beam YAML 还支持 SQL 等语言后端。默认情况下MapToFields只输出fields中声明的字段若想保留原行其他字段可像combine_multiple_aggregations.yaml中那样设置append: true。3.2Filter按条件保留行filter_callable.yaml 演示用Filter按条件过滤行。示例构造了 5 种植物每种带icon、name、duration字段只保留多年生perennial植物- type: Filter name: Filter perennials config: language: python keep: callable: lambda plant: plant.duration perennialkeep是 Filter 的条件入口布尔表达式为真则保留该行。预期输出只包含三种durationperennial的植物。值得注意的是simple_filter.yaml中还展示了一种更简洁的写法——keep直接写成裸表达式字符串category Electronics配合ReadFromCsv读取真实业务数据交易记录把category为 Electronics 的三条记录筛选出来并写入 CSV形成读-过滤-写的完整数据管道pipeline: transforms: - type: ReadFromCsv name: ReadInputFile config: path: gs://apache-beam-samples/beam-yaml-blog/products.csv - type: Filter name: FilterWithCategory input: ReadInputFile config: language: python keep: category Electronics - type: WriteToCsv name: WriteOutputFile input: FilterWithCategory config: path: output该示例同时展示了ReadFromCsv/WriteToCsv两个 I/O 变换以及通过nameinput显式连接变换的非 chain 写法。3.3Explode列表字段展开explode.yaml 演示如何把一行中的列表字段拆成多行。Create构造了两条园艺植物记录每条的produce是包含多个植物的列表- type: Explode name: Flatten lists config: fields: [produce]fields声明要展开的列表字段支持数组形式声明多个字段。展开后每个列表元素成为一行其余字段如season原样保留——从预期输出看produce[Strawberry, Potato]的一行变成了produceStrawberry与producePotato两行season字段保持不变。这正是 Wordcount 中把每行单词列表拆成每行一个单词的机制。四、Aggregation 变换Combine聚合全家桶README 指出聚合示例统一使用内置Combine变换完成 sum、mean、count 等基础聚合。这些示例位于 transforms/aggregation 目录覆盖了从单聚合到多聚合、从简单函数到自定义 CombineFn 的完整梯度。4.1 多指标聚合一次 Combine 输出 sum 与 meancombine_sum.yaml 是 Combine 用法的典型代表。它构造了 5 条食谱用料记录含recipe、fruit、quantity、unit_price字段按fruit分组后同时计算数量总和与价格均值- type: Combine name: Sum values per key config: language: python group_by: fruit combine: total_quantity: value: quantity fn: sum mean_price: value: unit_price fn: mean这里combine映射的每个键total_quantity、mean_price都是一个新输出字段各自独立声明聚合源字段value与聚合函数fn。预期输出验证了计算blueberry 的total_quantity312、mean_price2.5(2.003.00)/2与输入数据完全吻合。fn支持的内置函数包括sum、mean、count、min、max等仓库中对应的 combine_sum_minimal.yaml、combine_count_minimal.yaml、combine_mean_minimal.yaml、combine_max_minimal.yaml、combine_min_minimal.yaml 分别给出了每个函数的极简独立示例。4.2 极简写法直接以聚合函数作为字段值combine_count_minimal.yaml 展示了 Combine 的最简语法——当对分组键本身计数时combine的值可以直接写函数名- type: Combine name: Shortest names per key config: language: python group_by: season combine: produce: count这等价于按season分组统计每组行数输出字段produce。预期输出spring4, summer3, fall2, winter1与 10 条输入数据的分组结果一一对应。这种省略value/fn结构的糖衣语法让简单计数一目了然。4.3 全局聚合技巧伪造全局 keycombine_multiple_aggregations.yaml 演示了一个实用技巧Combine 必须指定group_by若要模拟不分组的全局聚合可以先用MapToFields给每行附加一个常量全局键聚合完成后再用MapToFields把它剔除# Simulates global GroupBy by creating global key - type: MapToFields name: Add global key config: append: true language: python fields: global_key: 0 - type: Combine name: Get multiple aggregations per key config: language: python group_by: global_key combine: min_price: value: unit_price fn: min mean_price: value: unit_price fn: mean max_price: value: unit_price fn: max # Removes unwanted global key - type: MapToFields name: Remove global key config: fields: min_price: min_price mean_price: mean_price max_price: max_price这个示例还同时展示了MapToFields的两个进阶用法append: true表示在原字段基础上追加新字段而非替换最后的MapToFields通过输出字段: 输入字段的映射实现字段裁剪仅保留三个聚合结果字段。最终输出单行Row(min_price1.0, mean_price2.5, max_price4.0)。4.4 自定义 CombineFnTop N 与分批聚合函数不仅限于内置的简单函数还可以直接引用 Beam Python SDK 的完整 CombineFn 类。top_largest_per_key.yaml 用apache_beam.transforms.combiners.TopCombineFn求每个 key 下最大的 2 个值- type: Combine config: language: python group_by: produce combine: biggest: value: amount fn: type: apache_beam.transforms.combiners.TopCombineFn config: n: 2fn从函数名字符串升级为对象描述type指向 SDK 中真实的 CombineFn 类路径config传入构造参数n: 2表示取 Top 2。预期输出中组得到[3, 2]组得到[5, 4]验证了按值降序取前 N 个的逻辑。group_into_batches.yaml 则展示了同一个TopCombineFn的另一种应用——利用其收集前 N 个元素的能力把每组元素聚合成固定大小的批次列表n: 3输出形如Row(seasonspring, produce[, , ])的行。仓库中另有 top_smallest_per_key.yaml 与 top_largest 形成对照。五、示例的可验证性仓库内置测试框架这些 YAML 示例并非孤立文件仓库通过 testing/examples_test.py 对它们进行了自动化验证这保证了示例的可运行性与输出正确性。测试框架的核心机制值得了解自动发现YamlExamplesTestSuite用glob扫描../*.yaml根示例、../transforms/elementwise/*.yaml逐元素示例和../transforms/aggregation/*.yaml聚合示例为每个 YAML 文件动态生成一个test_文件名测试方法。期望输出解析测试读取 YAML 文件定位# Expected:标记之后的所有注释行将其解析为期望的 Row 输出如果示例缺少该标记测试会直接抛出ValueError——这正是上一节所说自带验收标准被机器强制执行的方式。输入重定向通过register_test_preprocessor注册的预处理函数会把ReadFromText/ReadFromCsv中的gs://路径替换为本地临时文件。例如 Wordcount 测试会按期望输出的词频分布随机生成kinglear.txt写入本地simple_filter测试则注入内置的products.csv交易数据从而做到离线可测。执行与断言测试用yaml_transform.expand_pipeline把 YAML 规格展开为真实 Beam Pipeline 并运行最后用assert_matches_stdout比对实际输出与期望输出。从源码结构可以看出Beam YAML 示例与框架本身apache_beam/yaml包紧密耦合YAML 规格由yaml_transform解析展开SafeLineLoader负责安全加载 YAML 并剥离元数据apache_beam.yaml.main是命令行入口。因此这些示例同时充当了框架的回归测试集任何破坏示例输出的改动都会在 CI 中被捕获。六、从示例到实战Beam YAML 学习路径建议综合 README 的分类与示例布局推荐的学习路径如下入门先运行并逐行理解 wordcount_minimal.yaml掌握chain骨架、ReadFromText、MapToFields、Explode、Combine、LogForTesting六个基础构件。逐元素变换依次阅读 map_to_fields_callable.yaml、filter_callable.yaml、explode.yaml分别掌握字段计算、行过滤、列表展开三大操作。聚合变换从极简的 combine_sum_minimal.yaml 到多指标聚合 combine_sum.yaml再到自定义 CombineFn 的 top_largest_per_key.yaml循序渐进掌握Combine的全部语法层次。真实数据管道对照 simple_filter.yaml把学到的变换与ReadFromCsv/WriteToCsv组合成面向真实文件的数据处理流程。在动手仿写自己的 YAML 流水线时记住三个关键语法点type: chain声明串联执行顺序language: python启用表达式/函数后端callable支持 lambda 与多行函数combine下每个字段独立声明value聚合源与fn聚合函数可为内置名或 SDK CombineFn 类。只要遵循这些模式你就能把常见的读数据 → 清洗/过滤/展开 → 分组聚合 → 输出数据处理需求用纯 YAML 快速表达出来。赞分享批处理流处理大数据【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam15/beam点击查看免费下载相关推荐Apache Beam YAML 示例目录实战指南从 Wordcount 到 Kafka、Iceberg 与 ML 管道Apache Beam YAML 示例目录实战指南从 Wordcount 到 Kafka、Iceberg 与 ML 管道 本文基于 Apache Beam 仓大数据批处理流处理数据工程Apache Beam Kotlin 示例管线全解析从 WordCount 四部曲到 Cookbook 实战模式Apache Beam Kotlin 示例管线全解析从 WordCount 四部曲到 Cookbook 实战模式 Apache Beam 的 Kotlin 示批处理流处理大数据Apache Beam Python DataFrame API 示例管道实战从 Wordcount 到航班延误分析Apache Beam Python DataFrame API 示例管道实战从 Wordcount 到航班延误分析 导读 Apache Beam Pytho批处理流处理大数据上一篇flutter-examples中的疫情追踪应用covid19_mobile_app功能详解下一篇敏捷方法论知识体系实战指南基于 linkedin-skill-assessments-quizzes 的 140 题全景梳理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表