Apache Beam Python Kata 实战:用 AfterWatermark Early Triggers 在固定窗口中实现即时提前输出
大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载导读本文围绕 Apache Beam 官方 Katas 学习项目中的 Early Triggers提前触发练习展开以 task.md 为骨架结合 task.py、generate_event.py 及 Python SDK 的 trigger.py 源码完整演示如何在 1 天固定窗口中借助AfterWatermark(earlyAfterCount(1))让每个新元素被处理后立即输出一次计数结果。读完本文你将掌握 Beam 触发器的分类、AfterWatermark的 early/late 参数语义、累积模式与 allowed lateness 的配合方式并能独立复现该 Kata 的完整可运行代码与验证流程。一、为什么需要 Early Triggers默认触发器的延迟困境在 Beam 中窗口Window负责将数据按时间切分而触发器Trigger决定聚合结果在什么时候被输出为一个 pane面板/批次。默认情况下Beam 会在估算窗口内所有数据都已到达时才输出聚合结果即等到 watermark水位线越过窗口末尾并且此后到达的迟到数据默认会被丢弃。这套默认策略对批处理很合适但对低延迟流式场景却是个问题假设你使用 1 天固定窗口统计事件数如果必须等满 24 小时、watermark 越过窗口终点才能看到第一个结果那么实时监控滚动告警增量报表等需求就完全无法满足。Early Triggers 正是为此而生。原文档对它的定义非常直接Triggers allow Beam to emit early results, before all the data in a given window has arrived. For example, emitting after a certain amount of time elapses, or after a certain number of elements arrives.即触发器允许 Beam 在窗口内数据尚未全部到达时就提前产出结果——例如每隔一段时间发射一次或每到达 N 个元素就发射一次。这正是本文 Kata 要解决的问题。二、Beam 触发器体系速览在深入实现之前先厘清 Beam 提供的触发器分类对应 Event Time Triggers 练习文档 中的说明触发器类型触发依据典型内置实现事件时间触发器Event time triggers元素自带的事件时间戳AfterWatermark处理时间触发器Processing time triggers数据被处理时的系统时间AfterProcessingTime数据驱动触发器Data-driven triggers已到达窗口内的元素数量/数据本身AfterCount复合触发器Composite triggers上述多种条件的任意组合AfterWatermark(early..., late...)、AfterAll、AfterAny等其中AfterWatermark是事件时间触发器的代表。在 trigger.py 源码中它的 docstring 和构造函数清楚地说明了其语义class AfterWatermark(TriggerFn): Fire exactly once when the watermark passes the end of the window. Args: early: if not None, a speculative trigger to repeatedly evaluate before the watermark passes the end of the window late: if not None, a speculative trigger to repeatedly evaluate after the watermark passes the end of the window def __init__(self, earlyNone, lateNone): self.early self._wrap_if_not_repeatedly(early) ...不带任何参数时AfterWatermark()只在 watermark 越过窗口末尾时发射一次即默认行为early参数指定 watermark 到达之前的投机性speculative提前触发条件late参数指定 watermark 越过窗口之后、迟到数据到来时的迟发触发条件。Watermark 是 Beam 中衡量某个时刻输入完整程度的全局进度指标AfterWatermark正是基于这一指标来决定发射时机。三、Kata 目标拆解原文档给出的练习目标如下Given that sample events generated with one second granularity and a fixed window of 1-day duration, please implement an early trigger that emits the number of events count immediately after new element is processed.翻译成工程语言输入数据由工具以1 秒粒度生成每个事件时间戳相差 1 秒使用1 天时长的FixedWindows固定窗口实现一个提前触发每处理一个新的元素就立即输出一次当前窗口内的事件计数配合allowed_lateness 0与DISCARDING丢弃式累积模式。原文档给出了四个关键提示即实现所需的四个技术点AfterWatermark的early参数、1 天FixedWindows、allowed lateness 置 0 DISCARDING模式、CombineGloballyCountCombineFn计数。下面逐一落实。四、数据准备TestStream 生成 1 秒粒度事件Kata 使用 generate_event.py 中的GenerateEventPTransform 构造流式输入其底层是 Beam 测试常用的TestStreamclass GenerateEvent(beam.PTransform): staticmethod def sample_data(): return GenerateEvent() def expand(self, input): return (input | TestStream() .add_elements(elements[event], event_timestampdatetime(2021, 3, 1, 0, 0, 1, 0, tzinfopytz.UTC).timestamp()) .add_elements(elements[event], event_timestampdatetime(2021, 3, 1, 0, 0, 2, 0, tzinfopytz.UTC).timestamp()) # ... 每个事件时间戳递增 1 秒 ... .advance_watermark_to(...) .advance_watermark_to_infinity())要点事件时间戳从2021-03-01 00:00:01 UTC开始逐秒递增共 20 个event元素每隔几秒调用advance_watermark_to(...)手动推进 watermark模拟真实流式中水位线随进度前进的现象最后advance_watermark_to_infinity()将 watermark 推到无穷强制窗口完成最终发射时间戳全部使用pytz.UTC时区确保事件时间语义清晰无歧义。TestStream的价值在于它可以精确控制何时注入元素、何时推进 watermark从而在无真实消息系统的情况下验证触发器的发射时机这也是本 Kata 能在本地跑通的原因。五、核心实现AfterWatermark AfterCount 提前触发完整的解答位于 task.py关键代码如下import apache_beam as beam from apache_beam.options.pipeline_options import PipelineOptions from apache_beam.options.pipeline_options import StandardOptions from generate_event import GenerateEvent from apache_beam.transforms.window import FixedWindows from apache_beam.transforms.trigger import AfterWatermark from apache_beam.transforms.trigger import AfterCount from apache_beam.transforms.trigger import AccumulationMode from apache_beam.utils.timestamp import Duration from apache_beam.transforms.util import LogElements class CountEventsWithEarlyTrigger(beam.PTransform): def expand(self, events): return (events | beam.WindowInto(FixedWindows(1 * 24 * 60 * 60), # 1 Day Window triggerAfterWatermark(earlyAfterCount(1)), accumulation_modeAccumulationMode.DISCARDING, allowed_latenessDuration(seconds0)) | beam.CombineGlobally(beam.combiners.CountCombineFn()).without_defaults()) options PipelineOptions() options.view_as(StandardOptions).streaming True # Required to get multiple trigger firing outputs with beam.Pipeline(optionsoptions) as p: (p | GenerateEvent.sample_data() | CountEventsWithEarlyTrigger() | LogElements(with_windowTrue))逐参数拆解1. 窗口FixedWindows(1 * 24 * 60 * 60)将数据流切分为 1 天86400 秒的固定窗口。本例所有事件都落在2021-03-01 00:00:00Z至2021-03-02 00:00:00Z这一个窗口内。2. 触发器AfterWatermark(earlyAfterCount(1))这是提前触发的灵魂。AfterWatermark负责在 watermark 越过窗口末尾时进行最终发射earlyAfterCount(1)则声明在 watermark 到达之前每累积到 1 个元素就发射一次。也就是说每处理一个新事件窗口就会立即输出一次当前累计计数——恰好满足emit the number of events count immediately after new element is processed的要求。从源码看AfterCount 的实现语义是Fire when there are at least count elements in this window pane且构造时校验 count 必须是正整数class AfterCount(TriggerFn): Fire when there are at least count elements in this window pane. def __init__(self, count): if not isinstance(count, numbers.Integral) or count 1: raise ValueError(count (%d) must be a positive integer. % count) self.count count因此AfterCount(1)表示每来 1 个元素就触发是最高频的提前触发节奏。实际生产中可改为AfterCount(100)每 100 条发射一次或AfterProcessingTime(60)每 60 秒发射一次以平衡延迟与开销。3. 累积模式AccumulationMode.DISCARDING由于一个窗口会被触发多次本例高达 20 余次必须明确每次发射后如何处理已累积的窗口状态。AccumulationMode 定义了两个取值class AccumulationMode(object): Controls what to do with data when a trigger fires multiple times. DISCARDING beam_runner_api_pb2.AccumulationMode.DISCARDING ACCUMULATING beam_runner_api_pb2.AccumulationMode.ACCUMULATINGDISCARDING丢弃式每次发射后清空窗口状态下一次发射从零开始计数。本 Kata 期望输出是每个 pane 都是1正是丢弃式的效果——每个事件单独成 pane计数互不累积ACCUMULATING累积式保留每次发射的结果并继续累加后一次输出包含前一次的内容这是本系列下一个练习 Window Accumulation Modes 的主题。4. 允许迟到allowed_latenessDuration(seconds0)allowed_lateness控制 watermark 越过窗口末尾后窗口状态继续保留多久以接收迟到数据。置 0 表示watermark 越过窗口即视为窗口完结不再等待任何迟到元素从而与DISCARDING模式共同保证输出结果确定、干净、可预期。5. 计数CombineGlobally(beam.combiners.CountCombineFn()).without_defaults()CombineGlobally是全局聚合变换CountCombineFn()计算元素个数.without_defaults()表示空输入时不输出默认的 0否则空窗口会额外输出一个无意义的 0。这里省略默认值后计数只在真实有数据时产生。6. 流式模式开关options.view_as(StandardOptions).streaming True是必不可少的一行只有声明为流式管道触发器才会被反复调用、产生多个 pane 输出若省略管道按批处理语义执行触发器多次发射的行为将无法体现。六、运行与验证预期输出逐 pane 对照运行task.py后LogElements(with_windowTrue)会把每个 pane 的数值连同其所属窗口打印出来。参考 tests/test_task.py 中的断言预期的核心输出为1, window(start2021-03-01T00:00:00Z, end2021-03-02T00:00:00Z) # 第 1 个事件到达立即输出 1, window(start2021-03-01T00:00:00Z, end2021-03-02T00:00:00Z) # 第 2 个事件到达立即输出 ...共 20 行每个事件对应一次输出 0, window(start2021-03-01T00:00:00Z, end2021-03-02T00:00:00Z) # watermark 越过窗口末尾后的最终发射逐 pane 解读前 20 行AfterCount(1)使每个新元素到达后窗口立即发射一次且因DISCARDING模式每次计数都从 0 重新累积故每次都输出1。这 20 行分别对应generate_event.py中时间戳从 00:00:01 到 00:00:20 的 20 个事件最后一行advance_watermark_to_infinity()让 watermark 越过窗口末尾AfterWatermark的 on-time 部分执行最终发射。此时窗口状态已被全部丢弃计数为 0输出0——这正是提前触发 丢弃式 0 迟到三者组合产生的独特输出形态所有 pane 都标注了同一窗口2021-03-01T00:00:00Z~2021-03-02T00:00:00Z证明 1 天窗口确实把所有事件收拢在同一个窗口内。可以自行运行验证需先安装 apache-beam Python SDKcd learning/katas/python/Streaming/Triggers/Early Triggers python task.pytest_task.py中通过test_is_not_empty检查输出非空并通过test_output逐条断言上述 21 行输出均出现在结果中是判断实现是否正确的标准答案。七、纵向对比默认触发、事件时间触发与提前触发将本练习与本系列另外两个练习对照能更清晰地理解触发器参数的作用场景触发器配置输出形态默认行为无触发器不指定 triggerwatermark 越过窗口后才输出一次最终聚合Event Time TriggersAfterWatermark()5 秒窗口watermark 越过 5 秒窗口末尾时输出一次本文 Early TriggersAfterWatermark(earlyAfterCount(1))1 天窗口每个新元素到达立即输出一次 watermark 越过后的最终发射Window Accumulation ModesAfterWatermark(earlyAfterCount(1))ACCUMULATING每次发射输出累计值1、2、3、…而非恒为 1这一对照恰好展示了 Beam 触发器的核心设计哲学窗口负责分桶触发器负责何时交货。二者正交组合可以派生出丰富的实时语义如每分钟输出滚动累计值、达到阈值告警、窗口结束时输出最终值等。八、实战要点总结AfterWatermark的early参数是把等满窗口变成边到边出的关键AfterCount(n)、AfterProcessingTime(d)等可作为其提前触发子句分别对应数据驱动与处理时间驱动多次触发必须搭配accumulation_modeDISCARDING得到增量/独立 paneACCUMULATING得到累计 pane二者输出语义截然不同allowed_lateness0保证窗口在 watermark 越过后立即终结结果确定、无迟到数据干扰流式管道必须显式设置streamingTrue否则触发器不会产生多次发射本 Kata 位于 learning/katas/python/Streaming/Triggers/ 目录同系列还有 Event Time Triggers 与 Window Accumulation Modes 两个练习三者递进覆盖了 Beam 流式窗口触发的主要知识点建议按顺序完成以形成完整认知。通过本文的完整实现与源码级解读你已经掌握了在 Apache Beam Python SDK 中使用 Early Triggers 让窗口聚合即时响应的完整套路选对触发器、配好累积模式、控制允许迟到即可在保持窗口语义的同时把结果延迟从窗口结束压缩到每条数据到达。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam Kotlin Kata 实战用 Early Triggers 让 1 天固定窗口立即产出结果Apache Beam Kotlin Kata 实战用 Early Triggers 让 1 天固定窗口立即产出结果 导读 在 Apache Beam 中默大数据批处理流处理数据工程Apache Beam Java Kata用 Early Triggers 在窗口结束前提前输出事件计数Apache Beam Java Kata用 Early Triggers 在窗口结束前提前输出事件计数 本篇文章基于 learning/katas/java大数据批处理流处理数据工程Apache Beam Python 事件时间触发器Event Time Triggers实战AfterWatermark 与固定窗口计数 Kata 解析Apache Beam Python 事件时间触发器Event Time Triggers实战AfterWatermark 与固定窗口计数 Kata 解析大数据批处理流处理数据工程上一篇RPG Maker MV插件宝库300插件让你的游戏开发效率翻倍下一篇VisualCppRedist AIO终极运行库修复工具3分钟解决DLL缺失问题创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
上一篇/下一篇内容由系统自动关联
返回资讯列表 →