You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Akka Actor中Sink生产正常调用,测试环境未执行问题排查

排查Akka Stream在测试中未执行的问题

从你的现象来看,生产环境中流能正常执行,但测试里只有post source run输出,没有inside flow,说明流的run()方法被调用了,但没有处理任何元素。下面是几个最可能的排查方向:

1. AccountCache的borrowAccount未返回有效账户

你的代码里构建submittingPaymentsWithAccounts时依赖accountCache.borrowAccount.map(ps -> _):

val submittingPaymentsWithAccounts: Seq[(Seq[Payment], Account)] = payments.grouped(100).flatMap(ps => accountCache.borrowAccount.map(ps -> _)).toSeq

如果borrowAccount返回None,那这个Seq就是空的,对应的Source没有任何元素可以处理,自然不会进入Flow的map阶段(也就看不到inside flow)。

测试中你发送了UpdateAccount(account)消息,但需要确认:

  • UpdateAccount消息的处理逻辑是否真的把传入的account添加到了AccountCache中?
  • borrowAccount方法是否在账户被添加后返回Some(account)?

可以在测试中添加断言,或者在Actor内部添加调试输出,确认submittingPaymentsWithAccounts非空:

// 在Actor的processPayments方法里添加
println(s"Number of payment-account pairs: ${submittingPaymentsWithAccounts.size}")

2. ActorMaterializer的创建方式不正确

在Akka Actor中创建ActorMaterializer时,正确的做法是绑定到Actor的上下文,这样Materializer的生命周期会和Actor绑定。你的代码里是:

implicit private val materializer: ActorMaterializer = ActorMaterializer()

这种方式创建的Materializer没有关联到Actor的Context,在测试环境中可能因为生命周期管理问题,导致流还没执行就被终止了。

修改为使用Actor上下文创建Materializer:

implicit private val materializer: ActorMaterializer = ActorMaterializer(context)

(如果是Akka 2.6+版本,推荐使用Materializer(context)替代ActorMaterializer)

3. 流处理的Future未正确完成

你的Flow中使用了mapAsync,依赖前面的Future完成:

.map { case (ps, account) =>
  println("inside flow")
  // ... 返回Future[(TransactionResponse, Seq[Payment], Account)]的代码块
}
.mapAsync(parallelism = config.accounts.size)(_.map { ... })

如果这个Future永远不会完成(比如测试中mock的依赖没有正确返回已完成的Future),mapAsync会一直等待,流的处理就会卡在这里。

检查你返回的Future是否在测试环境中能正确完成:

  • 如果是调用外部服务的逻辑,测试中是否用mock返回了Future.successful(...)?
  • 有没有可能Future被阻塞或者进入了死锁?

4. Actor状态切换的时机问题

虽然你调试确认source.run()被调用了,但需要再确认触发ProcessPayments时,Actor的状态是否满足nextKnownPaymentDate.exists(_.isBefore(ZonedDateTime.now()))的条件:

  • 测试中你设置了when(repo.earliestTimeDue).thenReturn(Some(ZonedDateTime.now())),但ZonedDateTime.now()在设置mock和触发ProcessPayments时的时间差是否会导致条件不满足?
  • UpdateNextPaymentTime消息的处理逻辑是否正确更新了nextKnownPaymentDate?

可以在processPayments方法开头添加日志,确认进入了正确的分支:

case ProcessPayments =>
  println(s"nextKnownPaymentDate: $nextKnownPaymentDate, current time: ${ZonedDateTime.now()}")
  if (nextKnownPaymentDate.exists(_.isBefore(ZonedDateTime.now()))) {
    // ... 现有逻辑
  }

额外测试建议

  • 在测试中使用Akka的TestProbe监听Actor的内部状态,或者添加调试输出,确认submittingPaymentsWithAccounts的内容;
  • 直接测试流的逻辑:把paymentSink提取出来,单独编写单元测试,传入测试用的Source,验证流是否能正确处理元素,排除Actor环境的干扰;
  • 检查测试中的ExecutionContext是否有异常,比如是否有未捕获的异常导致流终止(可以添加一个watchTermination的回调来监控流的状态):
Source.fromIterator(...)
  .toMat(paymentSink)(Keep.both)
  .run()
  ._2.onComplete {
    case Success(_) => println("流正常完成")
    case Failure(e) => println(s"流失败: ${e.getMessage}")
  }

内容的提问来源于stack exchange,提问作者Synesso

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.12 04:30:20