Hotdry.

Article

Agent 异步 run 队列:PostgreSQL SKIP LOCKED 租约、心跳续租与至少一次语义

交互式 SSE 断连后 Agent 仍需在 Worker 上跑完多轮 tool 环。本文给出 PostgreSQL 上 pending/leased/done/dead 状态机、SELECT FOR UPDATE SKIP LOCKED 领取、locked_until 心跳与可见性超时、取消位轮询的可落地表,并说明至少一次重放与写工具幂等边界。

2026-07-20systems

生产 Agent 很少只在「HTTP 请求线程里同步跑完整个 tool 环」。用户刷新页面、移动网络闪断、网关 60s 空闲超时之后,run 仍可能要继续:后台批处理、长 shell、多轮检索与写文件往往超过单次请求生命周期。此时编排器需要把「一次用户触发的 run」落成可领取的作业:Worker 崩溃后能被另一节点接管,取消要能跨进程可见,租约过期不能静默丢任务。本文只讨论 用 PostgreSQL 做 Agent run 队列的租约模型与领取语义,不重复 run 级 max_turns/deadline、各工具沙箱,以及写工具幂等去重表的细节(见近期专文与文末交叉引用)。

问题背景:同步请求生命周期装不下 Agent run

典型交互路径:

  1. 客户端 POST /runs 创建 run(写入用户消息、策略、预算快照)。
  2. 若在请求内同步执行:SSE 推送 token 与 tool 事件;连接断开时,若生命周期绑定 socket,Worker 会中止或变成无人认领的孤儿进程。
  3. 若改为异步:API 只负责入队并返回 run_id;独立 Worker 领取后执行「模型 → tool_calls → 回灌」循环,客户端用 SSE/WebSocket 订阅 事件流。

在缺少显式队列与租约时,常见故障包括:

  • 多 Worker 双执行:两个节点同时读到「未完成」行,各开一轮 shell/SQL,产生重复副作用。
  • 假死占坑:Worker OOM 后行仍标记 running,任务永不被重试,也永不进死信。
  • 无心跳的长任务:可见性超时(visibility timeout)设成 30s,而单次 tool 要 2 分钟,任务会在中途被第二 Worker 再次领取。
  • 取消不可见:用户点 Stop 只断了前端 SSE;跑在另一进程的模型流与子进程组仍在消耗配额。
  • 至少一次被当成恰好一次:租约过期重放后,若写工具无幂等键,会重复发邮件、重复 git push

因此应把 run 建模为带租约的状态机,而不是进程内布尔标志:

enqueue(run) → status=pending
claim: SELECT … FOR UPDATE SKIP LOCKED → status=leased, locked_until=now()+lease
  loop work:
    heartbeat: locked_until = now()+lease
    poll cancel_requested_at
    model / tools(绑定 run 预算与取消令牌)
ack success → status=done
nack / exception:
  attempts++ ; if attempts >= max → dead else pending(或延后 available_at)
lease expired(无心跳)→ 其他 Worker 可再次 claim

状态流转示意(仓库内图):

Agent 任务队列状态流转:pending → leased(SKIP LOCKED)→ done /dead,以及 locked_until 过期回 pending

可落地实现:表结构、领取 SQL、心跳与取消

1. 最小表结构

一张 agent_runs 即可支撑「入队 + 租约 + 取消 + 死信」;事件与消息可另表,避免大字段拖慢领取扫描。

CREATE TABLE agent_runs (
  id                uuid PRIMARY KEY,
  tenant_id         text NOT NULL,
  status            text NOT NULL
                    CHECK (status IN ('pending', 'leased', 'done', 'dead', 'cancelled')),
  priority          int  NOT NULL DEFAULT 0,
  attempts          int  NOT NULL DEFAULT 0,
  max_attempts      int  NOT NULL DEFAULT 3,
  available_at      timestamptz NOT NULL DEFAULT now(),
  locked_until      timestamptz,
  locked_by         text,              -- worker_id,如 host:pid:boot_id
  cancel_requested_at timestamptz,
  budget_json       jsonb NOT NULL DEFAULT '{}'::jsonb,
  input_json        jsonb NOT NULL,    -- 用户消息引用、会话 id 等;大 blob 外置对象存储
  error_code        text,
  error_message     text,
  created_at        timestamptz NOT NULL DEFAULT now(),
  updated_at        timestamptz NOT NULL DEFAULT now(),
  finished_at       timestamptz
);

CREATE INDEX agent_runs_claim_idx
  ON agent_runs (priority DESC, available_at, created_at)
  WHERE status = 'pending';

CREATE INDEX agent_runs_lease_idx
  ON agent_runs (locked_until)
  WHERE status = 'leased';

建议默认参数:

参数 建议起点 作用
lease_duration_s 3060 单次租约长度;Worker 必须在过期前心跳
heartbeat_interval_s lease_duration_s / 3(如 10~20) 续租周期;过密增加写放大,过疏易误回收
max_attempts 35 含首次;超过进 dead 供人工 / 自动分析
claim_batch_size 1(交互式)或 416(批处理) 单次事务领取行数;交互式常 1 run/Worker 协程
backoff_base_s 530 失败重入队的 available_at 指数退避基数
cancel_poll_s 0.52 循环内轮询 cancel_requested_at;可叠加 LISTEN
stale_lease_reaper_s 3060 可选清道夫:把「已过期仍标 leased」修回 pending(见下)

说明:

  • 领取条件应是 status = 'pending' AND available_at <= now(),不要用「locked_until 为空」 alone 表示可领 —— 失败重试常需要延迟可见。
  • 租约过期后的所有权:严格模式可让 reaper 把过期 leased 改回 pending;也可在 claim 时允许 status = 'leased' AND locked_until < now() 被再次锁定(需同一事务内 CAS 式更新,避免双领)。

2. 领取:FOR UPDATE SKIP LOCKED

PostgreSQL 的 SELECT … FOR UPDATE SKIP LOCKED 会跳过已被其它事务行锁住的候选行,避免多 Worker 在同一批 pending 上阻塞等待。领取必须在短事务内完成:锁行 → 更新租约字段 → 提交;不要把整个 Agent 多轮循环放在持有行锁的事务里。

-- 单次领取 1 条(交互式 run)
BEGIN;

WITH cte AS (
  SELECT id
  FROM agent_runs
  WHERE status = 'pending'
    AND available_at <= now()
  ORDER BY priority DESC, available_at, created_at
  FOR UPDATE SKIP LOCKED
  LIMIT 1
)
UPDATE agent_runs AS r
SET status       = 'leased',
    locked_by    = $1,                    -- worker_id
    locked_until = now() + ($2 * interval '1 second'),  -- lease_duration_s
    attempts     = r.attempts + 1,
    updated_at   = now()
FROM cte
WHERE r.id = cte.id
RETURNING r.*;

COMMIT;

要点:

  • attempts 在 claim 时 +1,而不是成功时 +1:这样 Worker 崩溃后重试次数会前进,避免无限重放。
  • SKIP LOCKED 不是分布式锁服务的替代品:它只保证「当前持有行锁的事务」互斥;提交后互斥靠 status/locked_until 约定,必须配合心跳与 CAS 更新。
  • 高优先级插队用 priority;同一优先级 FIFO 用 available_at, created_at

Python 示意(psycopg3):

from __future__ import annotations

import os
import socket
from dataclasses import dataclass
from datetime import datetime, timedelta, timezone
from typing import Any
from uuid import UUID

import psycopg
from psycopg.rows import dict_row


def worker_id() -> str:
    return f"{socket.gethostname()}:{os.getpid()}:{os.getenv('BOOT_ID', 'dev')}"


CLAIM_SQL = """
WITH cte AS (
  SELECT id
  FROM agent_runs
  WHERE status = 'pending'
    AND available_at <= now()
  ORDER BY priority DESC, available_at, created_at
  FOR UPDATE SKIP LOCKED
  LIMIT 1
)
UPDATE agent_runs AS r
SET status = 'leased',
    locked_by = %(worker_id)s,
    locked_until = now() + (%(lease_s)s * interval '1 second'),
    attempts = r.attempts + 1,
    updated_at = now()
FROM cte
WHERE r.id = cte.id
RETURNING r.id, r.tenant_id, r.attempts, r.max_attempts,
          r.budget_json, r.input_json, r.cancel_requested_at
"""


@dataclass
class ClaimedRun:
    id: UUID
    tenant_id: str
    attempts: int
    max_attempts: int
    budget: dict[str, Any]
    input: dict[str, Any]
    cancel_requested_at: datetime | None


def claim_one(conn: psycopg.Connection, *, lease_s: int = 45) -> ClaimedRun | None:
    with conn.transaction():
        with conn.cursor(row_factory=dict_row) as cur:
            cur.execute(CLAIM_SQL, {"worker_id": worker_id(), "lease_s": lease_s})
            row = cur.fetchone()
            if not row:
                return None
            return ClaimedRun(
                id=row["id"],
                tenant_id=row["tenant_id"],
                attempts=row["attempts"],
                max_attempts=row["max_attempts"],
                budget=row["budget_json"],
                input=row["input_json"],
                cancel_requested_at=row["cancel_requested_at"],
            )

3. 心跳续租:只延长「自己的」租约

长 run 必须周期性续租。更新条件必须包含 locked_bystatus = 'leased',防止在租约已过期并被他人领取后,旧 Worker 的迟到心跳把新租约覆盖掉。

UPDATE agent_runs
SET locked_until = now() + ($2 * interval '1 second'),
    updated_at   = now()
WHERE id = $1
  AND status = 'leased'
  AND locked_by = $3
  AND (cancel_requested_at IS NULL)  -- 可选:已取消则拒绝续租,促使尽快收尾
RETURNING locked_until;
HEARTBEAT_SQL = """
UPDATE agent_runs
SET locked_until = now() + (%(lease_s)s * interval '1 second'),
    updated_at = now()
WHERE id = %(run_id)s
  AND status = 'leased'
  AND locked_by = %(worker_id)s
RETURNING locked_until
"""


def heartbeat(conn: psycopg.Connection, run_id: UUID, *, lease_s: int = 45) -> bool:
    with conn.cursor() as cur:
        cur.execute(
            HEARTBEAT_SQL,
            {"run_id": run_id, "lease_s": lease_s, "worker_id": worker_id()},
        )
        return cur.fetchone() is not None

Worker 主循环中:

on claim:
  start heartbeat task every heartbeat_interval_s
  link cancel token to run budget / model stream / tools
  while not terminal:
    if heartbeat failed → stop work(租约已丢,禁止继续写副作用)
    if cancel polled → controlled finalize(status=cancelled)
    do one model turn / tool batch
  ack done | nack retry | dead
  stop heartbeat task

心跳失败即停工是硬约束:若续租返回 0 行,说明 reaper 或其它 Worker 已接管(或状态已终态)。继续执行会造成双写。

4. 租约过期回收(reaper)与失败重入队

两种等价策略,二选一写进运维手册并测并发:

策略 A — 显式 reaper(推荐可读性):

UPDATE agent_runs
SET status       = 'pending',
    locked_by    = NULL,
    locked_until = NULL,
    available_at = now(),          -- 或 now() + jitter
    updated_at   = now()
WHERE status = 'leased'
  AND locked_until < now();

策略 B — claim 时包含过期 leased:

WHERE (
    (status = 'pending' AND available_at <= now())
    OR (status = 'leased' AND locked_until < now())
  )

失败(可重试异常)时:

UPDATE agent_runs
SET status = CASE
               WHEN attempts >= max_attempts THEN 'dead'
               ELSE 'pending'
             END,
    locked_by = NULL,
    locked_until = NULL,
    available_at = CASE
                     WHEN attempts >= max_attempts THEN available_at
                     ELSE now() + make_interval(secs => LEAST(300, 5 * power(2, attempts - 1)))
                   END,
    error_code = $2,
    error_message = $3,
    finished_at = CASE WHEN attempts >= max_attempts THEN now() ELSE NULL END,
    updated_at = now()
WHERE id = $1
  AND status = 'leased'
  AND locked_by = $4;

成功与用户取消:

-- 成功
UPDATE agent_runs
SET status = 'done', locked_by = NULL, locked_until = NULL,
    finished_at = now(), updated_at = now()
WHERE id = $1 AND status = 'leased' AND locked_by = $2;

-- API 侧请求取消(任意时刻)
UPDATE agent_runs
SET cancel_requested_at = COALESCE(cancel_requested_at, now()),
    updated_at = now()
WHERE id = $1 AND status IN ('pending', 'leased');

-- Worker 观察到取消并完成收尾后
UPDATE agent_runs
SET status = 'cancelled', locked_by = NULL, locked_until = NULL,
    finished_at = now(), updated_at = now()
WHERE id = $1 AND status = 'leased' AND locked_by = $2;

5. 取消与订阅:队列层只保证「可见」,不保证「瞬时停」

取消是协作式的:

  1. 控制面写 cancel_requested_at(上节 SQL)。
  2. Worker 在每个模型调用前每批 tool 边界读取该字段(或 LISTEN agent_run_cancel + 载荷 run_id)。
  3. 将取消传播到:LLM 流的 AbortController / HTTP 取消、工具进程组、SQL 会话终止(与 run 预算专文同一顺序)。
  4. 收尾写 status = cancelled,并持久化已完成的部分事件供 UI 展示。

pending 且尚未被领取的 run:claim 前即可被标 cancelled,领取逻辑应排除 cancel_requested_at IS NOT NULL,或领取后立即走取消收尾、不调用模型。

6. 事件流与「至少一次」交付的边界

队列保证的是 run 作业的至少一次执行尝试,不是工具副作用的恰好一次。

层次 语义 工程含义
队列领取 同时至多一个 有效租约持有者 SKIP LOCKED + locked_by CAS
租约过期重放 至少一次 run 尝试 崩溃后另一 Worker 从检查点或从头继续
工具副作用 默认非事务、不可自动回滚 写工具必须幂等键 / HITL / 去重表
客户端事件 可重复投递 SSE 重连用 Last-Event-ID 或事件序号去重展示

若要在崩溃后「从上一轮 tool 之后继续」而不是整 run 重放,需要额外 checkpoint(已提交的消息列表、tool_call_id 结果、预算累计)。队列租约只回答「谁有权跑」,不回答「跑到哪一步」;检查点应与消息表同一数据库事务写入,避免「结果已回灌模型但检查点未落盘」。

7. 与 LISTEN/NOTIFY 的配合(可选)

纯轮询 claim 在低 QPS 下足够;要降低空转延迟可用:

  • Worker 在空 claim 后 LISTEN agent_runs_new
  • 入队事务 COMMITNOTIFY agent_runs_new(或按 tenant_id 分通道);
  • NOTIFY 丢失时仍靠轮询兜底(PostgreSQL NOTIFY 不持久)。

不要用 NOTIFY 载荷传递完整 prompt;载荷只唤醒,状态以表为准。

风险与边界

至少一次 ≠ 恰好一次。 租约过期时旧 Worker 可能仍卡在 DNS 或无超时的 read() 上,新 Worker 已开始执行。必须:心跳失败停工、工具 I/O 带 deadline、写路径幂等。OWASP GenAI 的 Excessive Agency 风险在「自动重试写工具」时会被放大。

attempts 与业务重试不要混用。 模型对可恢复工具错误的「再试一次」应计入 run 内 max_tool_calls / max_model_turns,而不是每次都让队列 attempts+1。队列 attempts 只计量 Worker 级失败或租约丢失后的重新领取

长事务与连接池。 领取事务必须短;Agent 循环中的业务 SQL 使用独立连接。PgBouncer 事务池模式下,会话级 LISTEN 与长会话状态不兼容,LISTEN 应走 session 模式连接或独立连接。

时钟与 timestamptz 所有比较使用数据库 now(),不要用 Worker 本地时钟判断租约是否过期。跨区域多主或逻辑复制延迟场景下,队列应落在单一主库写入。

表膨胀。 高频更新 locked_until 会产生死元组;需合理 heartbeat_interval、为热表配置 autovacuum,并考虑将「活跃租约」与「历史 run」分区或归档。done/dead 行应从 claim 部分索引中排除(上文 partial index)。

优先级饥饿。 只按 priority DESC 时,低优先级可能长期得不到领取;可对同租户加公平性(例如 per-tenant 并发上限 + 老化 priority)。

多租户噪声邻居。 全局一个 claim 队列时,大租户批任务会占满 Worker。应按 tenant_id 配额(信号量或「每轮 claim 跳过已超额租户」)限制飞行 run 数。

取消宽限与误杀。 移动网络闪断不应立刻 cancel 后台报告任务;交互式「仅前端订阅断开」与「用户显式 Stop」要分字段或分 API。宽限期内允许同一 run_id 续订事件流。

安全。 error_message、事件表、trace 可能含密钥与绝对路径,需与工具结果脱敏同一套规则;locked_by 仅运维可见,避免把内部主机名暴露给终端用户。

不要把模型流式半包执行与队列重放混为一谈。 流式 tool_calls 参数拼装仍须在 committed 之后执行;队列重放的是已持久化的 run 状态,不是 TCP 半包。

参考来源

systems

内容声明:本文无广告投放、无付费植入。

如有事实性问题,欢迎发送勘误至 i@hotdrydog.com