后端并发编程异步编程【免费下载链接】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 完整实现了 Reactive Streams 标准能够以非阻塞背压back pressure的方式与任何遵循该标准的流式库互通。本篇指南基于官方文档 reactive-streams-interop.md结合仓库内 ReactiveStreamsDocSpec.scala 与 ReactiveStreamsDocTest.java 的真实测试示例以及Source、Sink、Flow工厂与JavaFlowSupport的底层源码完整讲解Publisher、Subscriber、Processor三类接口如何进出 Akka Streams 图帮助你把 Akka Streams 与 Reactor、RxJava 等第三方流式库对接起来。读完你将掌握Source 端接入外部 Publisher、Sink 端输出外部 Subscriber、Flow 与 Processor 双向转换、fanout 广播语义及背压缓冲调优。依赖配置Akka 官方组件发布在 Akka 的安全库仓库需要通过 https://account.akka.io/token 获取带 token 的安全 URL 访问。使用 Akka Streams 需要先在项目中加入akka-stream模块依赖官方推荐通过 BOMBill of Materials统一管理版本!-- Maven -- properties akka.version2.x.x/akka.version scala.binary.version2.13/scala.binary.version /properties dependencyManagement dependencies dependency groupIdcom.typesafe.akka/groupId artifactIdakka-bom_${scala.binary.version}/artifactId version${akka.version}/version typepom/type scopeimport/scope /dependency /dependencies /dependencyManagement dependency groupIdcom.typesafe.akka/groupId artifactIdakka-stream_${scala.binary.version}/artifactId /dependency// sbt val AkkaVersion 2.x.x libraryDependencies com.typesafe.akka %% akka-stream % AkkaVersion// Gradle def akkaVersion 2.x.x dependencies { implementation platform(com.typesafe.akka:akka-bom_2.13:${akkaVersion}) implementation com.typesafe.akka:akka-stream_2.13 }版本号请以你实际使用的 Akka 版本为准当前仓库对应 Akka 2.9 系列$akka.version$占位符由文档构建时自动替换。概述Akka Streams 对 Reactive Streams 标准的实现Reactive Streams 是异步流处理与非阻塞背压的行业标准。Akka Streams 实现了该标准因此可以把 Akka Streams 的图与任何同样遵循标准的流式库自由拼接。在 API 形态上Reactive Streams 存在两个版本org.reactivestreams包Java 8 时代的独立 artifactorg.reactivestreams:reactive-streams定义了Publisher、Subscriber、Subscription、Processor四个接口java.util.concurrent.Flow包自 Java 9 起这四个接口被并入 JDK 标准库位于java.util.concurrent.Flow命名空间下接口语义与org.reactivestreams完全一致。Akka Streams 对两套接口都提供了互操作能力org.reactivestreams接口直接通过常规Source/SinkAPI 上的工厂方法接入即Source.fromPublisher、Sink.fromSubscriber、Sink.asPublisher、Source.asSubscriber、Flow.toProcessor、Flow.fromProcessorJava 9 内置的java.util.concurrent.Flow.*接口则通过独立的工厂集合接入Scala 端为akka.stream.scaladsl.JavaFlowSupportJava 端为akka.stream.javadsl.JavaFlowSupport。下文示例均以org.reactivestreams工厂演示每一处调用都可以原样替换为JavaFlowSupport中的对应方法并使用 JDK 的java.util.concurrent.Flow.*接口。需要特别注意的是JavaFlowSupport无法在 Java 8 上使用因为所需接口不在 JDK 8 标准库中这一点在两个JavaFlowSupport的源码注释中也明确标注了 For use only with JDK 9参见 scaladsl/JavaFlowSupport.scala 与 javadsl/JavaFlowSupport.java。核心接口与一个完整的端到端示例Reactive Streams 中最重要的是两个接口Publisher数据发布者向下游推送元素并响应需求信号和Subscriber数据订阅者向上游提交需求并接收元素。先看导入Scala来自 ReactiveStreamsDocSpec.scalaimport org.reactivestreams.Publisher import org.reactivestreams.Subscriber import org.reactivestreams.ProcessorJava来自 ReactiveStreamsDocTest.javaimport org.reactivestreams.Publisher; import org.reactivestreams.Subscriber; import org.reactivestreams.Processor;假设某个第三方库对外提供了一个推文的Publisher// Scalalibrary provides a publisher of tweets def tweets: Publisher[Tweet]// Java PublisherTweet tweets();另一个库知道如何把作者账号存进数据库它接收一个Subscriber[Author]// Scala def storage: Subscriber[Author]// Java SubscriberAuthor storage();现在用 Akka Streams 的Flow做变换并把两者连接起来——先过滤出带有 Akka 标签的推文再取出作者// Scala过滤 取作者 val authors Flow[Tweet].filter(_.hashtags.contains(akkaTag)).map(_.author) // 连接Publisher 作为 SourceSubscriber 作为 Sink Source.fromPublisher(tweets).via(authors).to(Sink.fromSubscriber(storage)).run()// Java final FlowTweet, Author, NotUsed authors Flow.of(Tweet.class).filter(t - t.hashtags().contains(AKKA)).map(t - t.author); Source.fromPublisher(rs.tweets()).via(authors).to(Sink.fromSubscriber(rs.storage()));这里可以看到最核心的互操作模式Source.fromPublisher(publisher)把外部Publisher当作流的输入SourceSink.fromSubscriber(subscriber)把外部Subscriber当作流的输出Sink。Flow夹在中间完成元素变换与背压传递。该示例在测试中的完整断言ReactiveStreamsDocSpec.scala验证了 7 位作者按序到达并最终收到完成信号证明背压与完成传播都工作正常。将 Flow 转换为 Processor一个Flow还可以被转换为RunnableGraph[Processor[In, Out]]当调用run()时物化出一个Processor实例。run()可以被多次调用每次都会得到新的Processor实例即可重用// Scala val processor: Processor[Tweet, Author] authors.toProcessor.run() tweets.subscribe(processor) // 外部 Publisher 订阅 Processor 的输入端 processor.subscribe(storage) // Processor 的输出端被 Subscriber 订阅// Java final ProcessorTweet, Author processor authors.toProcessor().run(system); rs.tweets().subscribe(processor); processor.subscribe(rs.storage());Processor既是Publisher又是Subscriber所以可以直接用subscribe方法把它接到两端的标准接口上。其底层实现详见后文“源码级实现解析”正是Source.asSubscriberFlowSink.asPublisher(false)的拼接物化。将 Source 暴露为 Publisher把Source暴露给外部世界使用的是Publisher-Sink即Sink.asPublisher// Scala val authorPublisher: Publisher[Author] Source.fromPublisher(tweets).via(authors).runWith(Sink.asPublisher(fanout false)) authorPublisher.subscribe(storage)// Java final PublisherAuthor authorPublisher Source.fromPublisher(rs.tweets()) .via(authors) .runWith(Sink.asPublisher(AsPublisher.WITHOUT_FANOUT), system); authorPublisher.subscribe(rs.storage());注意runWith的物化值就是这个Publisher拿到后即可把它交给任何 Reactive Streams 兼容的消费者。单订阅者fanout false使用Sink.asPublisher(fanout false)Java 端为AsPublisher.WITHOUT_FANOUT创建的Publisher只支持一次订阅。后续再尝试订阅会被拒绝并抛出IllegalStateException。这一点在两个JavaFlowSupport的Sink.asPublisher文档注释中均有明确说明scaladsl/JavaFlowSupport.scala、javadsl/JavaFlowSupport.java。多订阅者fanout true广播如果需要支持多个订阅者就要使用 fan-out / 广播模式// Scala一个 Publisher 同时服务 storage 与 alert 两个 Subscriber def alert: Subscriber[Author] val authorPublisher: Publisher[Author] Source.fromPublisher(tweets).via(authors).runWith(Sink.asPublisher(fanout true)) authorPublisher.subscribe(storage) authorPublisher.subscribe(alert)// Java final PublisherAuthor authorPublisher Source.fromPublisher(rs.tweets()) .via(authors) .runWith(Sink.asPublisher(AsPublisher.WITH_FANOUT), system); authorPublisher.subscribe(rs.storage()); authorPublisher.subscribe(rs.alert());fanout trueAsPublisher.WITH_FANOUT下物化出的Publisher会向每个订阅者广播同样的元素流。测试 ReactiveStreamsDocSpec.scala 中两个订阅者都完整收到了全部作者且断言注释特别说明“该测试依赖 fanout Publisher 的缓冲区大小大于作者数量”。fanout 的背压语义与 input-buffer该运算符的输入缓冲区大小input buffer size决定了最快订阅者与最慢订阅者之间最多能拉开多大差距超出后整个流会被迫减速。也就是说fanout true的Publisher会以最慢的订阅者为节拍器快的订阅者最多领先慢的订阅者input-buffer个元素随后上游处理被背压抑制。若测试中的作者数量大于缓冲大小慢订阅者就会因缓冲区溢出而丢数据或报错——这正是测试注释的用意。将 Sink 暴露为 Subscriber反过来用 Subscriber-SourceSource.asSubscriber可以把一个Sink暴露为外部可订阅的Subscriber// Scala val tweetSubscriber: Subscriber[Tweet] authors.to(Sink.fromSubscriber(storage)).runWith(Source.asSubscriber[Tweet]) tweets.subscribe(tweetSubscriber)// Java final SubscriberTweet tweetSubscriber authors.to(Sink.fromSubscriber(storage)).runWith(Source.asSubscriber(), system); rs.tweets().subscribe(tweetSubscriber);这里Source.asSubscriber[Tweet]的物化值是一个Subscriber[Tweet]外部库可以把整个 Akka Streams 图当作一个“会消费推文的 Subscriber”直接subscribe。Source.asSubscriber的实现定义于 Source.scala包装内部SubscriberSource阶段是“Sink 端接收外部订阅”的镜像操作。将 Processor 重新包装为 FlowProcessor实例也可以重新包装回Flow但必须传入一个创建 Processor 的工厂函数而非直接传入 Processor 实例。只有工厂才能保证结果Flow的可重用性——因为一个Processor只能被物化使用一次而工厂可以在每次物化时生成新实例// Scala // An example Processor factory def createProcessor: Processor[Int, Int] Flow[Int].toProcessor.run() val flow: Flow[Int, Int, NotUsed] Flow.fromProcessor(() createProcessor)// Java final CreatorProcessorInteger, Integer factory new CreatorProcessorInteger, Integer() { public ProcessorInteger, Integer create() { return Flow.of(Integer.class).toProcessor().run(system); } }; final FlowInteger, Integer, NotUsed flow Flow.fromProcessor(factory);如果你需要保留物化值可以使用Flow.fromProcessorMat它接受返回(Processor, Mat)二元组的工厂函数Scala 端定义见 Flow.scala。源码级实现解析以上工厂并非黑盒魔法它们的底层实现在仓库里清晰可见理解后能帮你判断背压行为与资源开销。1.Sink.asPublisher的两个分支Scaladsl 的 Sink.asPublisher 根据fanout选择两个不同的内部阶段def asPublisherT: Sink[T, Publisher[T]] fromGraph( if (fanout) new FanoutPublisherSinkT) else new PublisherSinkT))在 impl/Sinks.scala 中可以看到两者的本质区别PublisherSink非 fanout创建一个VirtualPublisher作为物化值。此时流上还没有任何订阅者因此流在预取元素填满内部缓冲区后会一直保持背压直到订阅者接入并产生需求。源码注释还说明它同时注册了StreamSubscriptionTimeout用于处理“流已启动但迟迟无人订阅”的超时场景FanoutPublisherSinkfanout直接启动一个名为FanoutProcessorImpl的 actor通过context.materializer.actorOf(...)用ActorProcessor实现广播分发这正是它天然支持多订阅者的原因也是它具备独立输入缓冲区语义的根源。2.Flow.toProcessor的组装逻辑Flow.scala 中的 toProcessor 内部就是把Source.asSubscriber、Flow、Sink.asPublisher(false)串成一个图物化后把Subscriber端与Publisher端打包成一个匿名Processordef toProcessor: RunnableGraph[Processor[In uncheckedVariance, Out uncheckedVariance]] Source .asSubscriber[In] .via(this) .toMat(Sink.asPublisherOut)(Keep.both[Subscriber[In], Publisher[Out]]) .mapMaterializedValue { case (sub, pub) new Processor[In, Out] { override def onError(t: Throwable): Unit sub.onError(t) override def onSubscribe(s: Subscription): Unit sub.onSubscribe(s) override def onComplete(): Unit sub.onComplete() override def onNext(t: In): Unit sub.onNext(t) override def subscribe(s: Subscriber[_ : Out]): Unit pub.subscribe(s) } }javadsl.JavaFlowSupport.Flow.toProcessor做了同样的事但面向java.util.concurrent.Flow.Processorjavadsl/JavaFlowSupport.java。注意它固定使用WITHOUT_FANOUT因此toProcessor产出的Processor输出端同样只接受单个订阅者。3.Flow.fromProcessor与ProcessorModuleFlow.fromProcessor最终落到fromGraph(ProcessorModule(processorFactory))Flow.scala。ProcessorModule定义在 impl/StreamLayout.scala它把“按需调用工厂创建 Processor”的语义固化为一个图阶段由 PhasedFusingActorMaterializer.scala 中的ProcessorModulePhase负责物化——每次物化都调用一次工厂这正是“Flow 可重用”的保障。4.JavaFlowSupport的适配转换Java 9 的java.util.concurrent.Flow.*接口并非另起炉灶而是通过 impl/JavaFlowAndRsConverters 提供的asRs/asJava隐式转换在Flow接口与org.reactivestreams接口之间做零成本适配。例如// scaladsl/JavaFlowSupport.scala def fromPublisherT: Source[T, NotUsed] scaladsl.Source.fromPublisher(publisher.asRs) def asSubscriber[T]: Source[T, java.util.concurrent.Flow.Subscriber[T]] scaladsl.Source.asSubscriber[T].mapMaterializedValue(_.asJava)也就是说JavaFlowSupport只是org.reactivestreams工厂的适配壳Scala 端见 scaladsl/JavaFlowSupport.scalaJava 端见 javadsl/JavaFlowSupport.java其运行时行为与对应org.reactivestreams工厂完全一致。仓库还提供了 JavaFlowSupportCompileTest.java 编译级测试来保证这层适配的 API 完整性。与其他 Reactive Streams 实现互操作实现 Reactive Streams 标准的价值在于Akka Streams 可以与同样遵循该标准的其他流式库直接拼接无需胶水代码。官方文档列举的部分其他实现包括Reactor1.1RxJava通过 RxJavaReactiveStreams 桥接RatpackSlick无论对接哪种库交互方式都一致对方提供Publisher就用Source.fromPublisher接进来对方需要Subscriber就用Sink.fromSubscriber送出去需要双向能力就用Flow.toProcessor/Flow.fromProcessor。这也意味着你在 Akka Streams 图内享受到的全部能力背压、扇出/扇入、异步边界、故障恢复等在跨越库边界后依然由 Reactive Streams 的信号协议完整承载。小结场景工厂方法物化值外部 Publisher 进图Source.fromPublisher(pub)NotUsed外部 Subscriber 出图Sink.fromSubscriber(sub)NotUsedSource 暴露为 Publisher单订阅Sink.asPublisher(fanout false)Publisher[T]Source 暴露为 Publisher广播Sink.asPublisher(fanout true)Publisher[T]Sink 暴露为 SubscriberSource.asSubscriber[T]Subscriber[T]Flow 转为 Processor可多次 runflow.toProcessor.run()Processor[In, Out]Processor 包回 Flow需工厂Flow.fromProcessor(() proc)NotUsed对于 Java 9 用户将上表所有方法替换为JavaFlowSupport中的对应版本、接口换成java.util.concurrent.Flow.*即可获得完全等价的互操作能力。完整的可运行示例可继续阅读 ReactiveStreamsDocSpec.scala 与 ReactiveStreamsDocTest.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点击查看免费下载相关推荐Akka Streams Sink.asPublisher 完全指南将 Akka Stream 桥接到 Reactive Streams PublisherAkka Streams Sink.asPublisher 完全指南将 Akka Stream 桥接到 Reactive Streams Publisher后端并发编程异步编程Reactive Streams四大核心组件详解Publisher、Subscriber、Subscription、ProcessorReactive Streams四大核心组件详解Publisher、Subscriber、Subscription、Processor Reactive St后端Akka Streams 设计原理从不可变蓝图到 Reactive Streams 互操作Akka Streams 设计原理从不可变蓝图到 Reactive Streams 互操作 Akka Streams 是 Akka 平台中面向分布式有界流处理后端并发编程异步编程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考