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

如何用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)字段,通过状态标记实现锁控制:

实现步骤

  1. Worker在事务中更新10个未锁定、未处理的行,设置locked_at为当前时间
  2. 查询这些被锁定的行获取ID
  3. 处理成功则删除/标记为已处理;失败则重置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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 11:57:44