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