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

Apache Pekko + R2DBC + PostgreSQL 间歇性更新0行问题:记录已插入但UPDATE影响行数为0

Apache Pekko + R2DBC + PostgreSQL 间歇性更新0行问题:记录已插入但UPDATE影响行数为0

我之前在处理Pekko + R2DBC的分布式项目时,也碰到过类似的间歇性数据一致性问题,结合你的场景和代码来看,这个问题大概率和R2DBC多语句执行的结果集处理、事务边界控制以及Pekko Projection的事件时序配置有关,下面一步步拆解分析:

问题核心场景复盘

在低负载、无并发的环境下:

  1. Campaign记录成功插入数据库(日志确认插入后500-900ms触发后续操作)
  2. Pekko Projection处理MessageSent事件时,执行SELECT查询Campaign ID返回空
  3. 后续UPDATE语句因WHERE条件不匹配,影响行数为0
  4. 该问题非必现,偶尔能正常执行更新

环境配置梳理

你的基础配置如下:

  • 数据库:AWS RDS PostgreSQL
  • 驱动:r2dbc-postgresql
  • 框架:Apache Pekko + pekko-persistence-r2dbc
  • 关键R2DBC连接池配置:
pekko.persistence.r2dbc {
  connection-factory {
    initial-size = 100
    max-size = 100
    connect-timeout = 5 seconds
    acquire-timeout = 10 seconds
    acquire-retry = 1
    max-idle-time = 10 minutes
  }
  query {
    behind-current-time = 5000 millis
    refresh-interval = 20s
  }
}

代码问题深度分析

看你的CampaignTableStatisticsUpdaterProjectionHandler代码,有几个容易被忽略的问题:

1. 多语句执行的结果集丢失

你在SELECT语句前加了SET TRANSACTION ISOLATION LEVEL READ COMMITTED;,并将两个语句放在同一个createStatement中用分号分隔:

var selectStatement = session
        .createStatement("SET TRANSACTION ISOLATION LEVEL READ COMMITTED; SELECT id FROM public.campaigns WHERE id = $1;")
        .bind(0, evt.getCampaignId());

但R2DBC的PostgreSQL驱动默认只会返回第一个语句的结果集,而SET TRANSACTION语句没有返回结果,导致后续session.selectOne()拿到的是空值,即使后面的SELECT语句能查到数据,也会被判定为“记录不存在”,最终跳过UPDATE。

而且PostgreSQL默认的事务隔离级别就是READ COMMITTED,这个SET语句完全是多余的。

2. 事务边界未显式控制

你的SELECT和UPDATE操作在同一个Session中,但未显式开启事务,默认会以自动提交模式执行(每个语句都是独立事务)。虽然低负载下不会有并发干扰,但这种分离的事务无法保证读取和更新的原子性,极端情况下可能出现数据不一致。

3. Projection时序配置不匹配

你设置了behind-current-time = 5000 millis,这个配置是让Pekko Projection延迟5秒处理事件,避免读取数据库中尚未完全持久化的新数据。但你的Campaign插入后仅500-900ms就触发了Projection,这时候Projection可能还在等待“落后当前时间5秒”的条件,导致事件处理时,数据库的新数据可见性还没完全同步(AWS RDS PostgreSQL存在轻微的提交延迟)。

针对性解决方案

1. 移除多余的事务级别设置语句

直接删除SET TRANSACTION ISOLATION LEVEL READ COMMITTED;,仅保留SELECT语句,确保selectOne()能正确拿到查询结果:

var selectStatement = session
        .createStatement("SELECT id FROM public.campaigns WHERE id = $1;")
        .bind(0, evt.getCampaignId());

2. 显式开启事务包裹查询与更新

使用session.withTransaction()将SELECT和UPDATE包裹在同一个事务中,保证读取和更新的原子性,同时避免自动提交模式下的事务分离问题:

return session.withTransaction(txSession -> 
    txSession.selectOne(selectStatement, row -> row.get("id"))
        .thenCompose(row -> {
            log.info("Campaign exists with id {}", row);
            if (row == null) {
                return CompletableFuture.completedFuture(Done.done());
            }

            var updateStatement = txSession
                    .createStatement("UPDATE public.campaigns SET message_sent_count = message_sent_count + 1, total_cost = total_cost + $1 WHERE id = $2;")
                    .bind(0, evt.getMessageCost())
                    .bind(1, evt.getCampaignId());

            return txSession.updateOne(updateStatement)
                    .thenApply(rowsUpdated -> {
                        log.info("Successfully updated campaign {}, rows affected: {}", evt.getCampaignId(), rowsUpdated);
                        return Done.done();
                    });
        })
);

3. 调整Projection的时序配置

将behind-current-time调整为与你的插入延迟匹配的数值,比如1000ms,既保证数据完全可见,又不会过度延迟事件处理:

pekko.persistence.r2dbc {
  query {
    behind-current-time = 1000 millis
    refresh-interval = 20s
  }
}

4. 优化连接池配置

你的连接池初始大小和最大大小都设为100,在低负载环境下会造成大量空闲连接,可能导致连接状态残留。建议调小至更合理的数值:

pekko.persistence.r2dbc {
  connection-factory {
    initial-size = 5
    max-size = 20
    connect-timeout = 5 seconds
    acquire-timeout = 10 seconds
    acquire-retry = 1
    max-idle-time = 10 minutes
  }
}

5. 额外排查建议

  • 确认Campaign表的id字段已设置为主键索引,避免因查询性能问题导致的延迟
  • 开启PostgreSQL的慢查询日志,记录查询和更新语句的执行时间,排查是否存在数据库端的延迟
  • 在日志中加入数据库当前时间与事件时间戳的对比,验证Projection的时序配置是否合理

修改后的完整代码

public class CampaignTableStatisticsUpdaterProjectionHandler extends R2dbcHandler<EventEnvelope<OutboundCampaignMessageEntity.Event>> {

    @Override
    public CompletionStage<Done> process(R2dbcSession session, EventEnvelope<OutboundCampaignMessageEntity.Event> eventEventEnvelope) throws Exception {
        if (eventEventEnvelope.event().getTimestamp().isBefore(ServiceUtils.BEFORE_INCIDENT)) {
            return CompletableFuture.completedFuture(Done.done());
        }

        log.debug("Processing event: {}, sequence: {}", eventEventEnvelope.event(), eventEventEnvelope.sequenceNr());

        if (eventEventEnvelope.event() instanceof OutboundCampaignMessageEntity.MessageSent evt) {
            var selectStatement = session
                    .createStatement("SELECT id FROM public.campaigns WHERE id = $1;")
                    .bind(0, evt.getCampaignId());

            // 显式开启事务,保证查询与更新的原子性
            return session.withTransaction(txSession -> 
                txSession.selectOne(selectStatement, row -> row.get("id"))
                    .thenCompose(row -> {
                        log.info("Campaign exists with id {}", row);
                        if (row == null) {
                            log.warn("Failed to find campaign with id: {}", evt.getCampaignId());
                            return CompletableFuture.completedFuture(Done.done());
                        }

                        var updateStatement = txSession
                                .createStatement("UPDATE public.campaigns SET message_sent_count = message_sent_count + 1, total_cost = total_cost + $1 WHERE id = $2;")
                                .bind(0, evt.getMessageCost())
                                .bind(1, evt.getCampaignId());

                        return txSession.updateOne(updateStatement)
                                .thenApply(rowsUpdated -> {
                                    log.info("Successfully updated campaign {}, rows affected: {}", evt.getCampaignId(), rowsUpdated);
                                    return Done.done();
                                });
                    })
            );
        } else if (eventEventEnvelope.event() instanceof OutboundCampaignMessageEntity.MessageCreated evt) {
            var statement = session
                    .createStatement("UPDATE public.campaigns SET message_created_count = message_created_count + 1 WHERE id = $1;")
                    .bind(0, evt.getCampaignId());

            // 同样建议用事务包裹MessageCreated的更新
            return session.withTransaction(txSession -> 
                txSession.updateOne(statement).thenApply(rowsUpdated -> Done.done())
            );
        }

        return CompletableFuture.completedFuture(Done.done());
    }
}

内容来源于stack exchange

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.08 10:19:32