- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
本文以 Apache Beam Tour of Beam 课程中 Windowing(窗口化)模块的 motivating challenge 练习为核心,带领读者使用纽约出租车行程 CSV 数据,构建一个"计算过去 10 分钟最高打车价格、且每分钟刷新一次"的流水线。读完本文,你将理解滑动时间窗口(Sliding Time Window)的窗口长度与滑动周期的含义与搭配方式,并能在 Java、Python、Go 三种 SDK 中独立完成该挑战及其完整实现。
练习背景:一次基于真实数据的窗口化热身
本练习出自 windowing/motivating-challenge,属于 Tour of Beam 课程中 Windowing 模块的进阶实战环节。模块本身按 module-info.yaml 编排了从windowing-concept(窗口概念)到adding-timestamp(添加时间戳)、global-window(全局窗口)、fixed-time-window(固定时间窗口)、sliding-time-window(滑动时间窗口)、session-window(会话窗口)的完整学习路径,而 motivating challenge 正是要求你综合运用前面所学,独立完成一个带时间窗口的聚合任务。
练习使用的数据集是纽约出租车行程(NYC taxi trips)CSV 文件,包含若干行程字段,其中与本练习直接相关的是单次行程的价格(cost)。数据集的部分字段与示例记录如下表所示:
| cost | passenger_count | ... |
|---|---|---|
| 5.8 | 1 | ... |
| 4.6 | 2 | ... |
| 24 | 1 | ... |
练习要求如下:
编写一个流水线,返回过去 10 分钟内出租车行程的最高价格,并且计算结果需要每分钟更新一次。
这一需求对应的核心窗口语义是:窗口长度为 10 分钟,滑动周期(period)为 1 分钟。也就是说,每个窗口覆盖过去 10 分钟的数据,但每隔 1 分钟就产生一个新的窗口、输出一次"过去 10 分钟最大值"。这正是滑动时间窗口最适合解决的场景:窗口之间相互重叠,绝大多数元素会同时属于多个窗口,从而可以在保持窗口"足够长"的同时,以更细的粒度(周期)刷新聚合结果——这与滑动时间窗口常用于"滚动平均"(如过去 60 秒数据每 30 秒更新一次)的典型用途一脉相承,参见 sliding-time-window/description.md。
滑动时间窗口核心概念:窗口长度 × 滑动周期
在动手写代码之前,先厘清两个决定窗口行为的参数:
- 窗口长度(window duration / size):每个窗口覆盖的时间跨度,即一次聚合统计多少时间的数据。
- 滑动周期(period / frequency):新窗口开始的时间间隔,即聚合结果多久刷新一次。
以官方示例为例:每个窗口捕获 60 秒的数据,但每 30 秒就启动一个新窗口,那么窗口长度是 60 秒、周期是 30 秒。由于窗口相互重叠,数据集中大多数元素会同时属于多个窗口;这种窗口化方式非常适合"运行中的均值"之类的计算——你可以在 60 秒窗口上求均值,并每 30 秒刷新一次输出。
映射到本练习就是:窗口长度 10 分钟、周期 1 分钟。注意挑战中"过去 10 分钟"与"每分钟更新"是两个独立维度,分别对应SlidingWindows的两个参数,缺一不可。
挑战代码骨架:先把数据读进来、把价格提出来
挑战模板已经替你完成了数据接入与字段解析部分,三种 SDK 的挑战代码分别位于:
- Java:java-challenge/Task.java
- Python:python-challenge/task.py
- Go:go-challenge/main.go
它们的共同流程是:从gs://apache-beam-samples/nyc_taxi/misc/sample1000.csv读取 CSV 文本行 → 按逗号切分 → 取出第 17 个字段(索引 16,即 total amount 车费字段)解析为浮点数 → 得到PCollection<Double>类型的车费集合rideTotalAmounts/cost。
以 Java 为例,解析逻辑如下(节选自 java-challenge/Task.java):
static class ExtractTaxiRideCostFn extends DoFn<String, Double> { @ProcessElement public void processElement(ProcessContext c) { String[] items = c.element().split(","); Double totalAmount = tryParseTaxiRideCost(items); c.output(totalAmount); } } private static Double tryParseTaxiRideCost(String[] inputItems) { try { return Double.parseDouble(tryParseString(inputItems, 16)); } catch (NumberFormatException | NullPointerException e) { return 0.0; // 字段缺失或非法时兜底为 0.0 } }Python 版用生成器实现同类逻辑,Go 版则在ExtractCostFromFile中以ParDo内联函数完成strconv.ParseFloat(taxi[16], 64)解析。这些骨架代码特意把"窗口化 + 全局聚合"留空,正是你要补齐的部分。
参考答案:三种 SDK 的滑动窗口实现
根据练习提示 hint1.md,解题路径清晰且统一:先给 PCollection 施加滑动窗口,再做全局求最大值聚合。以下是三种 SDK 的完整解答与逐段讲解。
Java:SlidingWindows.of(...).every(...) + Combine.globally
完整解答见 java-solution/Task.java:
PCollection<Double> slidingWindowedItems = rideTotalAmounts.apply( Window.<Double>into(SlidingWindows.of(Duration.standardSeconds(10)).every(Duration.standardSeconds(5)))) .apply(Combine.globally((SerializableFunction<Iterable<Double>, Double>) ride -> { Iterator<Double> iterator = ride.iterator(); double firstValue = iterator.hasNext() ? iterator.next() : -Double.MAX_VALUE; for (double i : ride) { if (firstValue < i) { firstValue = i; } } return firstValue; }).withoutDefaults()); slidingWindowedItems.apply("Log words", ParDo.of(new LogOutput<>()));要点拆解:
SlidingWindows.of(Duration.standardSeconds(10))指定窗口长度为 10 秒;.every(Duration.standardSeconds(5))指定每 5 秒启动一个新窗口。官方示例(sliding-time-window/description.md)使用SlidingWindows.of(Duration.standardSeconds(30)).every(Duration.standardSeconds(5))表达"30 秒窗口、每 5 秒刷新一次",结构与本解答完全一致;挑战中只需将参数调整为 10 分钟与 1 分钟即可。- 在窗口化之后调用
Combine.globally(...)进行全局聚合。这里没有使用现成的Max.doubles(),而是手写了一个SerializableFunction:遍历当前窗口内全部车费,维护最大值并返回;窗口为空时返回-Double.MAX_VALUE作为下界兜底。 .withoutDefaults()至关重要:Combine.globally默认会对空输入 PCollection 输出默认值(如 0.0),这会在每个滑动窗口都可能为空时产生误导性结果;withoutDefaults()让空窗口不产生输出,语义更准确。
Python:SlidingWindows(size, period) + CombineGlobally(max)
完整解答见 python-solution/task.py:
with beam.Pipeline() as p: input = (p | 'Log lines' >> beam.io.ReadFromText('gs://apache-beam-samples/nyc_taxi/misc/sample1000.csv') | beam.ParDo(ExtractTaxiRideCostFn())) (input | 'window' >> beam.WindowInto(window.SlidingWindows(10, 5)) | 'Sum above cost' >> beam.CombineGlobally(max).without_defaults() | 'Log above cost' >> Output())要点拆解:
window.SlidingWindows(10, 5)的两个位置参数分别是窗口长度与滑动周期(均为秒);结合挑战需求,应写成window.SlidingWindows(600, 60)。注意 Python API 的参数顺序与 Java 的.of(size).every(period)相反,Go 的NewSlidingWindows(period, size)又与二者都不同(见下文),跨语言对照时务必留意。beam.CombineGlobally(max)直接复用 Python 内置max作为组合函数,简单直接;.without_defaults()与 Java 语义相同,避免空窗口输出默认值。
Go:beam.WindowInto + window.NewSlidingWindows(period, size) + stats.Max
完整解答见 go-solution/main.go:
slidingWindowedItems := beam.WindowInto(s, window.NewSlidingWindows(5*time.Second, 10*time.Second), cost) max := getMax(s, slidingWindowedItems) debug.Printf(s, "Above pCollection output", max) func getMax(s beam.Scope, input beam.PCollection) beam.PCollection { return stats.Max(s, input) }要点拆解:
window.NewSlidingWindows(5*time.Second, 10*time.Second)的第一参数是滑动周期、第二参数是窗口长度,与 Python 的SlidingWindows(size, period)恰好相反。挑战若改为"10 分钟窗口、每分钟刷新",应写作window.NewSlidingWindows(1*time.Minute, 10*time.Minute)。- 求最大值直接复用 Beam Go SDK 提供的统计变换
stats.Max(s, input),无需手写遍历逻辑——这是 Go 解法比 Java 解法更简洁的原因之一。 - Go 代码中需要显式导入
github.com/apache/beam/sdks/v2/go/pkg/beam/core/graph/window与.../beam/transforms/stats才能使用上述 API。
从答案反推原理:为什么窗口与聚合必须分开写
细看三种实现,会发现一个共同设计:窗口化(WindowInto)与聚合(Combine.globally / stats.Max)是两个串行的独立变换。
WindowInto/Window.into负责把无界数据流按事件时间切分为若干个带时间边界的窗口,将每个元素分发到它所属的一个或多个窗口中(滑动窗口下通常多于一个)。- 之后的
Combine.globally/stats.Max在窗口化的 PCollection上执行,Beam 的聚合语义天然按窗口进行——即每个窗口内部独立完成"求最大值",窗口之间互不干扰。
这正是练习要求逐字对应的两种能力:"过去 10 分钟"由窗口长度决定,"每分钟更新"由滑动周期决定,而"最高价格"由窗口内的全局最大值聚合决定。三者分别被建模为窗口参数与聚合变换,而不是混在同一个 API 调用里。
直接跑起来:挑战模板与解答的目录对应关系
本练习的挑战模板与官方解答按 SDK 分目录存放,方便对照学习:
| SDK | 挑战模板 | 官方解答 |
|---|---|---|
| Java | java-challenge/Task.java | java-solution/Task.java |
| Python | python-challenge/task.py | python-solution/task.py |
| Go | go-challenge/main.go | go-solution/main.go |
学习建议:先只阅读挑战模板,亲手补上窗口与聚合两段代码,再对照解答验证。两版代码的差异集中在窗口化与聚合部分,字段解析等骨架代码完全一致,方便逐行比对。
延伸:把本练习嵌入窗口化学习路径
本挑战是整个 Windowing 模块的收口练习,建议按模块编排顺序(module-info.yaml)依次掌握前置单元:
- windowing-concept:理解为什么无界数据需要按窗口切分。
- adding-timestamp:窗口依赖事件时间,先学会为数据附加时间戳。
- global-window:默认的全局窗口及其聚合限制。
- fixed-time-window:固定时间窗口,窗口互不重叠。
- sliding-time-window:本挑战直接依赖的滑动窗口,窗口相互重叠,是"过去 N 分钟、每 M 分钟刷新"类需求的标配。
- session-window:基于数据间隔切分的会话窗口。
如果对窗口化后的聚合输出感兴趣,可继续学习 triggers 模块——它决定了窗口结果在何时、以何种方式被触发输出,是与本练习"每分钟更新一次"目标最接近的进阶话题。
- 大数据
- 批处理
- 流处理
- 数据工程
【免费下载链接】beam
Apache Beam is a unified programming model for Batch and Streaming data processing.
相关推荐
对比A3C与PPO:Super-mario-bros-PPO-pytorch如何实现从19关到31关的突破?
对比A3C与PPO:Super mario bros PPO pytorch如何实现从19关到31关的突破? Super mario bros PPO pyto
人工智能强化学习深度学习用 Apache Flink 会话窗口分析纽约出租车行程:Data Engineering Zoomcamp PyFlink 流处理作业实战
用 Apache Flink 会话窗口分析纽约出租车行程:Data Engineering Zoomcamp PyFlink 流处理作业实战 本篇技术指南围绕
教程数据工程纽约市出租车与网约车数据分析项目完整指南
纽约市出租车与网约车数据分析项目是一个功能强大的开源工具集,专门用于处理和分析纽约市数十亿次出租车及网约车行程记录。该项目提供了从数据下载到深度分析的完整工作流
数据分析数据工程大数据
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考