PostgreSQL的SKIP LOCKED如何实现高性能任务队列?

来源:网络编程作者:天马头衔:网络博主
导读:本期聚焦于天马创作的《PostgreSQL的SKIP LOCKED如何实现高性能任务队列?》,敬请观看详情。多个进程同时从数据库抢单子时,你有没有遇到过任务被重复消费的问题?传统的先查询再更新写法在并发场景下会产生竞态条件,导致同一条记录被多个worker处理。PostgreSQL从9.5版本开始提供的SKIP LOCKED特性可以优雅地解决这个问题:配合FOR UPDATE跳过已被其他事务锁定的行,让每个消费者各自拿到不冲突的任务。本文将深入讲解SKIP LOCKED的底层锁机制,对比悲观锁、乐观锁、 advisory lock等常见方案的优劣,给出完整的任务队列表结构与消费SQL实现,并附上PHP和Go两种语言的实战代码,最后分析索引设计、批量取任务、失败重试等生产环境必须注意的细节。

在构建后台任务系统时,很多团队会直接用数据库表来做队列:一张tasks表,插入任务,然后由若干个worker进程不断轮询取任务执行。这种方案简单直接,不需要引入Redis或RabbitMQ等额外组件,在中小规模场景下非常实用。但一旦worker数量增多,最常见的问题就出现了:两个worker同时查到同一条任务,都去执行,任务被重复消费。要解决这个问题,PostgreSQL提供了一个专门为此而生的特性——SKIP LOCKED。

PostgreSQL的SKIP LOCKED如何实现高性能任务队列?

为什么传统写法会出现重复消费

先看一段典型的有问题的代码逻辑。worker先查询出一条待处理任务,然后把它标记为处理中,中间存在一个时间窗口:

-- worker A 和 worker B 几乎同时执行
SELECT id FROM tasks WHERE status = 'pending' LIMIT 1;
-- 两者都拿到了 id = 100 的任务
UPDATE tasks SET status = 'running' WHERE id = 100;
-- 结果:两个worker都认为自己抢到了任务100,重复执行

这个竞态条件的根源在于SELECT和UPDATE是两条独立的语句,在两个语句之间没有任何锁保护。即使把两条语句放在同一个事务里也无济于事,因为普通的SELECT不加任何行锁,两个事务都能读到同一行。

有人会想到用SELECT FOR UPDATE来解决,这个思路是对的:FOR UPDATE会对选中的行加行级排他锁,其他事务再执行FOR UPDATE时就会阻塞等待。但随之而来的新问题是性能——当有10个worker同时抢单时,9个worker会卡在锁等待上,等第一个事务提交后它们再醒来,又都去争抢同一条记录,造成大量的锁竞争和无效等待,吞吐量反而上不去。

SKIP LOCKED正是为了这个场景设计的。它的语义是:执行SELECT FOR UPDATE时,如果目标行已经被其他事务锁定,就跳过这一行,而不是阻塞等待。这样一来,10个worker并发抢单时,各自锁定不同的行,互不干扰,天然实现了任务的分发的原子性。

Skip LOCKED的锁机制原理

PostgreSQL的行锁是通过对行的物理版本(元组)打标记实现的。当事务A执行SELECT ... FOR UPDATE并选中某行时,会在该元组上记录锁定者信息,其他事务对这一行再加锁请求时,会检测到冲突。默认行为是等待,而指定SKIP LOCKED后,加锁函数会立即返回失败,查询规划器随即把这一行从结果集中剔除,继续处理后续的行。

需要理解的一个关键点是:SKIP LOCKED跳过的只是被锁的行,而不是整个查询。假设队列里有100条pending任务,worker A锁住了前10条,worker B执行同样的查询时,会自动从第11条开始返回,这就是多个worker能高效瓜分任务队列的原因。整个加锁、跳过、返回的过程在数据库引擎内部完成,客户端感知不到任何中间状态。

还要注意事务边界。FOR UPDATE加的行锁会一直持有到事务结束(提交或回滚)才释放。因此取任务的完整流程应该是:开启事务,用FOR UPDATE SKIP LOCKED取出任务并立即UPDATE其状态,然后提交事务。提交之后锁释放,但状态已经改成running,其他worker的查询条件里status = 'pending'自然不会再命中这条记录。

完整的表结构与核心SQL实现

先设计一张生产级可用的任务队列表,包含状态、重试计数、执行时间等字段:

CREATE TABLE tasks (
    id          BIGSERIAL PRIMARY KEY,
    queue       TEXT NOT NULL DEFAULT 'default',   -- 队列名称,支持多队列
    payload     JSONB NOT NULL,                    -- 任务内容
    status      SMALLINT NOT NULL DEFAULT 0,       -- 0待处理 1处理中 2完成 3失败
    run_at      TIMESTAMPTZ NOT NULL DEFAULT now(), -- 延迟执行时间
    attempts    INT NOT NULL DEFAULT 0,            -- 已重试次数
    created_at  TIMESTAMPTZ NOT NULL DEFAULT now()
);

-- 关键索引:按队列+状态查询,按时间排序
CREATE INDEX idx_tasks_poll ON tasks (queue, status, run_at)
    WHERE status = 0;

这里用了部分索引,只索引status = 0的待处理任务,因为队列表中绝大多数数据是已完成状态,部分索引能大幅减小索引体积,提升查询速度。核心的取任务SQL用一个CTE配合data-modifying语句实现:

WITH picked AS (
    SELECT id FROM tasks
    WHERE queue = 'default'
      AND status = 0
      AND run_at <= now()
    ORDER BY run_at
    LIMIT 10
    FOR UPDATE SKIP LOCKED
)
UPDATE tasks t
SET status = 1, attempts = attempts + 1
FROM picked
WHERE t.id = picked.id
RETURNING t.id, t.payload;

这个SQL的巧妙之处在于:CTE内部的SELECT负责加锁并挑选任务,外层的UPDATE负责变更状态,最后通过RETURNING把任务内容直接返回给客户端,三步在一次数据库交互中原子完成。LIMIT 10表示一次批量取10条,减少轮询频率。每次被取出的任务attempts加一,为后续的重试判断提供依据。

ORDER BY在这里也很重要。不加排序的话,PG返回的行顺序不确定,多个worker可能倾向于争抢同一批物理位置靠前的行,加上run_at排序后配合索引,扫描顺序稳定且可预测。

PHP与Go语言的消费端实战代码

以PHP为例,使用PDO实现一个worker循环。注意要设置search_path或显式schema,并正确处理异常回滚:

<?php
$pdo = new PDO('pgsql:host=127.0.0.1;dbname=app', 'app', 'secret', [
    PDO::ATTR_ERRMODE => PDO::ERRMODE_EXCEPTION,
]);

while (true) {
    $pdo->beginTransaction();
    try {
        $sql = "WITH picked AS (
                    SELECT id FROM tasks
                    WHERE status = 0 AND run_at <= now()
                    ORDER BY run_at LIMIT 5
                    FOR UPDATE SKIP LOCKED
                )
                UPDATE tasks t SET status = 1
                FROM picked WHERE t.id = picked.id
                RETURNING t.id, t.payload";
        $rows = $pdo->query($sql)->fetchAll(PDO::FETCH_ASSOC);

        if (!$rows) {
            $pdo->rollBack();
            usleep(500000); // 队列空了,歇半秒再轮询
            continue;
        }
        foreach ($rows as $task) {
            processTask(json_decode($task['payload'], true));
        }
        $pdo->commit(); // 提交后任务标记为处理中
    } catch (Exception $e) {
        $pdo->rollBack();
        // 回滚后attempts不会增加,任务回到pending,可重试
    }
}

Go的实现思路类似,使用database/sql标准库:

rows, err := db.QueryContext(ctx, `
    WITH picked AS (
        SELECT id FROM tasks
        WHERE status = 0 AND run_at <= now()
        ORDER BY run_at LIMIT 5
        FOR UPDATE SKIP LOCKED
    )
    UPDATE tasks t SET status = 1
    FROM picked WHERE t.id = picked.id
    RETURNING t.id, t.payload`)
if err != nil {
    log.Fatal(err)
}
defer rows.Close()

for rows.Next() {
    var id int64
    var payload []byte
    if err := rows.Scan(&id, &payload); err != nil {
        log.Fatal(err)
    }
    go handleTask(id, payload)
}

两种语言的差别主要在事务管理上:PHP脚本式执行每次循环显式开关事务即可,Go作为常驻服务,建议把取任务和执行解耦,取到任务后提交事务,再异步处理,避免长事务持有行锁时间过久。

生产环境必须注意的细节

第一是失败重试机制。任务处理失败后不能永远停留在running状态,通常有两种做法:处理失败时显式将status改回0并设置run_at = now() + interval '1 minute' * attempts实现指数退避;同时配合一个超时回收任务,定期把处理时间超过阈值的running任务重置回pending,防止worker进程崩溃导致的任务悬挂。

第二是长事务陷阱。如果在取到任务后的事务里执行很耗时的业务逻辑,行锁会一直被持有,极端情况下会把后续worker的SELECT变得低效——虽然SKIP LOCKED不会阻塞,但大量被锁住的行会占满LIMIT配额之前的扫描位置,增加无效扫描量。正确姿势是取任务、更新状态、提交这个动作越快越好,业务处理放在事务外,处理完成后再单独UPDATE为最终状态。

第三是监控表膨胀。队列表的写入和更新非常频繁,死元组增长快,要保证autovacuum正常工作,必要时对该表单独设置更激进的autovacuum参数,例如把autovacuum_vacuum_scale_factor调低到0.01,并定期归档或清理已完成的历史任务,避免表无限膨胀拖慢查询。

综合来看,SKIP LOCKED方案在worker数量几十个以内、任务量每秒数千级别的场景下表现优秀,部署简单,不依赖外部中间件,还能借助SQL的强大表达力实现延迟任务、优先级队列、按队列隔离等特性。当规模进一步扩大时再考虑迁移到专业消息队列也不迟,而且表结构设计得好的话,迁移成本也相对可控。

PostgreSQLSKIP LOCKED任务队列修改时间:2026-09-08 12:47:10

免责声明:已尽一切努力确保本网站所含信息的准确性。网站作品多为原创整理与精心创作,观点力求客观中立。本站旨在免费分享,内容仅供个人学习、研究或参考使用。若引用了第三方作品,版权归原作者所有。如内容涉及您的权益,请联系我们进行处理Email:chomcom@qq.com。
引用或转载本作品时,请注明当前出处:https://www.ipipp.com/html/20260908/52774.html,基于非商业用途的前提下,欢迎转载或二创本作品。
内容垂直聚焦
专注技术核心技术栏目,确保每篇文章深度聚焦于实用技能。从代码技巧到架构设计,为用户提供无干扰的纯技术知识沉淀,精准满足专业提升需求。
知识结构清晰
覆盖从开发到部署的全链路。AI、前端、编程、数据库、服务器、建站、系统层层递进,构建清晰学习路径,帮助用户系统化掌握开发与运维所需的核心技术。
深度技术解析
拒绝泛泛而谈,深入技术细节与实践难点。无论是数据库优化还是服务器配置,均结合真实场景与代码示例进行剖析,致力于提供可直接应用于工作的解决方案。
专业领域覆盖
精准对应开发生命周期。从前端界面到后端编程,从数据库操作到服务器运维,形成完整闭环,一站式满足全栈工程师和运维人员的技术需求。
即学即用高效
内容强调实操性,步骤清晰、代码完整。用户可根据教程直接复现和应用于自身项目,显著缩短从学习到实践的距离,快速解决开发中的具体问题。
持续更新保障
专注既定技术方向进行长期、稳定的内容输出。确保各栏目技术文章持续更新迭代,紧跟主流技术发展趋势,为用户提供经久不衰的学习价值。