- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
导读
本文以 Apache Beam 仓库中 Fixed Windows 学习 Kata 为核心,系统讲解固定时间窗口(Fixed Time Windows)的核心概念、Python SDK 中的实现原理,并手把手完成"按 1 天固定窗口统计事件数量"的编程练习。读完本文,你将掌握 FixedWindows 的窗口划分公式、TimestampedValue 构造带时间戳元素的方法,以及通过beam.WindowInto结合Count.PerElement实现按窗口聚合的完整套路。
一、为什么需要窗口:Windowing 的核心思想
在 Beam 模型里,无界数据集(unbounded PCollection,例如持续流入的日志、点击流、传感器读数)整体上是"无限"的,无法像有界数据那样一次性做完聚合。窗口(Windowing)正是解决这一问题的机制:它按照元素自带的时间戳,把 PCollection 细分为一个个逻辑窗口,每个窗口内的元素数量是有限的。
这一点在 Fixed Windows 的 task.md 中有明确表述:
- Windowing 依据每个元素的 timestamp 对 PCollection 进行细分;
GroupByKey、Combine等聚合型 Transform 天然按窗口工作——它们把 PCollection 当作一串连续的、有限的窗口来处理,即便整个集合本身无界;- 任意 PCollection(包括无界 PCollection)都可以被细分为逻辑窗口,每个元素按照窗口函数被分配到一个或多个窗口;
- 分组类 Transform 会按 key + window 双重维度隐式分组。
需要特别强调的是:聚合操作是"按窗口"发生的。GroupByKey隐式地按 key 和 window 对元素分组,这意味着窗口划分会直接改变聚合结果——这正是本 Kata 要验证的核心行为。
二、Beam 提供的四种窗口函数
task.md 列出了 Beam 内置的四类窗口函数:
| 窗口类型 | 特征 |
|---|---|
| Fixed Time Windows | 固定时长、互不重叠的时间区间 |
| Sliding Time Windows | 固定时长但允许重叠,元素可属于多个窗口 |
| Per-Session Windows | 按事件间隙(gap)动态切分,会话式窗口 |
| Single Global Window | 全局单窗口,所有元素归入同一窗口 |
其中固定时间窗口是最简单的一种:它代表数据流中时长一致、互不重叠的时间区间。例如 1 天窗口会把 3 月 1 日的事件全部归入[2020-03-01, 2020-03-02)区间,3 月 5 日的事件归入[2020-03-05, 2020-03-06)。
三、Kata 任务与解题思路
task.md 给出的 Kata 要求如下:
Kata:请按 1 天时长的固定窗口统计发生的事件数量(count the number of events based on fixed window with 1-day duration)。
并给出提示:使用FixedWindows,窗口时长为 1 天(即24 * 60 * 60秒);更详细的背景可参阅 Beam Programming Guide 的 "Fixed time windows" 一节。
本任务在仓库中的配套文件包括:
- 任务描述与骨架:task.md
- 待完成的实现:task.py
- 隐藏的单元测试:tests/test_task.py
- 课程元数据:task-info.yaml,其中
placeholder指向 task.py 中需要补全的TODO()占位
解题链路非常清晰,分三步:
- 用
beam.Create创建一批带时间戳的事件; - 用
beam.WindowInto(window.FixedWindows(24*60*60))把元素划分到 1 天固定窗口; - 用
beam.combiners.Count.PerElement()统计每个窗口内的事件数,并配合beam.LogElements(with_window=True)输出窗口信息。
四、源码级原理:FixedWindows 的窗口划分公式
打开 Python SDK 的窗口实现 sdks/python/apache_beam/transforms/window.py,可以看到FixedWindows(NonMergingWindowFn)的完整实现。它属于NonMergingWindowFn(非合并窗口函数),每个元素恰好被分配到一个时间区间。
窗口的落点由size和offset两个参数决定,docstring 给出了精确公式:
[N * size + offset, (N + 1) * size + offset)其中:
size:窗口时长(秒)。构造时若size <= 0会直接抛出ValueError('The size parameter must be strictly positive.');offset:窗口起点偏移(秒),以 UNIX 纪元t=0为基准,窗口起点为t = N * size + offset。offset必须在[0, size)区间内,超出会被自动归一化(源码中通过Timestamp.of(offset) % self.size取模完成)。
元素到窗口的实际映射发生在assign()方法:
def assign(self, context: WindowFn.AssignContext) -> list[IntervalWindow]: timestamp = context.timestamp start = timestamp - (timestamp - self.offset) % self.size return [IntervalWindow(start, start + self.size)]即:先算出该时间戳所属窗口的起点,再生成左闭右开的IntervalWindow。以本 Kata 为例,size = 86400秒(1 天)、offset = 0,那么 2020-03-05 00:00:00 UTC 的事件会被分配到[2020-03-05, 2020-03-06)窗口。
从源码结构还可以推断,FixedWindows通过to_runner_api_parameter/from_runner_api_parameter与 Runner API 的标准窗口 URN(common_urns.fixed_windows.urn)进行序列化互操作,因此窗口函数可跨 SDK/Runner 传输——这体现了 Beam 可移植性设计的一角。
五、完整实现:带时间戳的事件与逐窗口计数
下面是在 task.py 基础上补全后的完整可运行代码:
from datetime import datetime import pytz import apache_beam as beam from apache_beam.transforms import window with beam.Pipeline() as p: (p | beam.Create([ window.TimestampedValue("event", datetime(2020, 3, 1, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 1, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 1, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 1, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 5, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 5, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 8, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 8, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 8, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), window.TimestampedValue("event", datetime(2020, 3, 10, 0, 0, 0, 0, tzinfo=pytz.UTC).timestamp()), ]) | beam.WindowInto(window.FixedWindows(24*60*60)) | beam.combiners.Count.PerElement() | beam.LogElements(with_window=True))关键点拆解:
window.TimestampedValue(value, timestamp):把普通值包装成"带时间戳的元素"。其实现同样位于 sdks/python/apache_beam/transforms/window.py 的TimestampedValue类——窗口函数读取的正是这个 timestamp。这里用datetime(...).timestamp()把带pytz.UTC时区的日期时间转成 UNIX 时间戳(秒)。beam.WindowInto(window.FixedWindows(24*60*60)):将窗口函数应用到 PCollection。24*60*60 = 86400秒即 1 天,这是 task.md 提示的核心答案。beam.combiners.Count.PerElement():按元素(key)计数。因为窗口划分后每个窗口形成独立的"视图",该 Transform 会隐式地按 key + window 聚合,所以同一 key("event")在不同窗口内分别计数。beam.LogElements(with_window=True):打印元素时附带所属窗口信息,便于直接观察每个计数落在哪个窗口区间。
六、预期输出与测试验证
隐藏测试 tests/test_task.py 中校验了四个窗口的计数结果:
answers = [ "('event', 4), window(start=2020-03-01T00:00:00Z, end=2020-03-02T00:00:00Z)", "('event', 2), window(start=2020-03-05T00:00:00Z, end=2020-03-06T00:00:00Z)", "('event', 3), window(start=2020-03-08T00:00:00Z, end=2020-03-09T00:00:00Z)", "('event', 1), window(start=2020-03-10T00:00:00Z, end=2020-03-11T00:00:00Z)" ]对照输入数据可验证窗口语义:
- 2020-03-01 有4个事件 → 窗口
[03-01, 03-02)计数 4; - 2020-03-05 有2个事件 → 窗口
[03-05, 03-06)计数 2; - 2020-03-08 有3个事件 → 窗口
[03-08, 03-09)计数 3; - 2020-03-10 有1个事件 → 窗口
[03-10, 03-11)计数 1。
注意窗口的边界是左闭右开:例如 03-01 00:00:00 属于[03-01, 03-02),而 03-02 00:00:00 则属于下一个窗口。测试通过test_not_empty(输出非空)与test_output(包含上述四行答案)两项断言来完成 Kata 验收。
七、运行方式与课程背景
本项目 Kata 采用 PyCharm Education(或安装 EduTools 插件的 PyCharm)运行,具体初始化步骤见 learning/katas/python/README.md:选择learning/katas/python目录创建项目、配置解释器(如 virtualenv)、在 Course 视图中按课程顺序完成练习。若要在命令行直接运行,也可使用 Python SDK 执行task.py(需已安装apache-beam与pytz),输出会直接打印各窗口的计数与窗口区间。
本任务在课程结构中隶属于 Streaming/Windows 课程目录(content: - Fixed Windows),是窗口(Windowing)主题下的第一个练习。它属于整个 Python Katas 课程体系中"Streaming → Windows → Fixed Windows"这条学习路径,后续可继续学习 Sliding Windows、Per-Session Windows 以及 Triggers 等进阶内容。
八、小结:从 Kata 到生产实践的要点
通过本 Kata,你可以提炼出三条可迁移到真实流水线的经验:
- 聚合前先想清楚窗口:任何
GroupByKey/Combine类操作都按 key + window 生效,窗口定义直接决定聚合结果的分组粒度; - 固定窗口用
FixedWindows(size, offset=0):size为秒数且必须为正,offset可在[0, size)内调整窗口对齐位置(例如按某时区对齐每天零点); - 无界流的有界化:固定窗口是把无界 PCollection 切成有限"批"的最直接手段,配合后续的触发器(Triggers)与累计模式(Accumulation Modes),可以进一步控制窗口结果的产出时机——这正是 Streaming 课程后续练习(如 Early Triggers、Window Accumulation Modes)要深入的主题。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
Apache Beam Go SDK 实战:使用固定时间窗口(Fixed Time Window)处理带时间戳的 PCollection
Apache Beam Go SDK 实战:使用固定时间窗口(Fixed Time Window)处理带时间戳的 PCollection 固定时间窗口(Fixe
大数据批处理流处理数据工程Apache Beam 固定时间窗口(Fixed Time Window)实战:用 1 天窗口统计事件数量
Apache Beam 固定时间窗口(Fixed Time Window)实战:用 1 天窗口统计事件数量 窗口(Windowing)是 Apache Beam
大数据批处理流处理数据工程Apache Beam Katas 实战:用 Kotlin 实现 Fixed Time Window 固定时间窗口按天聚合计数
Apache Beam Katas 实战:用 Kotlin 实现 Fixed Time Window 固定时间窗口按天聚合计数 本指南以 Apache Beam
大数据批处理流处理数据工程
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考