尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Akka Streams 与 Typed Actor 集成:ActorFlow.ask 操作符实战指南

Akka Streams 与 Typed Actor 集成:ActorFlow.ask 操作符实战指南 Akka Streams 与 Typed Actor 集成ActorFlow.ask 操作符实战指南【免费下载链接】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 中将 Akka Streams 的每个流元素作为一次请求发送给 Typed Actor新 Actor API并期望得到回复后继续向下游发射是流 Actor 请求-响应最常见的集成方式。ActorFlow.ask正是为此设计的流操作符它把流式处理与 Typed Actor 的 Ask 模式结合让每个元素都能与 Actor 进行一次一问一答的交互同时保持背压backpressure语义。读完本文你将掌握ActorFlow.ask及其衍生变体的签名、参数、错误处理行为、Reactive Streams 语义并能用 Scala 与 Java 两套 API 写出可直接运行的集成代码。一、什么是 ActorFlow.askActorFlow.ask位于akka.stream.typed.scaladslScala DSL与akka.stream.typed.javadslJava DSL包中是 Akka Streams Actor 互操作操作符Actor interop operators家族的一员。它使用 Ask 模式Ask Pattern将每个流元素作为一条请求消息发送给目标 Typed Actorakka.actor.typed.ActorRef[_]Actor 回复后回复消息作为流的下游元素发射出去。官方文档将其归类为 Actor interop operators与Source.ask/Flow.ask经典 Actor API 变体、ActorFlow.askWithContext、ActorFlow.askWithStatus等并列。它面向的是 Typed Actor 的请求-响应交互模式即从 Actor 外部发起 Ask的场景。与经典 Actor 变体的区别在于经典版 Flow.askclassic actors 需要传入目标ActorRef、响应类型Class与超时而ActorFlow.ask面向akka.actor.typed.ActorRef[Q]类型安全更强通过makeMessage函数在编译期约束消息类型。二、依赖引入ActorFlow.ask由akka-stream-typed模块提供需要引入该依赖。Akka 依赖通过 Akka 官方安全库仓库分发需要使用带 token 的安全 URL引入方式如下sbt配合 Akka BOM 管理版本libraryDependencies com.typesafe.akka %% akka-stream-typed % AkkaVersionMavendependency groupIdcom.typesafe.akka/groupId artifactIdakka-stream-typed_2.13/artifactId version${akka.version}/version /dependencyGradledependencies { implementation com.typesafe.akka:akka-stream-typed_2.13:${akkaVersion} }该模块同时依赖akka-stream与akka-actor-typed因此你在项目中可以同时使用流 DSL 与 Typed Actor API。本仓库中该模块的主源码位于 akka-stream-typed/src/main/scala。三、签名与核心概念ActorFlow.ask的 Scala 签名如下摘自 scaladsl/ActorFlow.scaladef askI, Q, A(makeMessage: (I, ActorRef[A]) Q)( implicit timeout: Timeout): Flow[I, A, NotUsed] def askI, Q, A(ref: ActorRef[Q])(makeMessage: (I, ActorRef[A]) Q)( implicit timeout: Timeout): Flow[I, A, NotUsed]Java 签名摘自 javadsl/ActorFlow.scalapublic static I, Q, A FlowI, A, NotUsed ask( ActorRefQ ref, java.time.Duration timeout, BiFunctionI, ActorRefA, Q makeMessage) public static I, Q, A FlowI, A, NotUsed ask( int parallelism, ActorRefQ ref, java.time.Duration timeout, BiFunctionI, ActorRefA, Q makeMessage)三个泛型参数的含义类型参数含义I流进入该操作符的输入元素类型Q目标 Actor 能理解的消息类型问题消息类型AActor 回复的消息类型同时成为该 Flow 的输出元素类型ActorFlow.ask要求调用者提供三个要素actor ref目标 Typed Actor 的ActorRef[Q]makeMessage 函数接收(输入元素 I, 回复 ActorRef[A])构造并返回要发送给 Actor 的消息Q。ActorRef[A]就是 Ask 模式中的replyToActor 用它回发回复timeout隐式akka.util.TimeoutScala或显式java.time.DurationJava。特别注意Scala API 中必须显式指定响应类型A通过类型标注或类型参数否则编译器会推断为Nothing导致类型错误。源码注释中明确提醒了这一陷阱scaladsl/ActorFlow.scala#L59-L68推荐写法为flow.via(ActorFlow.askString, Asking, Reply((el, replyTo) Asking(el, replyTo))) // 或更简洁 flow.via(ActorFlow.askString, Asking, Reply(Asking(_, _)))四、使用示例4.1 消息与 Actor 定义文档示例与测试源码ActorFlowSpec.scala 与 ActorFlowCompileTest.java展示了标准写法。首先定义消息类型Asking携带流元素内容和replyTo引用Reply是 Actor 的回复// Scala final case class Asking(s: String, replyTo: ActorRef[Reply]) final case class Reply(msg: String)// Java static class Asking { final String payload; final ActorRefReply replyTo; public Asking(String payload, ActorRefReply replyTo) { this.payload payload; this.replyTo replyTo; } } static class Reply { public final String msg; public Reply(String msg) { this.msg msg; } }Actor 端使用Behaviors.receiveMessage处理Asking并通过replyTo回发Replyval ref spawn(Behaviors.receiveMessage[Asking] { asking asking.replyTo ! Reply(asking.s !!!) Behaviors.same })4.2 Scala 用法implicit val timeout: Timeout 1.second val askFlow: Flow[String, Reply, NotUsed] ActorFlow.ask(ref)(Asking.apply) // 显式构造消息的等价写法 val askFlowExplicit: Flow[String, Reply, NotUsed] ActorFlow.ask(ref)(makeMessage (el, replyTo: ActorRef[Reply]) Asking(el, replyTo)) val result: Future[Seq[String]] Source(1 to 50).map(_.toString).via(askFlow).map(_.msg).runWith(Sink.seq)4.3 Java 用法Duration timeout Duration.ofSeconds(1); // 方法引用写法 FlowString, Reply, NotUsed askFlow ActorFlow.ask(actorRef, timeout, Asking::new); // 显式构造消息的等价写法 FlowString, Reply, NotUsed askFlowExplicit ActorFlow.ask(actorRef, timeout, (msg, replyTo) - new Asking(msg, replyTo)); Source.repeat(hello).via(askFlow).map(reply - reply.msg).runWith(Sink.seq(), system);测试源码中验证了完整行为Source.repeat(hello).via(askFlow).take(3).runWith(Sink.seq)得到List(Reply(hello!!!), Reply(hello!!!), Reply(hello!!!))并且Source(1 to 50)场景下回复按提交顺序严格有序ActorFlowSpec.scala#L145-L169。五、底层实现与运行机制从源码看ask的全部重载最终收敛到私有方法askImplscaladsl/ActorFlow.scala#L25-L53其核心链路为ref.toClassic通过 Typed 到 Classic 的适配器akka.actor.typed.scaladsl.adapter把 TypedActorRef[Q]转换为 ClassicActorRef.watch(classicRef)让该流阶段监视目标 Actor一旦 Actor 终止流会以WatchedActorTerminatedException失败——这正是文档中failswhen the passed-in actor terminates的来源.mapAsync(parallelism)对每个元素调用akka.pattern.extended.ask为元素构造一次请求并拿到Future由mapAsync负责并发与背压.mapError做异常归一化AskTimeoutExceptionClassic Ask 的超时异常转换为java.util.concurrent.TimeoutException——Typed 生态统一使用 JDK 的TimeoutExceptionWatchedActorTerminatedException被重命名为ask()阶段名便于定位是哪一级流阶段失败.named(ask)为阶段命名方便日志与调试。由此可以看出文档中如果任何一次 ask 超时流将以AskTimeoutException实际表现为TimeoutException失败这一描述的实现方式。5.1 并行度parallelism默认值不带parallelism参数的重载默认并行度为2scaladsl/ActorFlow.scala#L96。源码注释解释了设计动机当第一个 ask 消息还在被处理时第二个消息已经在 Actor 邮箱中排队提前发送第二个消息能维持更健康的吞吐量。如果需要更高吞吐可显式传parallelism 4等更大的值mapAsync语义下并行度即同时在途in-flight的 ask 数量上限。5.2 监督策略ActorFlow.ask遵循ActorAttributes.SupervisionStrategy属性Scala/Java 两端源码 javadoc 均明确标注 Adheres to the ActorAttributes.SupervisionStrategy attribute。也就是说在流图中为阶段配置监督策略如Supervision.resume/Supervision.restart时该操作符会按策略处理失败。5.3 测试验证的错误场景测试源码还覆盖了两个典型失败路径ask 超时对Behaviors.ignore的 Actor 发送消息永不回复设置Timeout 10.millis断言下游收到以Ask timed out on [Actor...开头的错误ActorFlowSpec.scala#L172-L187Actor 终止向 Actor 发送Asking(TERMINATE, ...)使其Behaviors.stopped断言流以Actor watched by [ask()] has terminated!...失败ActorFlowSpec.scala#L189-L199。六、Reactive Streams 语义ActorFlow.ask的背压与完成语义与文档一致亦是源码 javadoc 原文emits发射当内部由 ask 模式创建的 Future按提交顺序完成时发射backpressures背压当在途 Future 数量达到配置的并行度且下游施加背压时进行背压completes完成当上游完成、所有 Future 均已完成且所有元素均已发射后完成fails失败当传入的 Actor 终止或任何一次 ask 超过超时时间时失败cancels取消当下游取消时取消。需要特别强调的是emits in submission ordermapAsync保证输出顺序与输入顺序一致因此尽管多个 ask 并发在途回复仍然按提交顺序有序地下发这与无序遍历的mapAsyncUnordered形成对比。七、衍生变体askWithStatus / askWithContext同属ActorFlow对象的相关操作符均有独立文档与实现在需要时可直接选用7.1 askWithStatus当 Actor 的回复类型是akka.pattern.StatusReply[A]时使用askWithStatus 文档。StatusReply.success(value)会将value解包后向下游发射StatusReply.error(err)则让 Future 以err失败进而导致流失败。实现上它基于ask再叠加一层map解包scaladsl/ActorFlow.scala#L141-L158final case class AskingWithStatus(s: String, replyTo: ActorRef[StatusReply[String]]) // Actor 端 case asking asking.replyTo ! StatusReply.success(asking.s !!!) Behaviors.same // 流端输出类型直接是解包后的 String Source.repeat(hello) .via(ActorFlow.askWithStatus(replierWithSuccess)((el, replyTo) AskingWithStatus(el, replyTo))) .take(3)7.2 askWithContext当流携带上下文如 Kafka 分区号、请求 ID且需要把上下文与回复重组输出时使用askWithContext 文档。上下文不会随消息发送给 Actor而是在本地与回复重新组合输入类型为(I, Ctx)输出类型为(A, Ctx)。Java 侧使用akka.japi.Pairjavadsl/ActorFlow.scala#L140-L151Source.repeat(hello).zipWithIndex .via(ActorFlow.askWithContext(replier)((el, replyTo) Asking(el, replyTo))) // 输出 (Reply, Long) 元组7.3 askWithStatusAndContext两者的组合回复为StatusReply同时保留并重组上下文。其实现先将StatusReply[A]与上下文打包再统一解包或抛错scaladsl/ActorFlow.scala#L192-L202。以上变体同样提供默认并行度 2 与显式并行度两种重载且均有对应的 Scala/Java 测试用例ActorFlowSpec.scala。八、与经典 Actor 变体 Flow.ask 的对比维度ActorFlow.askTypedFlow.ask / Source.askClassic目标 Actorakka.actor.typed.ActorRef[Q]akka.actor.ActorRef响应类型由makeMessage的类型参数A保证需显式传ClassJava或类型参数SScala做mapTo转换超时隐式Timeout/java.time.Duration隐式TimeoutScala/Timeout参数Java失败类型TimeoutExceptionask 超时、WatchedActorTerminatedExceptionActor 终止AskTimeoutException、akka.actor.Status.Failure携带的 cause文档ActorFlow/ask.mdSource-or-Flow/ask.md经典变体还允许 Actor 回复akka.actor.Status消息其中Status.Failure会使操作符以Failure消息携带的 cause 失败见 Source-or-Flow/ask.md而 Typed 生态更推荐用StatusReply表达成功/失败语义。九、最佳实践小结务必显式指定响应类型Scala 中写ActorFlow.askString, Asking, Reply(...)避免类型推断为Nothing根据响应速率调并行度默认 2 兼顾吞吐与邮箱负载Actor 处理快、下游消费快时可提高并行度但注意过高的并行度会放大请求洪峰为 ask 设置合理的超时超时既保护流不被无响应 Actor 拖死也会以TimeoutException终止整个流——如需容错可在外层叠加recover/recoverWithRetries等错误处理操作符利用StatusReply表达业务失败askWithStatus能把 Actor 的业务错误转化为流失败配合监督策略使用需要携带上下文时用askWithContext上下文不经过网络传输只做本地重组避免额外消息开销。十、进一步阅读完整操作符索引与相关变体 akka-docs/src/main/paradox/stream/operators/index.mdScala 实现 akka-stream-typed/src/main/scala/akka/stream/typed/scaladsl/ActorFlow.scalaJava 实现 akka-stream-typed/src/main/scala/akka/stream/typed/javadsl/ActorFlow.scalaScala 测试 akka-stream-typed/src/test/scala/docs/scaladsl/ActorFlowSpec.scalaJava 编译测试 akka-stream-typed/src/test/java/docs/javadsl/ActorFlowCompileTest.java经典 Actor 变体 akka-docs/src/main/paradox/stream/operators/Source-or-Flow/ask.md变体文档 askWithContext · askWithStatus【免费下载链接】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),仅供参考
返回列表