ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Apache Beam 复合触发器(Composite Trigger)实战指南:组合事件时间、处理时间与数据驱动触发策略

Apache Beam 复合触发器(Composite Trigger)实战指南:组合事件时间、处理时间与数据驱动触发策略 大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载Apache Beam 的复合触发器Composite Trigger允许把多种触发器组合使用当组合中的任一或全部子触发器满足条件时触发窗口发射从而构造出远超单一触发器的复杂触发策略。本文基于 Apache Beam 仓库中 Tour of Beam 的 Triggers 教学模块composite-trigger/description.md展开结合 Java / Python / Go 三种 SDK 的可运行示例与底层实现源码帮助你掌握复合触发器的定义方式、AfterAll / AfterFirst 等组合语义以及在实际管道中如何搭配窗口、迟到数据与累积模式使用。什么是复合触发器在 Apache Beam 中当数据被收集并按窗口分组时Beam 使用**触发器Trigger**来决定何时将每个窗口的聚合结果称为 pane发射出去。默认情况下Beam 在估计所有数据都已到达watermark 越过窗口结束点时输出一次聚合结果并丢弃该窗口后续到达的数据。复合触发器则是把多个触发器组合在一起使用的触发器。Apache Beam 的官方教程定义如下Acomposite triggerin Apache Beam allows you to specify multiple triggers to be used in combination. When any of the triggers fire, the composite trigger will fire. This allows you to combine different types of triggers to create more complex triggering strategies.翻译过来即复合触发器允许指定多个触发器组合使用组合中的触发器任一满足条件时复合触发器就会发射不同组合算子语义略有差异详见下文 AfterAll / AfterFirst 的精确区分从而把事件时间触发器、处理时间触发器、数据驱动触发器等不同类型组合成更复杂的触发策略。复合触发器非常适合以下场景既想在窗口结束前尽快看到部分结果early firing又不想放弃窗口结束后的完整结果既要处理时间兜底防止数据迟迟不来又要数据数量兜底达到 N 条就发射需要把“窗口结束 迟到数据 周期刷新”等多重条件叠加成一个整体策略。在 Tour of Beam 的 Triggers 模块中复合触发器被列为 ADVANCED高级难度的独立课程单元见 unit-info.yamlcomplexity: ADVANCED覆盖 Java、Python、Go 三种 SDK。三种 SDK 下的复合触发器定义官方教程 description.md 给出了 Java、Python、Go 三种 SDK 的复合触发器定义方式全部围绕AfterAll所有子触发器都就绪时才发射展开。JavaAfterAll.of(List)WindowString window Window.into(FixedWindows.of(Duration.standardMinutes(5))); Trigger processingTimeTrigger AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(1)); Trigger dataDrivenTrigger AfterPane.elementCountAtLeast(2); PCollectionString windowed input.apply(window.triggering(AfterAll.of(Arrays.asList(processingTimeTrigger,dataDrivenTrigger))).withAllowedLateness(Duration.ZERO).accumulatingFiredPanes());AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(1))处理时间触发器从 pane 内第一个元素到达起延时 1 分钟发射AfterPane.elementCountAtLeast(2)数据驱动触发器pane 内累积满 2 个元素即发射AfterAll.of(...)只有当上述两个子触发器都就绪时复合触发器才真正发射.withAllowedLateness(Duration.ZERO)不允许迟到数据.accumulatingFiredPanes()累积模式accumulating fired panes后续发射的 pane 会包含之前已发射过的元素。仓库中的完整可运行版本见 Task.java它把窗口内结果通过自定义LogOutputDoFn 输出到日志。Pythontrigger.AfterAll(...)processing_time_trigger trigger.AfterProcessingTime(60) event_time_trigger trigger.AfterWatermark(earlytrigger.AfterCount(100), latetrigger.AfterCount(200)) composite_trigger trigger.AfterAll(processing_time_trigger,event_time_trigger)AfterProcessingTime(60)处理时间触发从窗口开始或第一个元素起 60 秒后发射AfterWatermark(earlyAfterCount(100), lateAfterCount(200))事件时间触发器watermark 越过窗口结束前每累积 100 条元素提前发射一次watermark 之后迟到阶段每累积 200 条发射一次AfterAll(...)两者都就绪时复合触发。Python 的完整可运行示例见 task.py它把复合触发器挂到beam.WindowInto(FixedWindows(2), triggercomposite_trigger, accumulation_modetrigger.AccumulationMode.DISCARDING)上并使用 DISCARDING丢弃累积模式。Gotrigger.AfterAll([]trigger.Trigger{...})trigger : trigger.AfterAll([]trigger.Trigger{trigger.AfterEndOfWindow(). EarlyFiring(trigger.AfterProcessingTime(). PlusDelay(60 * time.Second)). LateFiring(trigger.Repeat(trigger.AfterCount(1))),trigger.AfterCount(2)})AfterEndOfWindow()窗口结束时发射 on-time pane.EarlyFiring(trigger.AfterProcessingTime().PlusDelay(60*time.Second))窗口结束前处理时间到达 60 秒时提前发射.LateFiring(trigger.Repeat(trigger.AfterCount(1)))窗口结束后每来 1 条迟到数据就发射一次trigger.AfterCount(2)数据驱动子触发器每累积 2 条发射整个AfterAll要求上述子触发器全部就绪才发射。Go 的完整可运行版本见 main.go它使用beam.WindowInto(s, window.NewFixedWindows(60*time.Second), input, beam.Trigger(trigger), beam.PanesDiscard())将复合触发器应用到 60 秒固定窗口上。Playground 练习AfterFirst 组合官方教程还给出了一个 Playground 练习场景在复合触发器中可以保证第一个触发器触发之后第二个触发器接着被触发即 AfterFirst 语义。JavaAfterFirst.of(...)Trigger trigger AfterFirst.of( AfterCount.of(100), AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(5)) );AfterCount.of(100)数据驱动子触发器累积 100 条发射AfterProcessingTime.pastFirstElementInPane().plusDelayOf(Duration.standardMinutes(5))处理时间子触发器pane 首个元素到达后 5 分钟发射AfterFirst.of(...)任一子触发器先满足即发射——数据多到 100 条立刻出结果或者即使数据迟迟凑不满5 分钟也兜底发射一次。PythonAfterFirst.of(...)input | WindowIntoFixedWindows beam.WindowInto(beam.window.FixedWindows(10)) | CountAndProcessTimeTrigger beam.Trigger( AfterFirst.of( AfterCount(100), AfterProcessingTime(5*60) # 5 minutes ) )固定窗口 10 秒AfterFirst组合数据量100 条与处理时间5 分钟两个兜底条件。精确理解组合语义AfterAll vs AfterFirst vs AfterEachBeam 的触发器教学模块 concept/description.md 系统列举了 Beam 内置的触发器家族其中与复合触发器直接相关的三个组合器语义必须区分清楚组合器语义子触发器定义方式典型用途AfterAll所有子触发器都 ready 时才发射of(ListTrigger)/ 构造函数传参要求多条件同时满足AfterEach子触发器按顺序逐个执行inOrder(ListTrigger)按阶段依次触发AfterFirst至少一个子触发器先满足即发射of(...)多个触发条件互为兜底需要特别指出的是官方教程 description.md 中的表述“When any of the triggers fire, the composite trigger will fire”描述的是 AfterFirst / AfterAny 一类语义而AfterAll在 Python SDK 源码 trigger.py 中明确实现为class AfterAll(_ParallelTriggerFn): Fires when all subtriggers have fired. Also finishes when all subtriggers have finished. combine_op all其姊妹类AfterAnycombine_op any才对应“任一子触发器触发即触发”。因此在选择组合器时务必按需求区分需要“全部满足”用 AfterAll需要“任一满足 / 互为兜底”用 AfterFirst / AfterAny。源码级实现印证GoAfterAll 的参数约束Go SDK 的 trigger.go 中AfterAll的构造函数对子触发器数量有强制约束// AfterAll returns a new AfterAll trigger with subtriggers set to the passed argument. func AfterAll(triggers []Trigger) *AfterAllTrigger { if len(triggers) 1 { panic(fmt.Sprintf(number of subtriggers to trigger.AfterAll() should be greater than 1, got: %v, len(triggers))) } return AfterAllTrigger{subtriggers: triggers} }即传给AfterAll的子触发器数量必须大于 1否则直接 panic从实现层面保证了“组合”二字名副其实。其余组合器的定义也在同一文件中AfterEach、AfterEndOfWindow、AfterProcessingTime、AfterCount等均位于该 trigger 包内。PythonAfterEach 的顺序状态Python 的AfterEach在 trigger.py 中通过INDEX_TAG状态记录当前执行到第几个子触发器on_element只把元素喂给当前索引对应的子触发器实现严格的“按顺序逐个执行”窗口合并时则取索引最靠后的状态——这印证了 concept 文档中“子触发器按 inOrder 顺序一个接一个执行”的描述。与累积模式Accumulation Mode配合在使用触发器时还必须同时设置窗口的累积模式见 concept/description.md。触发器每次发射都会把窗口当前内容作为 pane 输出由于触发器可能多次发射累积模式决定了系统是保留之前已发射的 pane 内容还是丢弃它们。累积模式accumulatingaccumulatingFiredPanes()Java/AccumulationMode.ACCUMULATINGPython/beam.PanesAccumulate()Go。每次发射都包含之前发射过的元素First trigger firing: [5, 8, 3] Second trigger firing: [5, 8, 3, 15, 19, 23] Third trigger firing: [5, 8, 3, 15, 19, 23, 9, 13, 10]丢弃模式discardingdiscardingFiredPanes()Java/AccumulationMode.DISCARDINGPython/beam.PanesDiscard()Go。每次发射只包含两次发射之间新到的元素First trigger firing: [5, 8, 3] Second trigger firing: [15, 19, 23] Third trigger firing: [9, 13, 10]本模块的 Java 示例使用累积模式accumulatingFiredPanes()Python 与 Go 示例使用丢弃模式AccumulationMode.DISCARDING/beam.PanesDiscard()正好覆盖两种模式的实际写法。迟到的数据怎么办如果希望管道处理 watermark 越过窗口结束点之后才到达的数据需要在设置窗口配置时指定允许迟到时间allowed lateness让触发器有机会对迟到数据做出反应一旦设置了 allowed lateness默认触发器会在迟到数据到达时立即发射新结果。三种 SDK 的写法出自 concept/description.mdPCollectionString input ...; input.apply(Window.Stringinto(FixedWindows.of(1, TimeUnit.MINUTES)) .triggering(AfterProcessingTime.pastFirstElementInPane() .plusDelayOf(Duration.standardMinutes(1))) .withAllowedLateness(Duration.standardMinutes(30));input [Initial PCollection] input | beam.WindowInto( FixedWindows(60), triggerAfterProcessingTime(60), allowed_lateness1800) # 30 minutes | ...allowedToBeLateItems : beam.WindowInto(s, window.NewFixedWindows(1*time.Minute), pcollection, beam.Trigger(trigger.AfterProcessingTime(). PlusDelay(1*time.Minute)), beam.AllowedLateness(30*time.Minute), )复合触发器常与AfterWatermark的late参数或LateFiring组合使用把迟到数据处理纳入整体触发策略如 Go 示例中的LateFiring(trigger.Repeat(trigger.AfterCount(1)))。小结复合触发器是 Apache Beam 触发器体系中能力最强的部分通过AfterAll全部就绪、AfterFirst/AfterAny任一就绪、AfterEach顺序执行三种组合方式可以把事件时间、处理时间、数据驱动三大类触发器自由拼接成贴合业务语义的复杂策略。实际使用时注意三点一是根据“全部满足还是任一满足”正确选择组合器二是为窗口配置合适的累积模式累积 / 丢弃三是用 allowed lateness 为迟到数据留出处理窗口。你可以直接运行仓库中 JavaTask.java、Pythontask.py、Gomain.go三个可执行示例或在 Tour of Beam 的 Playground 中完成 AfterFirst 练习来加深理解。赞分享大数据批处理流处理数据工程【免费下载链接】beamApache Beam is a unified programming model for Batch and Streaming data processing.项目地址https://gitcode.com/gh_mirrors/beam4/beam点击查看免费下载相关推荐Apache Beam 触发器实战用组合触发器按元素数与时间窗输出出租车订单数据Apache Beam 触发器实战用组合触发器按元素数与时间窗输出出租车订单数据 导读 本文围绕 Apache Beam 官方 Tour of Beam 学习大数据批处理流处理数据工程Apache Beam Java Katas 实战使用 AfterWatermark 事件时间触发器实现固定窗口事件计数Apache Beam Java Katas 实战使用 AfterWatermark 事件时间触发器实现固定窗口事件计数 导读 本篇文章围绕 Apache B大数据批处理流处理数据工程英雄联盟对局先知选人阶段胜负预测终极指南英雄联盟对局先知选人阶段胜负预测终极指南 在英雄联盟的激烈对局中你是否经常遇到这样的困惑选人阶段看不清队友实力进入游戏才发现队友是牛马对手是通天上一篇终极指南如何利用 awesome-static-analysis 提升代码质量与供应链安全 下一篇EventBus: 简洁高效的Java事件总线框架创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表