Akka StreamsgroupedWeighted操作符完全指南按元素权重聚合流【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址: https://gitcode.com/gh_mirrors/ak/akka-coregroupedWeighted是 Akka Streams 中用于按权重分批聚合元素的核心操作符它不会像grouped那样按元素个数切分而是由一个costFn为每个元素计算权重不断累积直到总权重达到或超过阈值minWeight才把这一批元素作为一个Seq/List向下游发射。本文基于 Akka 官方文档与仓库源码完整讲解其签名、分组语义、边界行为、背压特性、代码示例与源码级实现原理并对比grouped、groupedWithin、groupedWeightedWithin等邻近操作符帮助你在实际流式应用中正确选型与使用。概述什么时候需要按重量而不是按数量分组在很多真实场景中元素本身大小并不均匀例如一批日志行长短不一、一批网络包字节数不同、一批数据库记录大小各异。如果按固定条数分组可能造成每组负载严重不均而groupedWeighted允许你自定义每个元素的代价权重按累计代价切分批次从而实现更均匀、更可控的分组。它的定位属于文档中的 Simple operators 一族是Source与Flow共有的通用操作符。文档原文对它的定义是Accumulate incoming events until the combined weight of elements is greater than or equal to the minimum weight and then pass the collection of elements downstream.即持续累积流入的事件直到元素的合计权重大于或等于最小权重minWeight再把这批元素的集合传递给下游。签名Scala 与 Java 两种 API 形态groupedWeighted同时提供 Scala DSL 与 Java DSL 两套 API二者语义完全一致只是函数式接口与返回容器类型不同。Scala 签名见 Flow.scala 与 Source.scaladef groupedWeighted(minWeight: Long)(costFn: Out Long): Repr[immutable.Seq[Out]]Java 签名见 Flow.scala 与 Source.scalaFlowIn, ListOut, Mat groupedWeighted( long minWeight, java.util.function.FunctionOut, java.lang.Long costFn)参数含义参数类型说明minWeightLong触发分组的最小累计权重阈值必须大于 0否则流在初始化阶段即抛出IllegalArgumentExceptioncostFnOut LongScala/FunctionOut, LongJava为每个元素计算权重的函数返回值不允许为负数否则阶段会以IllegalArgumentException失败从 Java DSL 的实现看Flow.scala它内部委托给 Scala 的delegate.groupedWeighted(minWeight)(costFn.apply)再把结果immutable.Seq通过.map(_.asJava)转换为java.util.List输出源码注释中标注了// TODO optimize to one step说明该转换目前是一个独立的映射步骤。行为语义如何累积、何时发射、何时结束groupedWeighted的核心行为如下与文档中的 Reactive Streams semantics 一致emits发射当元素的累计权重大于或等于minWeight时或上游完成时此时剩余不足一组的数据也会作为一个缩小版分组被发射backpressures背压当一个分组已经组装完毕而下游仍在背压时completes完成当上游完成时。一个重要的边界细节是最后一组可能小于minWeight。因为流结束时即使累积权重未达阈值尚未发射的缓冲元素也必须被推送出去否则数据会丢失。文档明确说明the last group possibly smaller than requestedminWeightdue to end-of-stream。底层实现Ops.scala是一个GraphStage[FlowShape[T, immutable.Seq[T]]]其状态机非常直观构造时执行require(minWeight 0, minWeight must be greater than 0)校验内部维护一个Vector.newBuilder[T]缓冲和计数器left初始值为minWeight每次onPush()时取出元素调用costFn得到cost若cost 0直接failStage抛出IllegalArgumentException(Negative weight [...] for element [...] is not allowed)否则把元素加入缓冲、left - cost当left 0时取走缓冲内容并push(out, elements)同时重置left minWeightonUpstreamFinish()时如果缓冲非空先把剩余元素作为最后一组推送给下游再completeStage()。官方示例把多个Seq按内部长度分批文档给出的示例来自 GroupedWeighted.scalaScala与 SourceOrFlow.javaJava演示了将多个子集合按元素个数作为权重进行分组并展示了分组结果如何与runForeach等其他操作串联。Scala 示例import akka.actor.ActorSystem import akka.stream.scaladsl.Source import scala.collection.immutable implicit val system: ActorSystem ActorSystem() val collections immutable.Iterable(Seq(1, 2), Seq(3, 4), Seq(5, 6)) // 权重 每个子集合的长度minWeight 4 Source[Seq[Int]](collections).groupedWeighted(4)(_.length).runForeach(println) // Vector(Seq(1, 2), Seq(3, 4)) ← 2 2 4达到阈值 // Vector(Seq(5, 6)) ← 流结束剩余不足一组也发射 // 权重 每个子集合的长度minWeight 3 Source[Seq[Int]](collections).groupedWeighted(3)(_.length).runForeach(println) // Vector(Seq(1, 2), Seq(3, 4)) ← 2 2 4 ≥ 3 // Vector(Seq(5, 6)) ← 2 ≥ 3 不成立但流结束被迫发射Java 示例import akka.actor.ActorSystem; import akka.stream.javadsl.Source; import java.util.Arrays; ActorSystem system ActorSystem.create(); Source.from(Arrays.asList(Arrays.asList(1, 2), Arrays.asList(3, 4), Arrays.asList(5, 6))) .groupedWeighted(4, x - (long) x.size()) .runForeach(System.out::println, system); // [[1, 2], [3, 4]] // [[5, 6]] Source.from(Arrays.asList(Arrays.asList(1, 2), Arrays.asList(3, 4), Arrays.asList(5, 6))) .groupedWeighted(3, x - (long) x.size()) .runForeach(System.out::println, system); // [[1, 2], [3, 4]] // [[5, 6]]逐组推演以 Scala、minWeight 4为例读入Seq(1, 2)权重 2累计 2 4继续读入Seq(3, 4)权重 2累计 4 ≥ 4发射Vector(Seq(1, 2), Seq(3, 4))读入Seq(5, 6)权重 2累计 2 4缓冲等待上游完成把剩余的Vector(Seq(5, 6))作为最后一组发射。注意minWeight 3时输出与4完全相同因为前两个元素累计 4 已超过阈值 3触发点未变最后一个子集合权重 2 依然不足 3直到流结束才被发射。这说明阈值不是精确上限而是达到即发射的下限。源码验证边界行为与测试用例groupedWeighted的各类边界行为在 FlowGroupedWeightedSpec.scala 中有系统的测试覆盖可以作为行为契约的权威参考场景测试设置期望结果空源空输入任意大权重函数不产生任何分组输出空序列权重恒为 0costFn _ 0LminWeight 1整个源收敛为一个分组永不达阈值直到流结束一次性发射恰好达阈值16 个元素、每个权重 1、minWeight 16整个源恰好一个分组阈值超过总量minWeight Long.MaxValue整个源一个分组靠流结束触发正常分组元素1,2,1,2、权重即元素值、minWeight 2两组Seq(1,2)、Seq(1,2)未达阈值不发射权重恒 0、minWeight 10上游未完成等待期间不发射任何分组上游sendComplete()后才发射Seq(1)并完成minWeight为负groupedWeighted(-1)初始化即抛IllegalArgumentException: minWeight must be greater than 0minWeight为 0groupedWeighted(0)同上初始化即抛异常costFn返回负数costFn _ -1L阶段以IllegalArgumentException: Negative weight [-1] for element [1] is not allowed失败这些用例直接对应 Ops.scala 中的require校验与 Ops.scala 中的负权重检查是理解该操作符错误语义的第一手资料。与相邻操作符的对比与选型文档在 See also 中列出了三个邻近操作符选型要点如下操作符分组依据时间窗口典型场景grouped固定元素个数无每 N 条消息一批groupedWithin固定元素个数 时间窗口有每 N 条或超时即发groupedWeighted元素权重累计无按负载/字节数等不均等元素分批groupedWeightedWithin元素权重累计 时间窗口有按负载分批且需保证最大延迟如果你的场景还需要最晚等待时间约束——例如既要按字节数分批、又不能让第一批迟迟不发——则应选择groupedWeightedWithin其实现Ops.scala在GroupedWeighted基础上叠加了定时器GroupedWeightedWithinTimer通过scheduleWithFixedDelay实现周期触发。相关行为在 FlowGroupedWithinSpec.scala 中有完整覆盖。使用注意与最佳实践结合文档与源码使用groupedWeighted时请牢记以下几点minWeight必须为正数0或负数会在流初始化阶段直接抛出IllegalArgumentException而不是等到运行时因此应尽早校验。costFn返回值不得为负负权重会导致阶段失败并终止整个流权重为0是合法的但要注意它会使分组永不达标最终整个源被收敛为一个分组参见测试用例。最后一组可能不满足阈值这是设计使然保证上游完成时不丢数据下游处理时不要假设每组的累计权重都 ≥minWeight。背压语义分组一旦组装完成就会等待下游消费因此groupedWeighted不会无界缓冲但权重函数计算成本高时它会在每个元素到达时同步执行属于热路径上的 CPU 开销。内存权衡未达阈值前元素会累积在内部Vector中若权重普遍很小而minWeight很大可能积压大量元素需要结合groupedWeightedWithin或上游限流来兜底。延伸阅读操作符总览Stream operators index按个数分组grouped 按个数 时间窗分组groupedWithin 按权重 时间窗分组groupedWeightedWithin核心实现Ops.scalaGroupedWeightedGraphStageScala API 定义Flow.scala、Source.scalaJava API 定义Flow.scala、Source.scala官方示例Scala GroupedWeighted.scala Java SourceOrFlow.java行为测试FlowGroupedWeightedSpec.scala【免费下载链接】akka-coreA platform to build and run apps that are elastic, agile, and resilient. SDK, libraries, and hosted environments.项目地址: https://gitcode.com/gh_mirrors/ak/akka-core创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考