ARTICLE DETAIL

资讯详情

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

Akka Stream 的 StreamConverters.javaCollectorParallelUnordered:并行归约 Sink 的签名、实现与实战

Akka Stream 的 StreamConverters.javaCollectorParallelUnordered:并行归约 Sink 的签名、实现与实战 后端并发编程异步编程【免费下载链接】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点击查看免费下载导读StreamConverters.javaCollectorParallelUnordered是 Akka Stream 提供的一个将 Java 8Collector以并行方式接入响应式流的 Sink 操作符。它基于图阶段Balance将上游元素分发到多个异步 worker 中分别累积再用Collector.combiner归约合并最终物化为FutureScala或CompletionStageJava。读完本文你将掌握该操作符的精确 API 签名、并行执行的数据流拓扑、与顺序版javaCollector的取舍以及可复制的使用示例和仓库内的验证用例。操作符概览根据 操作符参考文档该操作符的功能定位是创建一个 Sink它会物化为一个 scala[Future] java[CompletionStage]并在其中完成 Java 8Collector的转换transformation与归约reduction操作的结果。换句话说它把标准 Java 8java.util.stream.Collector的能力带入了 Akka Stream 的响应式管线上游的元素流入该 Sink 后被累积进可变的中间结果容器所有元素处理完毕后再经由Collector的可选 finisher 变换成最终结果整个过程的产物是一个异步完成的结果句柄。它属于StreamConverters工具对象下 Additional Sink and Source converters 这一组转换器家族与 javaCollector、asJavaStream、fromJavaStream、asInputStream等并列。与顺序版的javaCollector不同这里的归约处理是基于图Balance并行执行的Reduction processing is performed in parallel based on graphBalance见 scaladsl/StreamConverters.scala 中的源码注释因此适合元素量大、累积计算相对重的场景。API 签名原文档给出的签名如下出自 javaCollectorParallelUnordered.mdScalascaladsl.StreamConvertersdef javaCollectorParallelUnorderedT, R( collectorFactory: () java.util.stream.Collector[T, _ : Any, R]): Sink[T, Future[R]]Javajavadsl.StreamConvertersdef javaCollectorParallelUnorderedT, R( collector: akka.japi.function.Creator[Collector[T, _ : Any, R]]): Sink[T, CompletionStage[R]]关键参数语义| 参数 | 类型 | 含义 | | -- | -- | -- | |parallelism|Int| 并行分片的数量即图中Balance的输出端口数也就是同时进行累积的 worker 个数 | |collectorFactory/collector|() Collector[T, _, R]/Creator[Collector[T, _, R]]| 一个工厂函数每次需要时生产一个新的Collector实例传入的必须是工厂而非Collector本身详见下文注意事项 |类型参数T是流入元素的类型R是最终归约结果的类型中间累积容器类型A即Collector[T, A, R]中的A被泛型擦除为Any由内部状态类持有。Java 侧传入akka.japi.function.Creator后内部会包装成 Scala 的函数字面量() collector.create()最终得到CompletionStage[R]形式的物化值通过.toCompletionStage()转换见 javadsl/StreamConverters.scala。并行执行的实现原理从源码结构看parallelism 1时该 Sink 的图拓扑是一个经典的分而治之结构实现位于 scaladsl/StreamConverters.scala上游元素进入BalanceT一个公平分发的扇出阶段把元素轮流/按需分配到parallelism个输出端口每个端口连接一个worker 流水线Flow[T].fold(...).async即以fold形式做局部的顺序累积把每个元素state.update(elem)进CollectorState并用.async标记异步边界使各 worker 可并行推进各 worker 的输出汇入Merge[CollectorState[T, R]](parallelism)Merge下游再叠加一个fold使用ReducerState对各个分片的累积结果做归约——这一步调用的是Collector的combiner函数BinaryOperatorA把多批局部结果两两合并归约完成后执行state.finish()调用Collector.finisher得到最终结果R交给Sink.head[R]由此物化出一个在流完成时完成的Future[R]。这个 Sink 的默认属性名为javaCollectorParallelUnordered注册在 impl/Stages.scala 中。内部的 CollectorState 与 ReducerState并行收集的关键在于两套内部状态类定义于 impl/Sinks.scala均为InternalApi private[akka]CollectorState[T, R]负责累积一侧。FirstCollectorState在收到第一个元素时才调用collectorFactory()创建真正可变的Collector取出supplier().get()得到累积容器、accumulator()得到累积函数然后用accumulator.accept(accumulated, elem)累积元素之后的元素交给MutableCollectorState原地更新finish()时用finisher().apply(...)收尾。把工厂调用延迟到首个元素到达、且每次update都返回新实例是为了保证不同 materialization 之间绝不共享同一个可变 Collector。ReducerState[T, R]负责归约一侧。FirstReducerState收到第一批局部结果时取出collector.combiner()之后MutableReducerState.update反复执行reduced combiner(reduced, batch)空流情况下finish()会以null作为累积值调用 finishercollector.finisher().apply(null)。由此可见该操作符能否正确工作强依赖于你提供的Collector本身是可组合的它必须实现了有意义的supplier、accumulator、combiner与finisher其中combiner是并行版本独有的、串行版本根本不会触碰的组件。parallelism 1 时的退化行为源码中有一个容易忽略的重要细节scaladsl/StreamConverters.scalaif (parallelism 1) javaCollectorT, R else { ... }当parallelism 1时javaCollectorParallelUnordered会直接委托给顺序版javaCollector走Flow.foldSink.head的简单管线不再构建Balance/Merge图。这保证了即使调用方传 1也不会出现并行度 1 却绕一圈分片归约的额外开销。也正因如此parallelism的合法下限是 1传入 0 或负数将没有意义会进入 else 分支构造出端口数非法/无意义的图实践中应从 2 开始体现并行收益。与顺序版 javaCollector 的对比与取舍| 维度 |javaCollector|javaCollectorParallelUnordered| | -- | -- | -- | | 归约方式 | 单条流水线顺序累积Reduction processing is performed sequentially | 基于Balance并行累积combiner归约performed in parallel based on graphBalance | | 物化值 |Future[R]/CompletionStage[R]| 相同 | | 是否使用combiner| 不使用 | 使用 | | 并行度参数 | 无 |parallelism: Int为 1 时退化为顺序版 | | 元素到达最终结果的次序 | 累积有序 |无序Unordered——多 worker 各自的局部结果以任意顺序到达Merge并被 combiner 合并 |两条实现的事实依据分别见 scaladsl/StreamConverters.scala顺序版与同文件 L121-L162并行版。选择建议元素量大、且你的Collector具备廉价且正确的combiner如Collectors.summingInt、Collectors.toList、Collectors.joining等 JDK 标准实现时用并行版摊薄累积成本对顺序敏感或依赖流的固有次序做折叠时用顺序版。实战示例仓库的文档代码示例位于 JavaCollectorDocExample.scala顺序版供对照与 JavaCollectorDocExamples.java。并行版可以直接按如下方式使用Scalaimport java.util.stream.Collectors import akka.stream.scaladsl.{ Source, StreamConverters } import akka.actor.ActorSystem implicit val system: ActorSystem ActorSystem(demo) // 并行收集为 List注意 parallelism4 val future: scala.concurrent.Future[java.util.List[String]] Source(List(one, two, three)) .runWith(StreamConverters.javaCollectorParallelUnordered(4)(() Collectors.toList[String]()))Javaimport akka.stream.javadsl.Source; import akka.stream.javadsl.StreamConverters; import java.util.concurrent.CompletionStage; import java.util.stream.Collectors; Source.from(java.util.List.of(one, two, three)) .runWith(StreamConverters.javaCollectorParallelUnordered(4, Collectors::toList), system);并行求和仓库测试中的写法见 StreamConvertersSpec.scalaval future Source(1 to 100) .runWith(StreamConverters.javaCollectorParallelUnordered(4)(() Collectors.summingIntInt)) future.futureValue.toInt should (5050) // 12...100注意 Java 8 的Collector接口本身对collector的Characteristics如CONCURRENT、UNORDERED、IDENTITY_FINISH有一定语义约定Akka 的实现默认按可并发累积、无序归约的方式分片使用因此像Collectors.joining(, )这类实现会通过combiner得到拼接结果但拼接次序不保证与上游元素原始顺序一致——这是名字里 Unordered 的含义若业务对顺序敏感需自行权衡。使用注意事项必须传工厂而非实例collectorFactory/Creator会被多次调用——每次流 materialization 都会重建状态且在并行图里每个 worker 各自需要一份 Collector。源码注释明确提醒a flow can be materialized multiple times, so the function producing theCollectormust be able to handle multiple invocationsscaladsl/StreamConverters.scala。若直接复用同一个可变Collector会造成状态跨流共享、结果错误。工厂调用是惰性的FirstCollectorState在收到首个元素时才实例化 Collector对空流finish()会单独调用一次工厂并直接对空容器应用 finisher见 impl/Sinks.scala。测试用例 work parallelly with an empty source 验证了空流下joining(, )得到StreamConvertersSpec.scala。可复用性已验证同一 Sink 可被多次runWith且各次结果互不影响——仓库测试 be reusable with parallel version 用同一个javaCollectorParallelUnordered(4)(...)先对 1..4 求和得 10、再对 4..6 求和得 15印证工厂模式隔离了状态StreamConvertersSpec.scala。异常传播supplier/accumulator/combiner/finisher中抛出的异常会沿流水线传播并导致物化出的Future/CompletionStage以失败结束对应测试见 StreamConvertersSpec.scala 附近。验证与测试仓库内对javaCollectorParallelUnordered的覆盖测试位于 akka-stream-tests/src/test/scala/akka/stream/scaladsl/StreamConvertersSpec.scala核心断言包括Source(1 to 100)Collectors.summingIntparallelism 4最终结果为 5050验证并行累积combiner 归约的正确性空源 Collectors.joining结果为验证空流路径与 finisher 对空容器的处理Sink 复用场景下两次runWith分别得到 10 与 15验证工厂隔离、无跨 materialization 状态泄漏。这些用例同时是理解该操作符行为边界的可直接运行的参考。相关参考操作符文档javaCollectorParallelUnordered.md顺序版对照文档javaCollector.md转换器家族索引operators/index.mdScala 实现scaladsl/StreamConverters.scalaJava 实现javadsl/StreamConverters.scala内部状态实现impl/Sinks.scala默认属性注册impl/Stages.scala行为验证测试StreamConvertersSpec.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点击查看免费下载相关推荐mold 内嵌 TBB 的 parallel_reduce 归约算法全解析函数签名、两种形式与并行求和实战mold 内嵌 TBB 的 parallel_reduce 归约算法全解析函数签名、两种形式与并行求和实战 导读 mold一个现代链接器在实现并行链接流水开发工具构建工具系统编程mold 项目内 oneTBB 并行规约模式实战parallel_reduce、combinable 与确定性归约的选型与实现mold 项目内 oneTBB 并行规约模式实战parallel_reduce、combinable 与确定性归约的选型与实现 本篇技术指南以 mold 仓库开发工具构建工具系统编程Akka Streams StreamConverters.asJavaStream 详解将 Akka Sink 物化为 Java 8 Stream 的桥接之道Akka Streams StreamConverters.asJavaStream 详解将 Akka Sink 物化为 Java 8 Stream 的桥接之后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表