【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
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 任务的文件结构"。
环境与依赖准备
在开始搭建之前,需要确认两件事:
- Python 解释器:课程声明
programming_language_version: 3.0,即需要 Python 3 环境(建议使用较新的 3.8+ 稳定版本)。 - 依赖包:课程根目录的 requirements.txt 指定了依赖版本:
apache-beam==2.41.0 apache-beam[test]==2.41.0 pytz~=2022.2.1apache-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 插件)按以下顺序操作:
- 创建新项目:选择 "Create New Project"(新建项目),并在项目路径处选择本目录
learning/katas/python(即本文档所在目录,包含course-info.yaml的目录)。 - 选择解释器:选择你的项目解释器(例如 virtualenv 虚拟环境),然后点击 "Create"(创建)。
- 确认从已有源码创建:当 PyCharm 提示 "Create project from existing sources"(从现有源码创建项目)时,点击 "Yes"(是)。
- 等待索引完成:等待 PyCharm 完成项目索引(Indexing),期间不要中断。
- 配置项目解释器:打开 "Preferences"(设置),搜索 "Project Interpreter"(项目解释器),按需选择/添加解释器(例如 virtualenv),确保与第 2 步一致。
- 切换到 Course 视图:打开 "Project"(项目)工具窗口,选择 "Course"(课程)视图。这是 EduTools 插件提供的专用视图,会以"章节 → 练习"的树形结构展示整门课程,并显示每道题的完成状态。
- 开始练习:项目就绪后,即可在 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.
相关推荐
Apache Beam Java Katas 实战指南:用 IntelliJ Education 搭建交互式 Beam 学习环境
Apache Beam Java Katas 实战指南:用 IntelliJ Education 搭建交互式 Beam 学习环境 Apache Beam 官方仓
批处理流处理大数据使用 IntelliJ EduTools 搭建 Apache Beam Katas Kotlin 交互式练习环境
使用 IntelliJ EduTools 搭建 Apache Beam Katas Kotlin 交互式练习环境 Beam Katas Kotlin 是 Apa
Audacity 4 选型指南:免费开源多轨音频编辑器与七个实战场景
Audacity 4 选型指南:免费开源多轨音频编辑器与七个实战场景 Audacity 是一款完全免费、开源的多轨音频编辑器,录音、降噪、混音、导出在一条流水线
音频处理桌面应用音视频
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考