使用TestKit测试AbstractPersistentActorWithAtLeastOnceDelivery时TestActorRef创建失败
解决AbstractPersistentActorWithAtLeastOnceDelivery的TestKit测试失败问题
我来帮你梳理下测试中可能踩的坑,毕竟用TestActorRef测试带持久化+至少一次投递的Actor,确实有几个容易忽略的关键点:
最可能的核心问题:缺失persistenceId实现
你的示例Actor代码里没写persistenceId()方法——这是AbstractPersistentActor的强制要求!如果没实现这个方法,Actor启动时直接会抛出异常,这大概率是测试失败的首要原因。先补上这个方法:
@Override public String persistenceId() { return "my-actor-test-id"; // 给测试用的唯一标识 }
测试环境的持久化配置问题
AbstractPersistentActor依赖持久化日志(Journal)和快照存储,默认配置下可能没有启用适合测试的内存型插件,导致persist操作失败。你需要在测试的ActorSystem配置里指定内存日志:
// 测试前初始化ActorSystem时添加配置 Config testConfig = ConfigFactory.parseString(""" akka.persistence.journal.plugin = "akka.persistence.journal.inmem" akka.persistence.snapshot-store.plugin = "akka.persistence.snapshot-store.local" akka.persistence.snapshot-store.local.dir = "target/test-snapshots" akka.test.single-expect-default = 3s """).withFallback(ConfigFactory.load()); ActorSystem system = ActorSystem.create("MyActorTest", testConfig);
TestActorRef与AtLeastOnceDelivery的状态初始化
AbstractPersistentActorWithAtLeastOnceDelivery启动时会尝试恢复未确认的投递消息,测试环境下如果没有历史持久化数据,恢复流程虽然会完成,但有时候需要手动触发状态初始化,避免投递ID混乱:
TestActorRef<MyActor> actorRef = TestActorRef.create(system, MyActor.props(), "test-actor"); MyActor underlyingActor = actorRef.underlyingActor(); // 重置投递ID,避免测试间状态污染 underlyingActor.setDeliveryId(0L); // 手动触发恢复完成(测试环境可能不会自动发送RecoveryCompleted信号) underlyingActor.onRecoveryCompleted();
异步persist操作的等待问题
persist()是异步操作,如果你发送消息后立即断言状态,大概率会因为操作未完成导致断言失败。可以用两种方式处理:
- 用TestProbe监听响应:在Actor的persist回调里给发送者返回确认消息,测试中用TestProbe等待这个消息:
TestProbe probe = new TestProbe(system); actorRef.tell("test-message", probe.ref()); // 等待persist完成后的确认 probe.expectMsg("Persisted: test-message");
- 用awaitCond等待状态变化:如果Actor有记录状态的字段,用等待条件判断状态更新:
actorRef.tell("test-message", system.deadLetters()); // 等待状态更新,超时时间2秒 awaitCond(() -> "test-message".equals(underlyingActor.getLastPersistedMessage()), Duration.ofSeconds(2));
Mock注入的正确时机
如果需要给Actor注入Mock依赖,最好在创建Props的时候就传入,避免事后修改underlyingActor字段导致的初始化问题:
// 假设你的Actor依赖一个Service接口 MockService mockService = Mockito.mock(MockService.class); Props actorProps = Props.create(MyActor.class, () -> new MyActor(mockService)); TestActorRef<MyActor> actorRef = TestActorRef.create(system, actorProps, "test-actor");
完整测试示例
这里给你一个可运行的简化测试代码,涵盖上面所有要点:
import akka.actor.ActorSystem; import akka.persistence.AbstractPersistentActorWithAtLeastOnceDelivery; import akka.testkit.TestActorRef; import akka.testkit.TestKit; import com.typesafe.config.Config; import com.typesafe.config.ConfigFactory; import org.junit.AfterClass; import org.junit.BeforeClass; import org.junit.Test; import java.time.Duration; import static org.junit.Assert.*; public class MyActorTest { static ActorSystem system; @BeforeClass public static void setup() { Config testConfig = ConfigFactory.parseString(""" akka.persistence.journal.plugin = "akka.persistence.journal.inmem" akka.persistence.snapshot-store.plugin = "akka.persistence.snapshot-store.local" akka.persistence.snapshot-store.local.dir = "target/test-snapshots" akka.test.single-expect-default = 3s """).withFallback(ConfigFactory.load()); system = ActorSystem.create("MyActorTestSystem", testConfig); } @AfterClass public static void teardown() { TestKit.shutdownActorSystem(system); system = null; } @Test public void testPersistAndDeliver() { TestActorRef<MyActor> actorRef = TestActorRef.create(system, MyActor.props(), "test-my-actor"); MyActor underlyingActor = actorRef.underlyingActor(); // 初始化投递状态 underlyingActor.setDeliveryId(0L); underlyingActor.onRecoveryCompleted(); // 发送测试消息 String testMsg = "Hello Persist"; actorRef.tell(testMsg, system.deadLetters()); // 等待persist完成并更新状态 awaitCond(() -> testMsg.equals(underlyingActor.getLastPersistedMessage()), Duration.ofSeconds(2)); assertEquals(testMsg, underlyingActor.getLastPersistedMessage()); } // 你的Actor实现(补全必要方法) public static class MyActor extends AbstractPersistentActorWithAtLeastOnceDelivery { private String lastPersistedMessage; public static Props props() { return Props.create(MyActor.class); } @Override public Receive createReceive() { return receiveBuilder() .match(String.class, msg -> { persist(msg, persistedMsg -> { lastPersistedMessage = persistedMsg; // 这里可以添加投递逻辑,比如deliver(...) }); }) .build(); } @Override public String persistenceId() { return "my-actor-test-id"; } public String getLastPersistedMessage() { return lastPersistedMessage; } } }
内容的提问来源于stack exchange,提问作者MrkK
相关产品推荐
相关产品推荐

