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

如何在自定义序列化器中复用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,无需新建。

具体实现步骤:

  1. Akka Persistence Cassandra的Session默认配置路径是akka.persistence.cassandra.session,如果你的应用没有自定义配置前缀,直接使用这个路径即可。
  2. 在仓库中获取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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 20:23:19