Apache Pekko + R2DBC + PostgreSQL 间歇性更新0行问题:记录已插入但UPDATE影响行数为0
我之前在处理Pekko + R2DBC的分布式项目时,也碰到过类似的间歇性数据一致性问题,结合你的场景和代码来看,这个问题大概率和R2DBC多语句执行的结果集处理、事务边界控制以及Pekko Projection的事件时序配置有关,下面一步步拆解分析:
问题核心场景复盘
在低负载、无并发的环境下:
- Campaign记录成功插入数据库(日志确认插入后500-900ms触发后续操作)
- Pekko Projection处理
MessageSent事件时,执行SELECT查询Campaign ID返回空 - 后续UPDATE语句因WHERE条件不匹配,影响行数为0
- 该问题非必现,偶尔能正常执行更新
环境配置梳理
你的基础配置如下:
- 数据库: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

