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

Akka Durable State整合PostgreSQL时表不存在异常求助

Akka Durable State(Play 2.8.18)与PostgreSQL集成问题排查及解决

问题背景

使用Play Framework 2.8.18集成Akka Durable State特性,选用PostgreSQL作为存储介质,但创建Actor时提示durable_state表不存在,同时存在依赖版本冲突问题。


1. Durable State Actor实现代码

public class CompanyActorTest extends DurableStateBehavior<CompanyActorTest.CompanyCommand, CompanyActorTest.CompanyState> {

    public interface CompanyCommand {
    }

    public static class GetCompany implements CompanyCommand {

        public final ActorRef<Object> replyTo;
        public final String userEmail;

        public GetCompany(ActorRef<Object> replyTo, String userEmail) {
            this.replyTo = replyTo;
            this.userEmail = userEmail;
        }
    }

    public static class CompanyState {
        private final Set<Company> companies;

        public CompanyState(Set<Company> companies) {
            this.companies = companies;
        }

        public Set<Company> getCompanies() {
            return companies;
        }
    }

    public static class AddCompany implements CompanyCommand {
        public final Company company;

        public AddCompany(Company company) {
            this.company = company;
        }
    }

    private final ActorContext<CompanyCommand> ctx;

    public static Behavior<CompanyCommand> create(PersistenceId persistenceId) {
        return Behaviors.setup(ctx -> new CompanyActorTest(persistenceId, ctx));
    }

    private CompanyActorTest(PersistenceId persistenceId, ActorContext<CompanyCommand> ctx) {
        super(persistenceId);
        this.ctx = ctx;
    }

    @Override
    public CompanyState emptyState() {
        return new CompanyState(Collections.emptySet());
    }

    @Override
    public CommandHandler<CompanyCommand, CompanyState> commandHandler() {
        return newCommandHandlerBuilder()
                .forAnyState()
                .onCommand(AddCompany.class, (state, command) -> {
                    state.companies.add(command.company);
                    return Effect().persist(new CompanyState(state.companies));
                })
                .onCommand(
                        GetCompany.class, (state, command) -> {
                            ActorRef<GetCompanyActor.GetCompany> ref = this.ctx.getSystem().systemActorOf(GetCompanyActor.create(), "Get-Company-Actor", Props.empty());
                            ref.tell(new GetCompanyActor.GetCompany(command.replyTo, command.userEmail));
                            return Effect().none();
                        })
                .build();
    }
}

2. application.conf配置

akka.persistence.state.plugin = "akka.persistence.r2dbc.state"

akka.persistence.r2dbc {

  # postgres or yugabyte
  dialect = "postgres"

  # set this to your database schema if applicable, empty by default
  schema = ""

  connection-factory {
    driver = "postgres"

    # the connection can be configured with a url, eg: "r2dbc:postgresql://<host>:5432/<database>"
    url = ""

    # The connection options to be used. Ignored if 'url' is non-empty
    host = "localhost"
    port = 5432
    database = "postgres"
    user = "postgres"
    password = "postgres"

    ssl {
      enabled = off
      # See PostgresqlConnectionFactoryProvider.SSL_MODE
      # Possible values:
      #  allow - encryption if the server insists on it
      #  prefer - encryption if the server supports it
      #  require - encryption enabled and required, but trust network to connect to the right server
      #  verify-ca - encryption enabled and required, and verify server certificate
      #  verify-full - encryption enabled and required, and verify server certificate and hostname
      #  tunnel - use a SSL tunnel instead of following Postgres SSL handshake protocol
      mode = ""

      # Server root certificate. Can point to either a resource within the classpath or a file.
      root-cert = ""

      # Client certificate. Can point to either a resource within the classpath or a file.
      cert = ""

      # Key for client certificate. Can point to either a resource within the classpath or a file.
      key = ""

      # Password for client key.
      password = ""
    }

    # Initial pool size.
    initial-size = 5
    # Maximum pool size.
    max-size = 20
    # Maximum time to create a new connection.
    connect-timeout = 3 seconds
    # Maximum time to acquire connection from pool.
    acquire-timeout = 5 seconds
    # Number of retries if the connection acquisition attempt fails.
    # In the case the database server was restarted all connections in the pool will
    # be invalid. To recover from that without failed acquire you can use the same number
    # of retries as max-size of the pool
    acquire-retry = 1

    # Maximum idle time of the connection in the pool.
    # Background eviction interval of idle connections is derived from this property
    # and max-life-time.
    max-idle-time = 30 minutes

    # Maximum lifetime of the connection in the pool.
    # Background eviction interval of connections is derived from this property
    # and max-idle-time.
    max-life-time = 60 minutes

    # Configures the statement cache size.
    # 0 means no cache, negative values will select an unbounded cache
    # a positive value will configure a bounded cache with the passed size.
    statement-cache-size = 5000

    # Validate the connection when acquired with this SQL.
    # Enabling this has some performance overhead.
    # A fast query for Postgres is "SELECT 1"
    validation-query = ""
  }

  # If database timestamp is guaranteed to not move backwards for two subsequent
  # updates of the same persistenceId there might be a performance gain to
  # set this to `on`. Note that many databases use the system clock and that can
  # move backwards when the system clock is adjusted.
  db-timestamp-monotonic-increasing = off

  # Enable this for testing or workaround of https://github.com/yugabyte/yugabyte-db/issues/10995
  # FIXME: This property will be removed when the Yugabyte issue has been resolved.
  use-app-timestamp = off

  # Logs database calls that take longer than this duration at INFO level.
  # Set to "off" to disable this logging.
  # Set to 0 to log all calls.
  log-db-calls-exceeding = 300 ms

}

3. PostgreSQL建表脚本

CREATE TABLE IF NOT EXISTS public.event_journal(
                                                   ordering BIGSERIAL,
                                                   persistence_id VARCHAR(255) NOT NULL,
    sequence_number BIGINT NOT NULL,
    deleted BOOLEAN DEFAULT FALSE NOT NULL,

    writer VARCHAR(255) NOT NULL,
    write_timestamp BIGINT,
    adapter_manifest VARCHAR(255),

    event_ser_id INTEGER NOT NULL,
    event_ser_manifest VARCHAR(255) NOT NULL,
    event_payload BYTEA NOT NULL,

    meta_ser_id INTEGER,
    meta_ser_manifest VARCHAR(255),
    meta_payload BYTEA,

    PRIMARY KEY(persistence_id, sequence_number)
    );

CREATE UNIQUE INDEX event_journal_ordering_idx ON public.event_journal(ordering);

CREATE TABLE IF NOT EXISTS public.event_tag(
                                               event_id BIGINT,
                                               tag VARCHAR(256),
    PRIMARY KEY(event_id, tag),
    CONSTRAINT fk_event_journal
    FOREIGN KEY(event_id)
    REFERENCES event_journal(ordering)
    ON DELETE CASCADE
    );

CREATE TABLE IF NOT EXISTS public.snapshot (
    persistence_id VARCHAR(255) NOT NULL,
    sequence_number BIGINT NOT NULL,
    created BIGINT NOT NULL,

    snapshot_ser_id INTEGER NOT NULL,
    snapshot_ser_manifest VARCHAR(255) NOT NULL,
    snapshot_payload BYTEA NOT NULL,

    meta_ser_id INTEGER,
    meta_ser_manifest VARCHAR(255),
    meta_payload BYTEA,

    PRIMARY KEY(persistence_id, sequence_number)
    );

CREATE TABLE IF NOT EXISTS public.durable_state (
                                                    global_offset BIGSERIAL,
                                                    persistence_id VARCHAR(255) NOT NULL,
    revision BIGINT NOT NULL,
    state_payload BYTEA NOT NULL,
    state_serial_id INTEGER NOT NULL,
    state_serial_manifest VARCHAR(255),
    tag VARCHAR,
    state_timestamp BIGINT NOT NULL,
    PRIMARY KEY(persistence_id)
    );
CREATE INDEX CONCURRENTLY state_tag_idx on public.durable_state (tag);
CREATE INDEX CONCURRENTLY state_global_offset_idx on public.durable_state (global_offset);

4. 错误信息

[error] c.k.a.c.CompanyActorTest - Supervisor StopSupervisor saw failure: Exception during recovery. PersistenceId [Company]. Relation «durable_state» existiert nicht
akka.persistence.typed.state.internal.DurableStateStoreException: Exception during recovery. PersistenceId [Company]. Relation «durable_state» existiert nicht
        at akka.persistence.typed.state.internal.Recovering.onRecoveryFailure(Recovering.scala:125)
        at akka.persistence.typed.state.internal.Recovering.onGetFailure(Recovering.scala:188)
        at akka.persistence.typed.state.internal.Recovering.onMessage(Recovering.scala:78)
        at akka.persistence.typed.state.internal.Recovering.onMessage(Recovering.scala:61)
        at akka.actor.typed.scaladsl.AbstractBehavior.receive(AbstractBehavior.scala:84)
        at akka.actor.typed.Behavior$.interpret(Behavior.scala:281)
        at akka.actor.typed.Behavior$.interpretMessage(Behavior.scala:237)
        at akka.actor.typed.internal.InterceptorImpl$$anon$2.apply(InterceptorImpl.scala:57)
        at akka.persistence.typed.state.internal.DurableStateBehaviorImpl$$anon$1.aroundReceive(DurableStateBehaviorImpl.scala:125)
        at akka.actor.typed.internal.InterceptorImpl.receive(InterceptorImpl.scala:85)
Caused by: io.r2dbc.postgresql.ExceptionFactory$PostgresqlBadGrammarException: Relation «durable_state» existiert nicht
        at io.r2dbc.postgresql.ExceptionFactory.createException(ExceptionFactory.java:96)
        at io.r2dbc.postgresql.ExceptionFactory.createException(ExceptionFactory.java:65)
        at io.r2dbc.postgresql.ExceptionFactory.handleErrorResponse(ExceptionFactory.java:132)
        at reactor.core.publisher.FluxHandleFuseable$HandleFuseableSubscriber.onNext(FluxHandleFuseable.java:176)
        at reactor.core.publisher.FluxFilterFuseable$FilterFuseableConditionalSubscriber.onNext(FluxFilterFuseable.java:337)
        at reactor.core.publisher.FluxContextWrite$ContextWriteSubscriber.onNext(FluxContextWrite.java:107)
        at reactor.core.publisher.FluxPeekFuseable$PeekConditionalSubscriber.onNext(FluxPeekFuseable.java:854)
        at reactor.core.publisher.FluxPeekFuseable$PeekConditionalSubscriber.onNext(FluxPeekFuseable.java:854)
        at io.r2dbc.postgresql.util.FluxDiscardOnCancel$FluxDiscardOnCancelSubscriber.onNext(FluxDiscardOnCancel.java:91)
        at reactor.core.publisher.FluxDoFinally$DoFinallySubscriber.onNext(FluxDoFinally.java:113)

(注:错误信息中的德语“existiert nicht”意为“不存在”)

5. 当前依赖配置

libraryDependencies += "com.typesafe.akka" %% "akka-serialization-jackson" % "2.8.0"
libraryDependencies += "com.lightbend.akka" %% "akka-persistence-r2dbc" % "1.0.0"
libraryDependencies += "com.github.dnvriend" %% "akka-persistence-jdbc" % "3.5.3"

调整版本时出现依赖冲突。


解决方案

一、解决durable_state表不存在问题

  1. 验证表的创建状态

    • 登录PostgreSQL数据库,执行\dt public.*命令,确认durable_state表是否存在。若不存在,重新执行建表脚本,注意脚本换行是否导致语句截断,确保完整执行。
    • 确认配置中的schema参数为空(默认使用public schema),与建表脚本中的public.durable_state匹配。
  2. 修复Actor状态修改的线程安全问题

    • 当前代码中直接修改state.companies(state.companies.add(command.company)),但Collections.emptySet()返回的是不可变集合,调用add会抛出UnsupportedOperationException,导致状态无法持久化。修改为创建新的不可变集合:
      .onCommand(AddCompany.class, (state, command) -> {
          Set<Company> newCompanies = new HashSet<>(state.getCompanies());
          newCompanies.add(command.company);
          return Effect().persist(new CompanyState(Collections.unmodifiableSet(newCompanies)));
      })
      
    • 同时修改CompanyState,确保内部集合不可变:
      public class CompanyState {
          private final Set<Company> companies;
      
          public CompanyState(Set<Company> companies) {
              this.companies = Collections.unmodifiableSet(new HashSet<>(companies));
          }
      
          public Set<Company> getCompanies() {
              return companies;
          }
      }
      

二、解决依赖版本冲突问题

  1. 移除不必要的依赖

    • 仅使用Durable State特性,无需akka-persistence-jdbc,直接删除该依赖:
      // 移除此行
      // libraryDependencies += "com.github.dnvriend" %% "akka-persistence-jdbc" % "3.5.3"
      
  2. 对齐Akka与Play版本

    • Play 2.8.18默认依赖Akka 2.6.20,当前使用的akka-serialization-jackson 2.8.0版本不兼容,需调整为对应版本:
      libraryDependencies += "com.typesafe.akka" %% "akka-serialization-jackson" % "2.6.20"
      libraryDependencies += "com.lightbend.akka" %% "akka-persistence-r2dbc" % "1.0.0"
      
  3. 添加R2DBC驱动依赖

    • 补充PostgreSQL的R2DBC驱动,确保数据库连接正常:
      libraryDependencies += "io.r2dbc" % "r2dbc-postgresql" % "0.8.13.RELEASE"
      

三、额外配置检查

  • 确认akka.persistence.state.plugin配置未被其他设置覆盖,确保生效的是akka.persistence.r2dbc.state。
  • 验证Docker容器的PostgreSQL端口映射(5432端口)、数据库名称、用户名、密码与配置完全匹配。

内容的提问来源于stack exchange,提问作者Kritiko

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.18 14:06:59