如何在自定义序列化器中复用Akka Persistence Cassandra的CqlSession
问题描述
我在Akka中为应用实现自定义序列化器,该序列化器需要偶尔访问Cassandra数据库获取数据。我通过Akka Extension将仓库注册到ActorSystem,并在序列化器的构造函数中获取该仓库实例。
当前存在的问题是:仓库依赖CqlSession连接Cassandra,我目前通过Alpakka Cassandra创建了新的CqlSession,但Cassandra官方明确建议每个应用仅应使用一个CqlSession。现在我的应用里同时存在两个活跃的CqlSession:一个由Akka Persistence Cassandra维护,另一个是自定义序列化器仓库创建的。请问如何复用Akka Persistence Cassandra的CqlSession,让整个应用仅保留一个活跃的CqlSession?
原代码示例
自定义序列化器代码
public class CustomSerializer extends AsyncSerializerWithStringManifestCS { private final CustomSerializerRepo customSerializerRepo; public CustomSerializer(ExtendedActorSystem system) { super(system); this.system = system; customSerializerRepo= customSerializerRepoExtension.get(Adapter.toTyped(system)).customSerializerRepo(); } @Override public CompletionStage<byte[]> toBinaryAsyncCS(Object unencryptedPayload) { // 使用仓库进行操作 } @Override public CompletionStage<Object> fromBinaryAsyncCS(byte[] bytes, String manifest) { // 反序列化逻辑 } }
仓库代码
public class CustomSerializerRepo { private final CassandraSession cassandraSession; public CustomSerializerRepo (ActorSystem<?> actorSystem) { var cassandraSession = CassandraSessionRegistry.get(actorSystem).sessionFor(CassandraSessionSettings.create()) } // Cassandra数据操作逻辑 }
解决方案
Akka Persistence Cassandra内部通过CassandraSessionRegistry管理着一个全局共享的CqlSession,你只需要通过指定与Akka Persistence Cassandra一致的配置路径,就能复用这个已存在的Session,无需新建。
具体实现步骤:
- Akka Persistence Cassandra的Session默认配置路径是
akka.persistence.cassandra.session,如果你的应用没有自定义配置前缀,直接使用这个路径即可。 - 在仓库中获取Session时,传入该配置路径创建
CassandraSessionSettings,CassandraSessionRegistry会自动返回已初始化的共享Session实例。
修改后的仓库代码如下:
public class CustomSerializerRepo { private final CassandraSession cassandraSession; public CustomSerializerRepo (ActorSystem<?> actorSystem) { // 加载Akka Persistence Cassandra的Session配置,复用已有Session var sessionSettings = CassandraSessionSettings.create("akka.persistence.cassandra.session"); this.cassandraSession = CassandraSessionRegistry.get(actorSystem).sessionFor(sessionSettings); } // Cassandra数据操作逻辑 }
额外注意事项:
- 确保Akka Persistence Cassandra的Session已完成初始化:如果序列化器在ActorSystem启动早期被加载(比如持久化Actor启动阶段),建议在执行数据库操作前调用
cassandraSession.ready(),等待Session就绪后再执行逻辑,避免连接未就绪的异常。 - 若你自定义了Akka Persistence Cassandra的配置前缀,需将上述代码中的配置路径替换为你实际使用的路径。
内容的提问来源于stack exchange,提问作者Milad Amery
相关产品推荐
相关产品推荐

