☰
Apache Beam 滑动时间窗口实战:用纽约出租车数据完成“过去 10 分钟最高车费“流式练习
2026/10/7 2:40:29 网站建设 项目流程
  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

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

本文以 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)。数据集的部分字段与示例记录如下表所示:

costpassenger_count...
5.81...
4.62...
241...

练习要求如下:

编写一个流水线,返回过去 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挑战模板官方解答
Javajava-challenge/Task.javajava-solution/Task.java
Pythonpython-challenge/task.pypython-solution/task.py
Gogo-challenge/main.gogo-solution/main.go

学习建议:先只阅读挑战模板,亲手补上窗口与聚合两段代码,再对照解答验证。两版代码的差异集中在窗口化与聚合部分,字段解析等骨架代码完全一致,方便逐行比对。

延伸:把本练习嵌入窗口化学习路径

本挑战是整个 Windowing 模块的收口练习,建议按模块编排顺序(module-info.yaml)依次掌握前置单元:

  1. windowing-concept:理解为什么无界数据需要按窗口切分。
  2. adding-timestamp:窗口依赖事件时间,先学会为数据附加时间戳。
  3. global-window:默认的全局窗口及其聚合限制。
  4. fixed-time-window:固定时间窗口,窗口互不重叠。
  5. sliding-time-window:本挑战直接依赖的滑动窗口,窗口相互重叠,是"过去 N 分钟、每 M 分钟刷新"类需求的标配。
  6. session-window:基于数据间隔切分的会话窗口。

如果对窗口化后的聚合输出感兴趣,可继续学习 triggers 模块——它决定了窗口结果在何时、以何种方式被触发输出,是与本练习"每分钟更新一次"目标最接近的进阶话题。

  • 大数据
  • 批处理
  • 流处理
  • 数据工程

【免费下载链接】beam

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

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

相关推荐

上一篇:一键 PaaS 部署实战:用 CloudBase、Vercel、Netlify 与 Zeabur 把你的 Web 应用发布上线
下一篇:Home Assistant 中 Matter 热水器 Boost 动作(matter.water_heater_boost)完整指南

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询