Akka Projection测试持续失败:测试投影与Actor间消息交互时遭遇超时问题求解
问题根源:断言函数被重复执行导致二次等待超时
你遇到的问题核心在于**ProjectionTestKit.run会反复调用传入的断言函数**,直到Projection运行结束或断言抛出异常。你的代码第一次执行时成功拿到了消息并打印,但第二次执行断言函数时,TestProbe已经没有新消息了,expectMessageType会一直阻塞等待,最终触发超时错误。
官方文档提到的"断言函数会被反复调用直至无错误完成",这里的"无错误"是指每次执行断言函数都不能抛出异常——而你的断言函数不是幂等的,第一次执行成功,第二次执行就会因为没有新消息而超时。
解决方案
这里提供几种可行的修复方式:
方案1:改用幂等的断言逻辑
不直接等待新消息,而是检查TestProbe是否已经收到目标消息,这样多次执行断言函数也不会阻塞:
projectionTestKit.run(projection) { // 先确认已经收到消息 assert(mailer.receivedMessages.nonEmpty, "必须收到至少一条SendEmailMessage") // 取出消息并验证内容 val message = mailer.receivedMessages.head.asInstanceOf[SendEmailMessage] println(s"Got $message") // 可选:添加对消息具体内容的断言 assert(message.recipient == EmailAddress.parse("user@example.org").right.get, "收件人地址不匹配") }
方案2:用标志位确保断言只执行一次
通过原子变量标记断言是否已经完成,避免重复执行等待逻辑:
import java.util.concurrent.atomic.AtomicBoolean val alreadyAsserted = new AtomicBoolean(false) projectionTestKit.run(projection) { // 只有第一次执行时才做断言 if (alreadyAsserted.compareAndSet(false, true)) { val message = mailer.expectMessageType[SendEmailMessage] println(s"Got $message") // 这里可以添加消息内容的验证逻辑 } }
方案3:断言成功后主动停止Projection
在拿到消息后直接停止Projection,让run方法结束运行:
projectionTestKit.run(projection) { val message = mailer.expectMessageType[SendEmailMessage] println(s"Got $message") // 停止Projection,触发run方法结束 projectionTestKit.stop(projection) }
额外注意事项
确保你的Projection在处理完所有事件后能自动停止:
- 你使用的
TestSourceProvider基于有限的Seq,发送完事件后Source会正常完成,Projection应该会自动停止 - 如果你的
MyEventHandler.handle方法包含异步逻辑,一定要返回正确的Future,确保Projection等待异步操作完成后再处理下一个事件,避免Projection提前停止但消息还未发送的情况
内容的提问来源于stack exchange,提问作者gervais.b
相关产品推荐
相关产品推荐

