如何将基于Actor的逻辑迁移至Akka Streams?附场景示例
Akka Streams 实现 CreateUser 流程思路
作为从纯 Actor 转到 Akka Streams 的开发者,我完全理解你想利用流处理的背压、并发控制优势来重构部分功能的需求。针对你描述的 CreateUser 流程,我会分步骤拆解实现思路,结合代码示例说明:
1. 先定义领域模型
首先要明确流程中的数据结构和错误类型,这是流处理的基础:
// 请求与结果模型 case class CreateUserRequest(email: String, name: String, extraInfo: String) case class UserCreated(userId: String) // 错误类型 sealed trait CreateUserError case object EmailAlreadyInUse extends CreateUserError case object NoUserInInternalDB extends CreateUserError case class ThirdPartyServiceFailed(msg: String) extends CreateUserError
2. 封装原有服务为异步调用
你的 Actor 服务可以封装成返回 Future 的接口,方便在 Akka Streams 中使用(流处理天生支持异步操作):
// 假设这些是你原有 Actor 服务的封装 trait EmailCheckService { def isEmailRegistered(email: String): Future[Boolean] } trait ThirdPartyUserService { def fetchExternalUser(email: String): Future[Option[ExternalUser]] } trait InternalDbService { def isUserExists(userId: String): Future[Boolean] def createUser(request: CreateUserRequest, extUser: ExternalUser): Future[String] }
3. 构建核心流处理流程
Akka Streams 的核心是 Source→Flow→Sink 的流水线,我们可以通过分支(Branch)来处理不同的错误分支:
第一步:邮箱校验分支
先创建一个流来处理邮箱校验,将结果拆分为「邮箱已存在」和「校验通过」两个分支:
val emailCheckFlow: Flow[CreateUserRequest, Either[CreateUserError, CreateUserRequest], _] = Flow[CreateUserRequest].mapAsync(parallelism = 4) { req => emailCheckService.isEmailRegistered(req.email).map { case true => Left(EmailAlreadyInUse) case false => Right(req) } } // 拆分分支:错误流和正常流 val (invalidEmailSink, validRequestSource) = emailCheckFlow.branch( { case Left(_) => true }, // 左值(错误)进入第一个分支 { case Right(_) => true } // 右值(正常)进入第二个分支 ) // 处理邮箱已存在的情况:日志、通知等 invalidEmailSink.to(Sink.foreach { case Left(EmailAlreadyInUse) => println(s"邮箱 ${req.email} 已被注册") })
第二步:第三方调用与内部DB校验
接下来处理校验通过的请求,调用第三方服务并检查内部DB是否存在该用户,同样拆分错误分支:
val thirdPartyCheckFlow: Flow[CreateUserRequest, Either[CreateUserError, (CreateUserRequest, ExternalUser)], _] = Flow[CreateUserRequest].mapAsync(parallelism = 4) { req => thirdPartyUserService.fetchExternalUser(req.email).flatMap { case Some(extUser) => internalDbService.isUserExists(extUser.userId).map { case true => Right((req, extUser)) case false => Left(NoUserInInternalDB) } case None => Future.successful(Left(ThirdPartyServiceFailed("第三方服务未找到用户"))) } } // 拆分分支:内部DB无用户的错误流,以及正常流 val (noUserInDbSink, validUserSource) = thirdPartyCheckFlow.branch( { case Left(_) => true }, { case Right(_) => true } ) // 处理内部DB无用户的情况:比如触发内部用户创建逻辑 noUserInDbSink.via(handleNoUserInDbFlow).to(Sink.foreach { case _ => println("已触发内部用户创建流程") })
第三步:最终创建用户
最后处理所有校验通过的请求,完成用户创建:
val createUserFlow: Flow[(CreateUserRequest, ExternalUser), UserCreated, _] = Flow[(CreateUserRequest, ExternalUser)].mapAsync(parallelism = 4) { case (req, extUser) => internalDbService.createUser(req, extUser).map(UserCreated) } // 将创建结果输出到日志或事件总线 validUserSource.via(createUserFlow).to(Sink.foreach { user => println(s"用户创建成功,ID:${user.userId}") })
4. 进阶:使用 GraphDSL 处理复杂流
如果你的流程有更多分支、合并逻辑,建议使用 GraphDSL 来构建更清晰的流拓扑:
val createUserGraph = RunnableGraph.fromGraph(GraphDSL.create() { implicit builder => import GraphDSL.Implicits._ // 定义各个节点 val source = Source.queue[CreateUserRequest](bufferSize = 100, OverflowStrategy.backpressure) val emailCheck = builder.add(emailCheckFlow) val emailBranch = builder.add(Balance[Either[CreateUserError, CreateUserRequest]](2)) val thirdPartyCheck = builder.add(thirdPartyCheckFlow) val thirdPartyBranch = builder.add(Balance[Either[CreateUserError, (CreateUserRequest, ExternalUser)]](2)) val createUser = builder.add(createUserFlow) val invalidEmailSink = Sink.foreach[Either[CreateUserError, CreateUserRequest]](...) val noUserSink = Sink.foreach[Either[CreateUserError, (CreateUserRequest, ExternalUser)]](...) val successSink = Sink.foreach[UserCreated](...) // 连接拓扑 source ~> emailCheck ~> emailBranch emailBranch.out(0) ~> invalidEmailSink emailBranch.out(1).map(_.right.get) ~> thirdPartyCheck ~> thirdPartyBranch thirdPartyBranch.out(0) ~> noUserSink thirdPartyBranch.out(1).map(_.right.get) ~> createUser ~> successSink ClosedShape }) // 运行流 createUserGraph.run()
关键注意点
- 并发控制:
mapAsync的parallelism参数可以控制异步操作的并发数,避免压垮下游服务 - 背压处理:Akka Streams 自动处理背压,当下游处理不过来时,上游会自动减速
- 错误处理:除了分支拆分,还可以使用
recover、rescue操作符来捕获并处理流中的异常,或者自定义SupervisionStrategy来决定错误发生时的行为(继续、停止、重启) - Actor 集成:如果需要调用原有 Actor,可以使用
AskPattern将 Actor 消息转换为Future,无缝集成到流中
内容的提问来源于stack exchange,提问作者Alex Fruzenshtein
相关产品推荐
相关产品推荐

