支持:多时间维度流处理的水印与合并机制深度解析)
缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址https://gitcode.com/gh_mirrors/ha/hazelcast点击查看免费下载本文基于 Hazelcast 仓库中的设计文档 docs/design/sql/14-keyed-watermark-support.md深入讲解 Hazelcast Jet 引擎为支撑流到流 JOINstream-to-stream join而引入的多键水印keyed watermark机制包括Watermark类的键字段设计、ProcessorAPI 的扩展、不同键水印的独立合并语义以及按键 idle message 的处理规则。读完本文你将理解 Jet 中水印的完整生命周期生成、传播、合并、延迟、丢弃、三条水印约束规则以及该机制在滑动窗口、异步变换、流合并等处理器中的落地方式。背景为什么需要多个水印在流处理中单条流只有一个水印时间戳timestamp所有事件共享同一条水印序列。但当我们要实现流到流 JOIN例如订单流orders与投递流deliveries进行时间窗口内的关联时JOIN 两侧输入各自有独立的时间进度——左侧按order_time推进、右侧按delivery_time推进。如果仍使用单一水印就无法分别追踪两条流的进度导致 JOIN 处理器无法判断某条流的事件是否已经全部到齐。关于流到流 JOIN 为何需要多水印的完整动机可参见 docs/design/sql/15-stream-to-stream-join.md。本文聚焦的是Jet 引擎层面为支持多水印所需的改动不讨论 JOIN 本身的算法。本文的目标可以概括为三点定义按键水印keyed WM的合并coalescing行为定义不处理水印的处理器的默认行为由 Jet 引擎实现定义需要处理水印的处理器所需的改动。术语watermark 与 watermark instance在 Jet 中watermark一词被用于指代两种不同的东西一种广播流条目broadcast stream item它携带一个水印值一条流中所有水印的集合引入按键水印后指同一 key 的所有水印的集合。为消除歧义本文约定当强调携带值的流条目时使用术语watermark instance水印实例当说流中有两个水印时指的是由 key 区分的两套相互独立的水印集合。当前行为单键水印的语义水印的含义一个值为M的水印实例表示任何时间戳小于M的事件都可以视为迟到late并被忽略。这在流处理中极其重要——它让处理器能够推断未来还可能出现什么事件从而安全地清理状态、关闭窗口并输出结果。我们说一条流处于水印值N当且仅当最后收到的一个水印实例的值为N。由于每个新水印实例的值必须严格大于前一个流的水印值始终单调递增。水印合并WM coalescing所谓水印合并发生在两条流合并为一条的场合合并后的流使用所有被合并流中当前水印值的最小值。这是因为合并流无法保证任何一条输入后续不会再送来更早的事件只有取最小值才是安全保守的。原文给出的合并示例只展示水印实例事件载荷与合并无关in: wm(0, 10) in: wm(1, 12) out: wm(10) in: wm(0, 11) out: wm(11) in: wm(1, 13) ## 没有新输出因为最小值仍是 11 in: wm(0, 14) out: wm(13)其中in: wm(N, M)表示来自输入N、值为M的水印实例out: wm(M)表示转发到合并输出流的水印实例。这里的合并发生在Jet 引擎内部的流合并层面多输入汇聚到一条 DAG 边并非某个具体处理器。这一取最小值逻辑在源码中由 WatermarkCoalescer.java 实现其StandardImpl为每个输入队列维护queueWms[]数组在checkObservedWms()中遍历找出所有非 idle 输入队列的最小观测水印值min仅当min大于上次已发射值时才发射否则返回NO_NEW_WM。该合并器还针对 0 输入源处理器ZeroInputImpl与单输入直接转发SingleInputImpl做了专门优化。空闲消息IDLE_MESSAGE某条输入可能收到特殊的IDLE_MESSAGE表示该输入流应从合并中排除。收到IDLE_MESSAGE后该输入被排除出合并计算水印基于剩余非空闲输入的最小值输出一旦该输入再次送来任何新水印或事件空闲状态被清除重新参与合并。IDLE_MESSAGE在实现上是一个值为Long.MAX_VALUE的水印实例见 WatermarkCoalescer.java 中的IDLE_MESSAGE_TIME Long.MAX_VALUE与IDLE_MESSAGE new Watermark(IDLE_MESSAGE_TIME, (byte) 0)它并非真正的水印只是搭便车复用了水印广播机制。水印的生成任何处理器都可以基于任意逻辑生成水印它可以生成新的水印转发输入的水印丢弃输入的水印做其他任何事情。典型规则如下源source根据其观测到的最新事件之后的滞后量lag发出水印若源从多个远端分区读取则对每个分区的水印进行合并非源处理器直接转发输入水印延迟输入条目的处理器如异步映射处理器、窗口聚合处理器对水印做同等量的延迟。水印规则WM rules无论处理器如何生成/转发水印都必须遵守以下规则单调性严格规则水印永远不能回退。Jet 要求严格单调递增因为重复发射相同值的水印是资源浪费违反将导致作业失败Job failure。在 WatermarkCoalescer.java 中observeWm()会在新值不大于旧值时抛出JetException(Watermarks not monotonically increasing on queue: ...)这正是该规则的落地校验。不超车约定规则输入时未迟到的事件不应在处理器输出时变成迟到事件。无多余延迟约定规则处理器应在确认不会有迟到事件紧随其后时尽快发射水印。只有第 1 条是硬性要求第 2、3 条是约定convention处理器理论上可以违背——目前 Jet 中还没有这样的处理器。AbstractProcessor为直接立即转发水印的简单处理器实现了默认情况其tryProcessWatermark实现仅将收到的水印原样发射见 AbstractProcessor.java注释明确写道 This basic implementation only forwards the passed watermark。面向多键的扩展Watermark 类新增 key 字段为支持多键Watermark类新增了一个byte key字段。key 不同的两个水印实例被完全独立处理。这一点在源码中已落地实现——Watermark.java 提供了两个构造器Watermark(long timestamp)默认 key 为0Watermark(long timestamp, byte key)显式指定 key。类的equals/hashCode同时比较timestamp与keytoString()在值为IDLE_MESSAGE_TIME时输出Watermark{IDLE_MESSAGE}否则输出Watermark{ts..., key...}。该类的 Javadoc 将 key 描述为用于区分来自不同源的水印标识符。Processor API 的扩展当前Processor接口原有方法boolean tryProcessWatermark(Nonnull Watermark watermark);由于参数已包含 key方法签名本身无需修改。但该 API 并非向后兼容——处理器现在可能需要检查watermark的 key如果不检查它可能先观察到key1的较新水印实例随后又观察到key2的较旧水印实例在忽略 key 的情况下这构成非单调序列。AbstractProcessor无需改动——它把收到的每个水印直接发射不受此变化影响见 Processor.java 对两个重载版本的详细说明。然而上述 API 对JOIN 处理器是不够的JOIN 处理器从左右两个输入接收不同 key的水印。例如订单流与投递流 JOIN 时每个输入使用不同的水印 key处理器才能独立追踪每条流的进度。但无 ordinal 版本的方法只有收到两个输入上相同 key 的水印之后才会被调用——于是这种情况下手处理器永远收不到任何水印实例。为此新增方法boolean tryProcessWatermark(int ordinal, Nonnull Watermark watermark);该方法收到的是汇聚到输入边ordinal的多个上游处理器实例合并后的水印但多个输入边之间不合并、水印立即送达。该方法在 Processor.java 中以默认返回true的形式存在Javadoc 标注since 5.2。处理器可以根据自身需求选择覆盖哪个版本JOIN 处理器需要区分各输入进度 → 覆盖带ordinal的版本合并处理器期望所有输入具有相同水印 → 覆盖无ordinal的版本。需要特别注意的是若某个 key 的水印未从所有输入边收到那么无ordinal版本的方法对该 key永远不会被调用参见 Processor.java 的说明。此外对于 at-least-once 处理保障的作业重启后同一水印可能被重复投递处理器可能被要求处理比重启前更旧的水印——实现者需要容忍这一点。不同 key 水印的合并不同 key 的水印完全独立地合并keyK的水印只与来自其他输入的、同为keyK的水印合并就像它们是两条独立的流一样。这一语义在源码中由 KeyedWatermarkCoalescer.java 实现它内部维护MapByte, WatermarkCoalescer coalescers按 key 惰性创建各自的单键合并器coalescer(byte key)使用computeIfAbsent当一个 key 的合并器首次被创建时会把此前已完成的队列doneQueues与已处于 idle 的队列idleQueues状态迁移进去保证新 key 合并器从一开始就拥有正确的输入状态。observeWm()对每个存在的 key 调用其合并器并收集所有发射值返回一个ListWatermark。对 union/merge 处理器的改动Jet 中没有独立的 union 处理器union 是通过map(identity())配合多输入实现的。这个设置隐含假设所有输入具有相同的水印且相同的水印 key 对应事件中的同一字段。如果某个 key 只出现在一个输入上而其他输入没有该 key 的水印将永远不会被转发——因为合并会无限等待另一侧传来相同 key 的水印。因此用户必须正确选择水印 key同一字段用同一 key不同字段用不同 key才能保证功能正确。对异步变换处理器的改动Ordered有序模式该处理器内部维护一个队列事件与水印按收到顺序共同入队从而保持了事件相对于水印的原始顺序因此几乎无需改动。唯一需要注意的是一个优化当队尾元素是水印时会将其替换而非追加。在多键场景下需要额外校验——被替换的水印必须与待替换水印key 相同。该优化逻辑的源码位于 AsyncTransformUsingServiceOrderedP.java。Unordered无序模式无序处理器按异步响应到达的顺序发射结果。为处理水印它记录每个水印之间收到的事件计数并为每个有待处理响应的事件记录它是在哪个水印之后收到的响应到达时递减计数当计数归零时按从旧到新的顺序检查各水印的计数发射直至第一个非零计数处的水印。多键改造的关键点是不再按单一水印值存储计数而是按元组{wmKey, wmValue}存储。对应实现见 AsyncTransformUsingServiceUnorderedP.java。对窗口聚合处理器的改动窗口聚合处理器在转发输入水印之前会先发射所有事件已到齐的窗口。处理器接收一个时间戳函数列表每个输入 ordinal 对应一个用这些函数从事件中提取时间戳并为事件分配窗口。考虑一个带两个水印时间戳time1、time2的流下面的查询是合法的select window_start window_start_inner, window_end window_end_inner, time2, count(*) from tumble(my_stream, descriptor(time1), 5 seconds) group by 1, 2, 3它按time1计算窗口同时按time2分组——time2应作为输出水印。可以看出处理器只能按一个水印值time1来计算窗口另一个水印值time2可以作为分组键的一部分参与分组如果它不属于分组键则其值不会出现在输出上对应的水印可以丢弃。实现层面若要在发射窗口后输出time2的水印处理器需要遍历所有进行中的窗口及其 key找出最小值再与输入水印取最小值作为输出——而且在条目被移出集合后还需重新迭代计算最小值代价很高。这正是设计文档明确决定不实现的原因按水印列分组的场景非常罕见。因此滑动窗口处理器会丢弃所有输入水印输出只携带窗口边界上的水印。对应实现可见 SlidingWindowP.java该类同样实现了tryProcessWatermark。按键空闲消息keyed idle message的处理idle 消息的初衷idle 消息最初是为了处理某些分区没有数据、其他分区有数据的情况。例如有 5 个分区但只有 4 个不同的 key——由于水印合并需要等待所有输入任何水印都不会被转发窗口聚合永远无法产生结果。这种场景通常是用户错误例如 key 划分不当但在真实世界尤其是玩具项目中经常发生会让用户困惑。该超时在Core API 中可配置但在Pipeline API 中硬编码为 60 秒——这直接导致了Jet pipeline 启动 60 秒后才开始产出结果这类 issue 报告。idle 消息的两处处理点idle 消息在两个位置被处理源处理器若在配置的超时时间内没有看到任何事件则生成 idle 消息水印合并器当所有被合并的输入都处于 idle 状态时转发 idle 消息。考虑这样的 DAGscan1 - map1 --\ -- aggr scan2 - map2 --/若scan1无数据idle 超时后scan1发送 idle 消息map1转发它因为其所有输入都已 idleaggr输入上的合并器将map1的输入从水印合并中排除转而转发另一路输入的水印于是aggr可以开始产生结果。为什么 idle 消息要忽略 key由于Watermark类引入了 key 字段idle 消息如今也携带 key。当前语义是收到某流的 idle 消息即标记该流 idle收到任何事件或水印即标记其 active——idle 状态按流stream追踪。这意味着 idle 消息有 key但反向激活消息事件没有 key。问题在于如果 idle 消息带 key那么每条流的状态可能变成队列对 key0 active、对 key1 idle也就是部分 idle。这在概念上不成立一个事件携带全部时间戳只要有事件到达流就既不 idle 也非部分 idle而是完全 active。因此设计结论是idle 消息中应忽略 key——任何 key 的 idle 消息都会把整条流标记为 idle并从所有 key的合并计算中排除。实现上可以只允许发送key0的 idle 消息这样设计意图在代码中更加显式。源码中 KeyedWatermarkCoalescer.java 正是这样做的IDLE_MESSAGE固定为new Watermark(IDLE_MESSAGE_TIME, (byte) 0)见 WatermarkCoalescer.javaKeyedWatermarkCoalescer.observeWm()中收到IDLE_MESSAGE时直接置idleQueues[queueIndex] true并遍历所有已存在的 key 合并器、同时在所有队列都 idle时额外发射IDLE_MESSAGE而普通水印则仅作用于其自身 key 的合并器。示例分析验证语义正确性合并算子Merge operator-aggr --merge ---join ----scan1 ----scan2 ---scan3若scan1与scan2都 idle则join被标记为 idle——尽管scan1、scan2使用了不同的水印 key。merge算子将join从水印合并中排除直接转发scan3的水印从而使aggr能够利用这些水印产生输出。注意merge算子是合并同质流的原型。类似的合并也发生在多个流汇聚到同一条 DAG 边时——例如scan1顶点由多个处理器支撑。该示例说明该算法在这种情况下同样有效。JOIN 算子-joinOuter --joinInner ---scan1 ---scan2 --scan3若joinInner的两个输入扫描都 idle则joinInner被标记为 idle——尽管两个扫描使用不同的水印 key。joinInner会发送key0的 idle 消息idle 消息唯一允许的 key。joinOuter有两个输入joinInner和scan3。joinInner被标记 idle 后不再产生任何事件或水印joinOuter将无法推进——它被joinInner不产出数据卡住这是流到流 JOIN流到批、批到批 JOIN 不使用水印。更糟的是joinOuter会无限缓冲scan3的数据直到joinInner重新 active 并送来水印。这个场景看似有问题但并非 idle 机制引起的——根因是 JOIN 的一侧不推进即使没有任何 idle 流处理也会如此。它只是说明把 idle 流从合并中排除在这里帮不上忙。结合测试验证KeyedWatermarkCoalescer 的语义上述设计在仓库中不仅有实现还有对应的单元测试 KeyedWatermarkCoalescerTest.java测试名直接反映了设计语义when_Q1ProducesWmPlusEventAndQ2IsIdle_then_forwardWmAndEventFromQ1一条队列 idle 时另一条队列的水印与事件照常转发when_noWmKeyKnown_and_idleMessageReceived_then_idleForwardedWhenAllIdle在所有队列都收到 idle 消息时转发 idle且keys()为空尚无任何 key 被观测到test_initialScenario2idle 状态与 key 的交互——先收到 idle、再收到key42的水印、随后另一队列 idle 时先前的水印被转发test_initialScenario3_idleStatusTransferredToNewWmKeyidle 状态被正确迁移到新出现的水印 key 上队列 0 已 idle 时队列 1 上新 key 的水印立即被转发test_initialScenario4队列 0 先发wm(42)再 idle、队列 1 idle 时同时输出wm(42)与IDLE_MESSAGE。这些测试覆盖了本文设计文档中idle 状态按流追踪、按 key 独立合并、idle 消息忽略 key的全部关键结论。总结按键水印keyed watermark是 Hazelcast Jet 支撑流到流 JOIN 的关键引擎级能力。核心设计可归纳为Watermark携带byte key不同 key 的水印实例彼此独立各自维护单调递增的水印序列实现见 Watermark.javaProcessorAPI 新增带 ordinal 的重载tryProcessWatermark(int ordinal, Watermark)让 JOIN 类处理器能分别追踪各输入的进度实现见 Processor.java合并按 key 独立进行KeyedWatermarkCoalescer按 key 惰性创建单键合并器并迁移 idle/done 状态实现见 KeyedWatermarkCoalescer.javaidle 消息忽略 key任何 key 的 idle 消息都把整条流标记为 idle并从所有 key 的合并中排除实现上固定使用key0。对普通转发型处理器而言AbstractProcessor的默认实现直接转发完全不受影响需要按 key 区分语义的只有 JOIN、无序异步变换等少数处理器。单调性违规水印回退会直接导致作业失败这是所有水印相关实现都必须守住的底线。上述语义均已在仓库源码与 KeyedWatermarkCoalescerTest.java 中得到验证后续如需深入流到流 JOIN 的完整设计可继续阅读 docs/design/sql/15-stream-to-stream-join.md。赞分享缓存KV存储消息队列流处理后端【免费下载链接】hazelcastHazelcast is a unified real-time data platform combining stream processing with a fast data store, allowing customers to act instantly on>项目地址https://gitcode.com/gh_mirrors/ha/hazelcast点击查看免费下载相关推荐Rufus工具拯救老旧电脑的Windows 11安装神器轻松绕过硬件限制Rufus工具拯救老旧电脑的Windows 11安装神器轻松绕过硬件限制 还在为老旧电脑无法安装Windows 11而烦恼吗TPM 2.0和安全启动要求让桌面应用开发工具Hazelcast Jet 流水线事件保序机制preserveOrder 设计与实现深度解析Hazelcast Jet 流水线事件保序机制preserveOrder 设计与实现深度解析 Hazelcast Jet现为 Hazelcast 统一实时数缓存KV存储消息队列流处理后端Element Plus水印组件Watermark页面水印与版权保护Element Plus水印组件Watermark页面水印与版权保护 还在为网站内容被轻易复制而烦恼吗还在担心敏感数据被截图泄露吗Element Plus前端UI组件上一篇3步掌控窗口尺寸Window Resizer让你的屏幕空间利用率提升60%下一篇如何快速掌握Medusa成本管理3步构建完整电商财务系统创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考