news 2026/10/9 12:28:15

Apache Beam Python Katas 上手指南:用 PyCharm Education 搭建交互式 Beam 练习环境

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Beam Python Katas 上手指南:用 PyCharm Education 搭建交互式 Beam 练习环境

【免费下载链接】beam

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

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

Beam Katas 是 Apache Beam 官方仓库中内置的一套交互式编码练习(Code Katas),它基于 JetBrains Educational Products 构建,帮助学习者通过"边读题目、边补代码、边跑测试"的方式逐步掌握 Beam 的批处理与流处理编程模型。本指南以 learning/katas/python/README.md 的官方设置流程为主线,结合仓库中真实的课程目录、任务文件与测试代码,完整讲解如何在 PyCharm Education 中搭建 Python 版 Katas 练习环境,读完即可开始逐题闯关。

Beam Katas 是什么

在 learning/katas/README.md 中对 Beam Katas 有明确定义:Beam Katas 是交互式的 Beam 编码练习(即 code katas),其目标是提供一系列结构化的动手学习体验,让学习者通过解决难度逐渐递增的练习来理解 Apache Beam 及其 SDK。它基于 JetBrains Educational Products 构建,并且同时提供Java、Python 和 Go三种 SDK 版本,本指南聚焦其中的 Python 版本(目录位于 learning/katas/python)。

Katas 的核心理念来自经典编程练习法:重复练习、小步递进。每道题(Kata)都给出一个最小化的 Beam 编程任务,例如"用beam.Create创建一个硬编码元素"、"读取文本文件并把国名转成大写"、"实现一个统计词频的 Pipeline",你需要在给定的骨架代码中补齐关键逻辑,然后由仓库内置的单元测试自动校验结果。

Python Katas 课程全景

Python 版 Katas 的课程定义在 course-info.yaml 中:

type: marketplace title: Beam Katas - Python language: English programming_language: Python programming_language_version: 3.0 environment: unittest content: - Introduction - Core Transforms - Common Transforms - IO - Streaming - Examples

课程由六个章节(Section)组成,难度从入门到综合逐步递进:

章节主题覆盖内容
Introduction入门Hello Beam:创建你的第一条 Pipeline
Core Transforms核心变换Map、FlatMap、ParDo、GroupByKey、CoGroupByKey、Combine、Flatten、Partition、Side Input、Side Output、Branching、Composite Transform
Common Transforms常用变换Filter、WithKeys,以及 Count、Largest、Mean、Smallest、Sum 等聚合操作
IO数据读写TextIO 的 ReadFromText(从文本文件读取 PCollection)
Streaming流处理Timestamps(Add Timestamps)、Windows(Fixed Windows)、Triggers(Early / Event Time / Window Accumulation Modes)
Examples综合示例Word Count 词频统计完整 Pipeline

各章节的内容清单可以在对应的section-info.yaml中确认,例如 Core Transforms/section-info.yaml 列出了 Map、GroupByKey、CoGroupByKey、Combine、Flatten、Partition、Side Input、Side Output、Branching、Composite Transform 共 10 个练习组。

目录结构约定

每个练习(Lesson)都是一个独立的子目录,包含固定的一组文件。例如Introduction/Hello Beam/Hello Beam/目录下:

Hello Beam/ ├── __init__.py ├── task-info.yaml # EduTools 任务元数据:文件可见性、占位符位置 ├── task.md # 题目说明与提示 ├── task.py # 待你补全的骨架代码 └── tests/ ├── __init__.py └── test_task.py # 校验答案的单元测试

理解这一结构是完成练习的前提,详见下文"一个 Kata 任务的文件结构"。

环境与依赖准备

在开始搭建之前,需要确认两件事:

  1. Python 解释器:课程声明programming_language_version: 3.0,即需要 Python 3 环境(建议使用较新的 3.8+ 稳定版本)。
  2. 依赖包:课程根目录的 requirements.txt 指定了依赖版本:
apache-beam==2.41.0 apache-beam[test]==2.41.0 pytz~=2022.2.1
  • apache-beam==2.41.0:Beam Python SDK 本体,Katas 骨架代码(如import apache_beam as beam)依赖它;
  • apache-beam[test]==2.41.0:SDK 的测试扩展,为 Katas 的单元测试环境提供支持(environment: unittest);
  • pytz~=2022.2.1:时间处理依赖,流处理相关的练习题(Timestamps、Windows、Triggers)会用到。

在 PyCharm 中创建项目时选择虚拟环境(virtualenv / venv)后,可在 PyCharm 的 Terminal 中执行pip install -r requirements.txt安装全部依赖;PyCharm 通常也会在检测到requirements.txt时主动提示安装。

分步搭建:导入 Katas 项目(官方流程)

根据 learning/katas/python/README.md 的官方说明,设置过程一共分为 7 步。请使用PyCharm Education(或在 PyCharm 中安装EduTools 插件)按以下顺序操作:

  1. 创建新项目:选择 "Create New Project"(新建项目),并在项目路径处选择本目录learning/katas/python(即本文档所在目录,包含course-info.yaml的目录)。
  2. 选择解释器:选择你的项目解释器(例如 virtualenv 虚拟环境),然后点击 "Create"(创建)。
  3. 确认从已有源码创建:当 PyCharm 提示 "Create project from existing sources"(从现有源码创建项目)时,点击 "Yes"(是)。
  4. 等待索引完成:等待 PyCharm 完成项目索引(Indexing),期间不要中断。
  5. 配置项目解释器:打开 "Preferences"(设置),搜索 "Project Interpreter"(项目解释器),按需选择/添加解释器(例如 virtualenv),确保与第 2 步一致。
  6. 切换到 Course 视图:打开 "Project"(项目)工具窗口,选择 "Course"(课程)视图。这是 EduTools 插件提供的专用视图,会以"章节 → 练习"的树形结构展示整门课程,并显示每道题的完成状态。
  7. 开始练习:项目就绪后,即可在 Course 视图中逐题打开任务,阅读题目、补全代码、运行测试。

关于 EduTools 的进一步使用说明,原文档指向了 JetBrains 官方的教育产品帮助页面;你也可以直接参考仓库中其他语言的 Katas 结构(learning/katas/java、learning/katas/go)来横向对照理解课程组织方式。

一个 Kata 任务的文件结构

每一道练习题都由四个核心文件组成,理解它们的分工能让你更快上手:

task.md:题目说明

采用 Markdown 格式,前半部分是知识背景,后半部分是任务描述(Kata: ...)与 Hint 提示。例如 Hello Beam/task.md 先介绍 Beam 是"用于定义批处理和流处理数据并行管线的统一开源模型,可运行于 Apache Flink、Apache Spark、Google Cloud Dataflow 等分布式后端",随后给出第一道题:

Kata:Your first kata is to create a simple pipeline that takes a hardcoded input element "Hello Beam". (创建一个接收硬编码元素 "Hello Beam" 的简单管线。)

Hint 中提示使用beam.Create从内存数据创建 PCollection。

task.py:待补全的骨架代码

预置了 Pipeline 框架,把关键逻辑留空。以 Hello Beam 为例(task.py):

import apache_beam as beam with beam.Pipeline() as p: (p | beam.Create(['Hello Beam']) | beam.LogElements())

注意骨架代码顶部还带有# beam-playground:注释块(如name、description、complexity: BASIC、tags),这是用于将练习同步到 Beam Playground 的元数据,本地练习时无需改动。

task-info.yaml:占位符定义

以 Hello Beam/task-info.yaml 为例:

type: edu files: - name: task.py visible: true placeholders: - offset: 1157 length: 26 placeholder_text: TODO() - name: tests/test_task.py visible: false - name: __init__.py visible: false - name: tests/__init__.py visible: false

它告诉 EduTools:task.py是可见的待编辑文件,其中从offset(字符偏移 1157)开始、长度为 26 的区间是需要学习者替换的占位符(默认内容是TODO());tests/test_task.py与__init__.py均不可见(visible: false),即隐藏实现细节,避免直接给出答案。

tests/test_task.py:答案校验器

每道题都配套一个基于 Python 标准库unittest的测试文件。以 Hello Beam/tests/test_task.py 为例:

import unittest from test_helper import get_file_output, test_is_not_empty class TestCase(unittest.TestCase): def test_not_empty(self): self.assertTrue(test_is_not_empty(), 'The output is empty') def test_output(self): output = get_file_output(path='task.py') # Remove warning line about docker and Python versions output = [x for x in output if not x.startswith("WARNING")] self.assertIn('Hello Beam', output, 'The input element should contain "Hello Beam".') self.assertEqual(1, len(output), 'The output should contain a single element.')

它做了三件事:检查task.py输出非空;执行task.py并断言输出包含Hello Beam;断言输出恰好只有 1 个元素(即没有多打内容)。只有通过这些断言,练习才算通过。

test_helper.py:测试工具函数

课程根目录的 test_helper.py 提供了三个被测试文件复用的工具:

  • get_file_text(path):读取指定文件文本;
  • get_file_output(path, arg_string="", encoding="utf-8"):用subprocess以sys.executable执行目标脚本,并把标准输出按行解析成字符串列表(支持通过arg_string向脚本标准输入写入参数);
  • test_is_not_empty():基于get_file_text判断文件非空。

从源码实现看,Katas 的测试机制是直接运行你的task.py并解析其标准输出,因此你的补全必须让 Pipeline 正确运行并打印出预期结果(beam.LogElements()负责把元素写入日志/标准输出)。

代表性任务实战剖析

下面挑选三个不同章节的练习,展示从题目到解法的完整思路。

入门:Hello Beam(Introduction)

上文已展示完整骨架。这道题的考点是:beam.Pipeline()的上下文管理器用法、beam.Create从内存数据创建 PCollection、beam.LogElements()打印元素。完成它即可理解 Beam Pipeline 的最小闭环:创建 → 变换 → 输出。

综合:Word Count(Examples)

Word Count/task.py 是一个完整的词频统计 Pipeline,几乎串联了课程中所有核心变换:

import apache_beam as beam lines = [ "apple orange grape banana apple banana", "banana orange banana papaya" ] with beam.Pipeline() as p: (p | beam.Create(lines) | beam.FlatMap(lambda sentence: sentence.split()) | beam.combiners.Count.PerElement() | beam.MapTuple(lambda k, v: k + ":" + str(v)) | beam.LogElements())
  • beam.Create(lines):创建输入 PCollection;
  • beam.FlatMap(...):把每行按空格拆成单词并扁平化(一行输入可产出多个元素);
  • beam.combiners.Count.PerElement():统计每个单词出现次数,输出(word, count)键值对;
  • beam.MapTuple(...):把元组映射为"word:count"字符串,便于输出。

这道题对应 Katas 中"Combiners / count / map / combine"等标签,是检验你对 Map 家族、Combine 家族掌握程度的经典综合题。

数据读取:ReadFromText(IO)

ReadFromText/task.md 讲解的是:当 Pipeline 需要从外部源(文件、数据库)读取数据时,Beam 提供了一系列内建的 read/write 变换;读取一个或多个文本文件为 PCollection 的标准做法是使用beam.io.ReadFromText并指定文件路径。任务描述:

Kata:Read the 'countries.txt' file and convert each country name into uppercase. (读取countries.txt文件,并把每个国家名转换为大写。)

练习目录内预置了数据文件 countries.txt,解法思路为:p | beam.io.ReadFromText('countries.txt')读出每行国名,再接一个beam.Map(str.upper)完成大写转换。若内建变换无法覆盖某种存储格式,task.md 也提示可以自行实现 read/write 变换。

流处理:Windows 与 Triggers(Streaming)

Streaming 章节覆盖三类流处理主题:Timestamps(Add Timestamps)练习如何为元素附加时间戳;Windows(Fixed Windows)练习固定窗口聚合;Triggers 章节包含 Early Triggers、Event Time Triggers、Window Accumulation Modes 三道题,练习目录内预置了 generate_event.py 这类事件生成脚本,用于模拟带时间戳的流式数据。这部分是理解 Beam 流处理窗口与触发语义的关键实践。

常见问题与排查建议

  • 运行测试报错找不到test_helper:测试通过from test_helper import ...导入课程根目录的工具模块,请确保以learning/katas/python作为项目根目录导入,并保持目录结构不被改动。
  • 输出与预期不符:从 test_helper.py 的实现可以看到,判定依据是task.py的标准输出行。检查你是否使用了beam.LogElements()(或等价的打印方式),且没有额外输出多余元素。
  • 依赖缺失:先确认虚拟环境已安装 requirements.txt 中的apache-beam==2.41.0与apache-beam[test]==2.41.0;注意[test]扩展依赖是测试环境正常运行的前提。
  • Course 视图不显示课程:确认导入的是包含course-info.yaml的目录(即learning/katas/python),而不是某个练习子目录;同时确认 PyCharm 版本为 Education 版或已安装 EduTools 插件。

下一步学习路径

完成 Python Katas 后,你可以:

  • 横向对比 learning/katas/java 与 learning/katas/go 中同名练习,理解同一概念在不同 SDK 下的表达差异;
  • 将 Katas 题目与仓库源码对应学习,例如在 sdks/python/apache_beam/transforms 中查看beam.Create、FlatMap、Combine等变换的真实实现;
  • 把练习代码同步到 Beam Playground(利用task.py顶部beam-playground注释块),在没有本地环境时在线运行调试。

Beam Katas 的价值在于"用测试驱动的节奏消化概念":每通过一题,你就把 Map、GroupByKey、Combine、窗口、触发器等核心机制在真实代码中亲手验证了一遍,这是阅读文档无法替代的体验。现在,按照本文的 7 步流程导入项目,从你的第一个 "Hello Beam" 开始吧。

【免费下载链接】beam

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

项目地址:https://gitcode.com/gh_mirrors/beam18/beam
点击查看免费下载
上一篇:Composer 提及弹窗(`@` 文件引用 / `/` 技能)在 Apache Maka 中的设计与实现全解析
下一篇:Lightdash 地图可视化的 GeoJSON 数据源:内置区域边界数据与自定义地图配置指南

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

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/9 12:26:22

网络安全系统上线安全检测与安全措施有效性验证报告模板:五段式结构与WAF绕过验证实战

简介:这份《系统上线安全检测和安全措施有效性验证报告模板》面向网络安全评估、系统运维、应用开发及安全合规管理人员,尤其适合参与系统上线前安全评审的技术人员使用。模板围绕网络安全技术、API接口安全、网站应用IPv6支持度三大方向,提供…

作者头像 李华
网站建设 2026/10/9 12:25:45

充电宝危险品识别工程实战:SSD300与样本不均衡处理全解析

简介:这是一套面向毕业设计/课程设计的机器学习危险物品识别项目,聚焦充电宝检测场景。项目完整交付代码与数据集,训练集覆盖带电芯充电宝与不带电芯充电宝两类样本,按1:10比例分布(500:5000),并…

作者头像 李华
网站建设 2026/10/9 12:20:58

AWS EventBridge 事件驱动架构实战:从同步雪崩到事件路由解耦

从一次凌晨三点的告警风暴说起。某个支付平台在上线前一天晚上,下游订单服务的状态变更像推倒了多米诺骨牌一样,一路击穿库存、账单、通知、对账等多个服务。所有团队都在抢修,但根因并不复杂:订单完成这个业务动作,被…

作者头像 李华
网站建设 2026/10/9 12:20:24

基于Java+MySQL的会议预约管理系统数据库课程设计

简介:一款面向数据库课程设计的会议预约管理系统完整资源包,以Java语言结合MySQL数据库和Swing图形界面实现,适合高校学生作为课程设计参考或二次开发的起点。系统覆盖会议预约的核心业务,从前端操作界面到后端数据处理均有完整源…

作者头像 李华
网站建设 2026/10/9 12:16:19

前端打包工具核心原理与选型指南:从依赖图到Tree Shaking

1. 打包工具到底在解决什么问题前端打包工具这个概念,刚入行的朋友经常把它和构建工具、脚手架混为一谈。我刚开始写页面那会儿,也觉得这些东西离自己很远——不就是写几个HTML、CSS、JS文件,浏览器直接打开就能跑吗?直到项目里模…

作者头像 李华