Akka DurableStateBehavior测试问题:求替代PersistenceTestKit的测试套件及示例
Akka DurableStateBehavior 测试方案及代码示例
核心结论
Akka Durable State 对应的测试套件是DurableStateTestKit,原有的PersistenceTestKit仅适配Event Sourced场景,无法用于验证DurableState的状态持久化,这就是你测试失败的原因。
步骤1:确保依赖正确
如果使用Maven,确保引入与Akka版本匹配的测试依赖:
<dependency> <groupId>com.typesafe.akka</groupId> <artifactId>akka-persistence-testkit_2.13</artifactId> <version>${akka.version}</version> <scope>test</scope> </dependency>
Gradle则配置:
testImplementation "com.typesafe.akka:akka-persistence-testkit_2.13:$akkaVersion"
步骤2:修改测试代码
替换原有PersistenceTestKit为DurableStateTestKit,并调整状态验证逻辑:
import akka.persistence.testkit.DurableStateTestKit; import akka.persistence.testkit.DurableStateTestKitPlugin; // 保留原有其他import @ExtendWith(MockitoExtension.class) class InterimActorTest { @ClassRule public static final TestKitJunitResource testKit = new TestKitJunitResource( DurableStateTestKitPlugin.getInstance() .config() .withFallback(ConfigFactory.defaultApplication().resolve())); // 替换PersistenceTestKit为DurableStateTestKit private DurableStateTestKit durableStateTestKit = DurableStateTestKit.create(testKit.system()); private static ClusterSharding sharding; private TestProbe<ParentActor.Command> parentActorProbe; private TestProbe<ChildActor.Command> childActorProbe; @Mock private MyClient myClient; private static final String jobId = "test-job-id"; // 自行定义DEFAULT_ACCOUNT_ID、TWO等常量 @BeforeAll public static void init() { if (sharding == null) { Cluster cluster = Cluster.get(testKit.system()); cluster.manager().tell(new Join(cluster.selfMember().address())); sharding = ClusterSharding.get(testKit.system()); sharding.init( Entity.of( InterimActor.TYPE_KEY, entityContext -> ProbedInterimActor.create( myClient, parentActorProbe, childActorProbe, entityContext.getShard(), cleanup))); } } @BeforeEach public void setUp() throws JsonProcessingException { // 清理所有测试状态 durableStateTestKit.clearAll(); parentActorProbe = testKit.createTestProbe(ParentActor.Command.class); childActorProbe = testKit.createTestProbe(ChildActor.Command.class); reset(myClient); } @Test void testStartToInterimActor() { EntityRef<InterimActor.Command> underTest = sharding.entityRefFor(InterimActor.TYPE_KEY, jobId); Item item = ItemObject.getFiniteItem(); when(myClient.streamContacts(DEFAULT_ACCOUNT_ID, item.getContactList(), null)) .thenReturn(responseFlux); underTest.tell(new InterimActor.Start(jobId, item, DEFAULT_ACCOUNT_ID)); // 验证持久化状态:替换原expectedPersisted,匹配Actor实际持久化的状态 durableStateTestKit.expectState(jobId, state -> { // 根据实际状态类型做断言,比如验证类型或字段值 if (!(state instanceof InterimActor.Started)) return false; InterimActor.Started startedState = (InterimActor.Started) state; return startedState.jobId().equals(jobId) && startedState.item().equals(item); }); // 保留原有业务逻辑验证 underTest.tell(new InterimActor.Pace(2, 2)); childActorProbe.expectMessageClass(ChildActor.Start.class); StepVerifier.create(responseFlux) .expectSubscription() .thenRequest(2) .expectNext(jsonNode, jsonNode) .thenRequest(2) .expectNext(jsonNode, jsonNode) .expectComplete(); } }
关键差异说明
- 套件替换:用
DurableStateTestKit替代PersistenceTestKit,对应插件换成DurableStateTestKitPlugin - 状态验证:原
expectPersisted针对事件,DurableState用expectState(persistenceId, predicate)验证状态对象 - 清理逻辑:用
durableStateTestKit.clearAll()清理测试状态
额外注意事项
- 确保Actor的
persistenceId与测试中使用的jobId一致(或匹配Actor实际的persistenceId生成规则) - 若需要更松散的验证,可简化断言为仅判断状态类型:
state instanceof InterimActor.Started
内容的提问来源于stack exchange,提问作者Nilendra Jain
相关产品推荐
相关产品推荐

