Akka Streams Source.mergePrioritizedN:按优先级合并多个数据源的加权扇入操作符实战指南
【免费下载链接】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
本篇技术指南聚焦 Akka Streams 的Source.mergePrioritizedN操作符:它以加权概率的方式将多个数据源(Source)合并为单个数据流,当多个源同时就绪时按优先级"偏向"高优先级源。你将掌握该操作符的 Scala / Java 完整用法、eagerComplete完成语义的取舍、底层MergePrioritizedGraphStage 的加权随机选择算法,以及与merge、mergePreferred、mergePrioritized等同类扇入操作符的选型差异,可直接用于真实流式应用的流量合并与分级调度场景。
操作符概览与定位
mergePrioritizedN属于 Akka Streams 的 Fan-in(扇入)操作符 家族,功能是按优先级合并多个数据源(Merge multiple sources with priorities)。与无差别合并的merge不同,当多个输入源同时有元素就绪时,mergePrioritizedN会依据各源配置的优先级整数进行加权随机选择,从而让高优先级源获得更高的输出占比。
从源码结构看,它是mergePrioritized(仅支持两个源)的 N 源推广版本。核心实现位于 Source.scala:
def mergePrioritizedNT], eagerComplete: Boolean): Source[T, NotUsed] = { sourcesAndPriorities match { case immutable.Seq() => Source.empty case immutable.Seq((source, _)) => source.mapMaterializedValue(_ => NotUsed) case sourcesAndPriorities => val (sources, priorities) = sourcesAndPriorities.unzip combine(sources.head, sources(1), sources.drop(2): _*)(_ => MergePrioritized(priorities, eagerComplete)) } }注意签名约定:sourcesAndPriorities中源与优先级的数量必须一致且顺序一一对应;优先级必须为正整数。当传入 0 个源时返回Source.empty,传入 1 个源时原样透传(仅将物化值统一为NotUsed)。
优先级如何起作用:加权概率模型
理解该操作符的核心是它的选择模型。文档明确给出了三源场景下的概率公式:
当三个源
sourceA、sourceB、sourceC同时就绪时,sourceA被选中的概率为priorityOfA / (priorityOfA + priorityOfB + priorityOfC),其余源同理。
几个关键事实需要掌握:
- 只在多个源同时就绪时才谈优先级:如果某一时刻只有一个源有元素,该元素会直接输出,不存在优先级竞争;
- 子集加权:如果只有部分源就绪,则用"就绪子集"的相对优先级进行加权。例如
sourceB与sourceC就绪而sourceA未就绪时,两者按priorityOfB : priorityOfC的比例竞争; - 必须是正整数:优先级取值为正整数,
0或负数会在底层 GraphStage 构造时被拒绝(见下文源码校验)。
也就是说,优先级并不是"绝对抢占",而是"加权随机偏向"——高优先级源被选中概率更高,但低优先级源在竞争中也不会完全饿死。
完整示例:三个源按 9900 : 99 : 1 合并
Scala 示例
以下代码摘自 FlowMergeSpec.scala 的测试用例:
import akka.stream.scaladsl.{ Sink, Source } val sourceA = Source(List(1, 2, 3, 4)) val sourceB = Source(List(10, 20, 30, 40)) val sourceC = Source(List(100, 200, 300, 400)) Source .mergePrioritizedN(List((sourceA, 9900), (sourceB, 99), (sourceC, 1)), eagerComplete = false) .runWith(Sink.foreach(println)) // prints e.g. 1, 100, 2, 3, 4, 10, 20, 30, 40, 200, 300, 400 since both sources have their first element ready and // the left sourceA has higher priority - if both sources have elements ready, sourceA has a 99% chance of being picked next // while sourceB has a 0.99% chance and sourceC has a 0.01% chance该示例把"概率"落实为直观数字:9900 / (9900 + 99 + 1) = 99%、99 / 10000 = 0.99%、1 / 10000 = 0.01%。输出1, 100, 2, 3, 4, ...说明:前三轮中sourceA以压倒性概率连续胜出,但sourceC也在第 2 轮抢到一次输出——这正是加权随机的体现,每次运行结果并不确定,注释中的 "prints e.g." 即表明仅为一次可能的运行结果。
Java 示例
对应的 Java 用法摘自 SourceOrFlow.java,使用Pair列表承载"源 + 优先级":
import akka.japi.Pair; import akka.stream.javadsl.Source; import akka.NotUsed; import java.util.Arrays; import java.util.List; Source<Integer, NotUsed> sourceA = Source.from(Arrays.asList(1, 2, 3, 4)); Source<Integer, NotUsed> sourceB = Source.from(Arrays.asList(10, 20, 30, 40)); Source<Integer, NotUsed> sourceC = Source.from(Arrays.asList(100, 200, 300, 400)); List<Pair<Source<Integer, ?>, Integer>> sourcesAndPriorities = Arrays.asList(new Pair<>(sourceA, 9900), new Pair<>(sourceB, 99), new Pair<>(sourceC, 1)); Source.mergePrioritizedN(sourcesAndPriorities, false).runForeach(System.out::println, system);Java 侧的 API 定义在 javadsl/Source.scala:它接收java.util.List[Pair[Source[T, _], Integer]],内部转换为 Scala 的Seq[(Source, Int)]后委托给 Scala 版实现,最终物化值统一为NotUsed(输入源各自的物化值被丢弃)。
eagerComplete 参数:完成语义的选择
mergePrioritizedN的第二个参数eagerComplete: Boolean决定上游完成时合并流的行为:
| eagerComplete | 完成行为 |
|---|---|
false(默认) | 等待所有上游完成,合并流才 complete |
true | 只要任意一个上游完成,立即取消其余上游并 complete |
对应文档中的 Reactive Streams 语义即为:"completes when all upstreams complete (or when any upstream completes ifeagerComplete=true.)"。
该逻辑在 Graph.scala 的onUpstreamFinish中实现:eagerComplete=true时取消所有输入并直接completeStage();否则递减runningUpstreams计数,直到全部上游关闭才完成。需要提醒的是,eagerComplete=true意味着未消费完的元素会被丢弃,适合"任一数据源结束即可停止整体"的场景;而默认false更贴近"必须等所有源都发完"的完整合并语义。
底层原理:MergePrioritized GraphStage 的加权随机选择
mergePrioritizedN最终通过combine构造一个 MergePrioritized 的GraphStage[UniformFanInShape[T, T]]。构造时的前置校验(require)直接决定了上文"正整数优先级"的约束:
require(priorities.nonEmpty, "A Merge must have one or more input ports") require(priorities.forall(_ > 0), "Priorities should be positive integers")其选择算法位于select()方法,逻辑分两步:
- 求和:遍历所有输入,对处于 available(就绪)状态的输入累加其优先级得到
tp;若tp == 0(无输入就绪)返回null,等待下游再次 pull; - 加权随机命中:用
SplittableRandom生成[0, tp)的随机数r,再次遍历就绪输入,依次r -= priorities(ix),当r < 0时即选中该输入——这等价于把区间[0, tp)按各就绪源的优先级比例切分,随机落点落在哪段就选哪个源。
此外,preStart中会对所有输入tryPull预取,onPush时若下游可用且无其他就绪输入则立即转发,避免无谓的竞争延迟。这些实现细节印证了文档对概率模型的描述,也解释了为何输出顺序具有随机性。
与同类扇入操作符的选型对比
| 操作符 | 输入源数量 | 选择策略 | 适用场景 |
|---|---|---|---|
merge | 多个 | 完全随机、无差别 | 不需要区分来源的普通合并 |
mergePreferred | 2 | 硬性偏向(preferred 源总是优先) | 严格"主从"分流,但可能饿死非优先源 |
mergePrioritized | 2 | 按优先级加权随机 | 双源按比例分级调度 |
mergePrioritizedN | N(≥2) | 按优先级加权随机 | 多源按比例分级调度(本文主题) |
四者的完整 Reactive Streams 语义可归纳为(mergePrioritizedN专属语义见下节):
- emits:当某个输入有元素可用时;若多个输入同时就绪,优先选择高优先级输入;
- backpressures:当下游背压时;
- completes:所有上游完成(若
eagerComplete=true则任一上游完成即完成); - cancels:下游取消时。
mergePrioritizedN的返回类型为Source[T, NotUsed],即输入源的物化值(如Future、Ref等)不会透传,统一映射为NotUsed。
使用要点与限制
- 优先级为正整数:传
0或负数会触发require异常(IllegalArgumentException),务必校验业务侧传入的优先级; - 顺序对应:
sourcesAndPriorities的源与优先级必须同序,源码注释明确要求 "same size and order"; - 输出非确定性:加权随机意味着输出序列每次运行都可能不同,需要确定性输出的场景请改用
mergeSorted或自定义GraphStage; - 低优先级不会饿死:只要低优先级源有元素且下游持续 pull,它仍会按比例被选中;若下游吞吐远低于上游总供给,高优先级源会占据绝大部分输出;
- 物化值丢弃:如果依赖某个输入源的物化值(如
Source.queue的SourceQueue),应在合并前通过其它途径持有引用。
参考实现路径
- 操作符定义:scaladsl/Source.scala
- Java API:javadsl/Source.scala
- 底层 GraphStage(选择算法、完成逻辑):scaladsl/Graph.scala
- 双源版
mergePrioritized(Flow API):scaladsl/Flow.scala - Scala 测试与运行示例:FlowMergeSpec.scala
- Java 文档示例:SourceOrFlow.java
【免费下载链接】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),仅供参考