TypeORM与PostgreSQL中悲观写锁结合聚合函数报错的解决
问题:TypeORM结合PostgreSQL使用悲观写锁时出现聚合函数与FOR UPDATE冲突错误
我是TypeORM和PostgreSQL的新手,尝试使用typeorm select and update lock(悲观写锁)实现并发控制,但运行代码时出现错误。
我的代码:
const queryRunner = this.dataSource.createQueryRunner(); // 初始化事务 await queryRunner.connect(); await queryRunner.startTransaction(); // 计算总积分 const pointData = await queryRunner.manager .getRepository(PointsLedger) .createQueryBuilder('points_ledger') .useTransaction(true) .setLock("pessimistic_write") .where(`points_ledger.customerId = :customerId`, { customerId }) .select('SUM(points_ledger.credit) - SUM(points_ledger.debit)', 'points') .getRawOne(); // 业务逻辑判断 const pointsLedgerData = { customerId, debit, credit }; // 尝试保存积分记录,但执行失败,错误信息见下文 const savePoints = await queryRunner.manager .getRepository(PointsLedger) .save(pointsLedgerData); // 获取积分汇总记录 const pointsSummaryData = await queryRunner.manager .getRepository(PointsSummary) .createQueryBuilder('pointsSummary') .where(`pointsSummary.customerId = :customerId`, { customerId }) .getRawOne(); // 如果积分汇总记录为空则新增,否则更新 // 更新积分汇总 const response = await queryRunner.manager .getRepository(PointsSummary) .update( { customerId }, { totalPoints: pointData.points + debit } );
错误信息:
QueryFailedError: FOR UPDATE is not allowed with aggregate functions at PostgresQueryRunner.query (project/src/driver/postgres/PostgresQueryRunner.ts:299:19) at processTicksAndRejections (node:internal/process/task_queues:96:5) at SelectQueryBuilder.loadRawResults (project/src/query-builder/SelectQueryBuilder.ts:3555:25) at SelectQueryBuilder.getRawMany (project/src/query-builder/SelectQueryBuilder.ts:1553:29) at SelectQueryBuilder.getRawOne (project/src/query-builder/SelectQueryBuilder.ts:1530:17) at LedgerService.getPoints (project/src/server/app/ledger/ledger.service.ts:226:20) at LedgerService.debit (project/src/server/app/ledger/ledger.service.ts:91:27) at LedgerController.burnPoints (project/src/server/app/ledger/ledger.controller.ts:60:22) { query: 'SELECT SUM("points_ledger"."credit") - SUM("points_ledger"."debit") AS "points" FROM "points_ledger" "points_ledger" WHERE "points_ledger"."customer_id" = $1 FOR UPDATE', parameters: [ 'some-customer-id' ], driverError: error: FOR UPDATE is not allowed with aggregate functions at Parser.parseErrorMessage (project/node_modules/pg-protocol/src/parser.ts:369:69) at Parser.handlePacket (project/node_modules/pg-protocol/src/parser.ts:188:21) at Parser.parse (project/node_modules/pg-protocol/src/parser.ts:103:30) at Socket.<anonymous> (project/node_modules/pg-protocol/src/index.ts:7:48) at Socket.emit (node:events:527:28) at addChunk (node:internal/streams/readable:324:12) at readableAddChunk (node:internal/streams/readable:297:9) at Socket.Readable.push (node:internal/streams/readable:234:10) at TCP.onStreamRead (node:internal/stream_base_commons:190:23) at TCP.callbackTrampoline (node:internal/async_hooks:130:17) }
解决方案
核心原因
PostgreSQL的FOR UPDATE锁仅对查询返回的原始表行生效,而聚合函数(如SUM)返回的是计算后的聚合结果,不是表中的实际行,因此无法对聚合结果加锁,导致报错。
方法1:先锁定用户的所有ledger记录,再计算总和
先执行一次不带聚合的查询,锁定该用户的所有PointsLedger记录,之后再计算总和,这样就能保证并发操作时数据的一致性:
const queryRunner = this.dataSource.createQueryRunner(); await queryRunner.connect(); await queryRunner.startTransaction(); // 第一步:锁定该用户的所有points_ledger记录,阻止并发修改 await queryRunner.manager .getRepository(PointsLedger) .createQueryBuilder('points_ledger') .setLock("pessimistic_write") .where(`points_ledger.customerId = :customerId`, { customerId }) .getMany(); // 执行查询并加锁 // 第二步:计算总积分,此时数据已被锁定,不会出现并发修改 const pointData = await queryRunner.manager .getRepository(PointsLedger) .createQueryBuilder('points_ledger') .where(`points_ledger.customerId = :customerId`, { customerId }) .select('SUM(points_ledger.credit) - SUM(points_ledger.debit)', 'points') .getRawOne(); // 后续业务逻辑保持不变 const pointsLedgerData = { customerId, debit, credit }; const savePoints = await queryRunner.manager .getRepository(PointsLedger) .save(pointsLedgerData); const pointsSummaryData = await queryRunner.manager .getRepository(PointsSummary) .createQueryBuilder('pointsSummary') .where(`pointsSummary.customerId = :customerId`, { customerId }) .getRawOne(); // 根据积分汇总记录是否为空执行新增或更新 if (!pointsSummaryData) { await queryRunner.manager.getRepository(PointsSummary).save({ customerId, totalPoints: pointData.points + debit }); } else { await queryRunner.manager .getRepository(PointsSummary) .update({ customerId }, { totalPoints: pointData.points + debit }); } // 提交事务(出错时需回滚) await queryRunner.commitTransaction();
方法2:使用嵌套查询锁定表行(多表场景适配)
如果查询涉及多张表,可以通过嵌套查询先锁定原始行,再执行聚合计算,单表场景下也适用:
const pointData = await queryRunner.manager .getRepository(PointsLedger) .createQueryBuilder('pl') .where(`pl.customerId = :customerId`, { customerId }) .setLock("pessimistic_write") .select((subQuery) => { return subQuery .select('SUM(pl2.credit) - SUM(pl2.debit)', 'points') .from(PointsLedger, 'pl2') .where('pl2.customerId = :customerId', { customerId }); }, 'points') .getRawOne();
方法3:改用乐观锁(可选)
如果业务场景中并发冲突概率较低,可以用乐观锁替代悲观锁:
- 在
PointsLedger或PointsSummary表中添加version字段(整数类型,默认值1) - 更新时将
version作为条件,每次更新后版本号自增 - TypeORM可通过
@Version()装饰器自动管理版本
示例:
// PointsSummary实体添加版本字段 @Entity() export class PointsSummary { // 其他字段... @Column() totalPoints: number; @Version() version: number; } // 更新逻辑 const response = await queryRunner.manager .getRepository(PointsSummary) .update( { customerId, version: pointsSummaryData.version }, { totalPoints: pointData.points + debit, version: pointsSummaryData.version + 1 } ); // 如果更新结果的affected为0,说明存在并发修改,需重试或返回错误
内容的提问来源于stack exchange,提问作者sabbir
相关产品推荐
相关产品推荐

