如何用Postgres/Prisma实现队列,避免Worker重复选取已处理行
解决方案
1. 基于Postgres行级锁的原子选取(推荐)
利用Postgres的SELECT ... FOR UPDATE SKIP LOCKED特性,实现Worker间无冲突获取未处理ID,无需高隔离级别,也不会锁定整张表。
核心逻辑
在单个事务中,原子性选取10个未处理ID并锁定对应行,自动跳过已被其他Worker锁定的行。处理成功后删除(或标记已处理),失败则回滚事务释放锁,让其他Worker可以重新选取。
Prisma实现示例
由于Prisma ORM暂不直接支持FOR UPDATE SKIP LOCKED,可通过原生SQL实现:
import { PrismaClient } from '@prisma/client'; const prisma = new PrismaClient(); async function processBatch() { return prisma.$transaction(async (tx) => { // 1. 原子获取并锁定10个降序排列的未处理ID const batch = await tx.$queryRaw` SELECT id FROM queue ORDER BY id DESC LIMIT 10 FOR UPDATE SKIP LOCKED; `; if (batch.length === 0) return null; const ids = batch.map((item: any) => item.id); // 2. 替换为你的业务处理逻辑 await handleBusinessLogic(ids); // 3. 处理成功后删除ID await tx.queue.deleteMany({ where: { id: { in: ids } } }); return ids; }); }
优势
- 彻底避免重复选取:
SKIP LOCKED自动过滤已锁定行,每个Worker拿到唯一的未处理ID - 低冲突风险:无需
SERIALIZED隔离级别,默认READ COMMITTED即可,不会触发写入冲突报错 - 完美支持降序:直接通过
ORDER BY id DESC实现排序需求 - 适配超大规模数据:Postgres可轻松承载10亿级表,只要给
id列加主键/索引,就能保证查询性能
2. 基于状态标记的锁机制
若不想用原生SQL,可给queue表添加locked_at(timestamp)和processed(boolean)字段,通过状态标记实现锁控制:
实现步骤
- Worker在事务中更新10个未锁定、未处理的行,设置
locked_at为当前时间 - 查询这些被锁定的行获取ID
- 处理成功则删除/标记为已处理;失败则重置
locked_at为null,释放锁
Prisma代码示例
async function fetchLockedBatch() { const lockTime = new Date(); // 锁定10个未处理、未锁定的降序ID const updateResult = await prisma.queue.updateMany({ where: { processed: false, locked_at: null }, data: { locked_at: lockTime }, orderBy: { id: 'desc' }, take: 10, }); if (updateResult.count === 0) return []; // 查询刚锁定的ID const batch = await prisma.queue.findMany({ where: { processed: false, locked_at: lockTime }, }); return batch.map(item => item.id); } // 处理完成后清理 async function markBatchAsProcessed(ids: number[]) { await prisma.queue.deleteMany({ where: { id: { in: ids } } }); } // 定时解锁超时的任务(避免Worker崩溃导致锁永久持有) async function unlockStaleTasks() { await prisma.queue.updateMany({ where: { processed: false, locked_at: { lt: new Date(Date.now() - 5 * 60 * 1000) }, // 解锁5分钟前的锁定任务 }, data: { locked_at: null }, }); }
3. 替代消息队列方案
若仍考虑消息队列,可选择以下适配需求的方案:
- Postgres原生消息队列:结合
LISTEN/NOTIFY和上述行锁方案,新ID入队时发送通知唤醒Worker,减少轮询开销(也可保留现有2秒轮询逻辑) - Apache Kafka:支持持久化磁盘存储(无需内存),可通过分区配置实现近似降序消费,适合超大规模数据场景,但需额外部署维护
关键注意事项
- 必须给
queue表的id列加主键或索引,否则ORDER BY id DESC LIMIT 10会触发全表扫描,性能极差 - 缩短事务时长:业务处理逻辑尽量精简,避免长时间持有行锁影响并发效率
- 定期清理:对已处理的ID及时删除或归档到历史表,保持主表的查询性能
内容的提问来源于stack exchange,提问作者Dominic Sore
相关产品推荐
相关产品推荐

