
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.scaladef 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; SourceInteger, NotUsed sourceA Source.from(Arrays.asList(1, 2, 3, 4)); SourceInteger, NotUsed sourceB Source.from(Arrays.asList(10, 20, 30, 40)); SourceInteger, NotUsed sourceC Source.from(Arrays.asList(100, 200, 300, 400)); ListPairSourceInteger, ?, 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默认等待所有上游完成合并流才 completetrue只要任意一个上游完成立即取消其余上游并 complete对应文档中的 Reactive Streams 语义即为completes when all upstreams complete (or when any upstream completes ifeagerCompletetrue.)。该逻辑在 Graph.scala 的onUpstreamFinish中实现eagerCompletetrue时取消所有输入并直接completeStage()否则递减runningUpstreams计数直到全部上游关闭才完成。需要提醒的是eagerCompletetrue意味着未消费完的元素会被丢弃适合任一数据源结束即可停止整体的场景而默认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多个完全随机、无差别不需要区分来源的普通合并mergePreferred2硬性偏向preferred 源总是优先严格主从分流但可能饿死非优先源mergePrioritized2按优先级加权随机双源按比例分级调度mergePrioritizedNN≥2按优先级加权随机多源按比例分级调度本文主题四者的完整 Reactive Streams 语义可归纳为mergePrioritizedN专属语义见下节emits当某个输入有元素可用时若多个输入同时就绪优先选择高优先级输入backpressures当下游背压时completes所有上游完成若eagerCompletetrue则任一上游完成即完成cancels下游取消时。mergePrioritizedN的返回类型为Source[T, NotUsed]即输入源的物化值如Future、Ref等不会透传统一映射为NotUsed。使用要点与限制优先级为正整数传0或负数会触发require异常IllegalArgumentException务必校验业务侧传入的优先级顺序对应sourcesAndPriorities的源与优先级必须同序源码注释明确要求 same size and order输出非确定性加权随机意味着输出序列每次运行都可能不同需要确定性输出的场景请改用mergeSorted或自定义GraphStage低优先级不会饿死只要低优先级源有元素且下游持续 pull它仍会按比例被选中若下游吞吐远低于上游总供给高优先级源会占据绝大部分输出物化值丢弃如果依赖某个输入源的物化值如Source.queue的SourceQueue应在合并前通过其它途径持有引用。参考实现路径操作符定义scaladsl/Source.scalaJava APIjavadsl/Source.scala底层 GraphStage选择算法、完成逻辑scaladsl/Graph.scala双源版mergePrioritizedFlow APIscaladsl/Flow.scalaScala 测试与运行示例FlowMergeSpec.scalaJava 文档示例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),仅供参考