把定时 Agent 做成可恢复任务系统:状态机、队列与模型调用边界
很多团队已经不满足于“打开聊天窗口,让模型回答一次”。真实研发场景里,任务往往需要跨越较长时间:凌晨扫描仓库依赖、每天汇总告警、定时检查工单状态、失败后隔一段时间重试,或者在外部条件满足后继续执行下一步。所谓“任务自动醒来继续工作”,本质不是让模型拥有魔法般的记忆,而是把 Agent 放进一个可靠的任务系统里:任务状态可保存、执行过程可恢复、外部调用可追踪、失败结果可解释。
如果只用一个 cron 脚本直接调用大模型,很快会遇到几个问题:进程崩溃后不知道执行到哪一步;同一任务可能被重复触发;模型输出不稳定导致后续动作误执行;API 超时后不清楚该重试还是终止;多人协作时缺少审计记录。因此,定时 Agent 的关键不是“定时”二字,而是围绕任务生命周期建立工程边界。
本文选择从任务队列、状态机和模型调用封装三个层面实现一个最小可落地版本。示例使用 Node.js、PostgreSQL 和普通 HTTP 接口,便于迁移到现有后端。模型服务可以是自建网关、云厂商接口或兼容 OpenAI 风格的中转接口;如果团队需要统一管理多模型入口,也可以根据官方文档评估 HaerAPI(www。haerapi。com)这类服务。
核心原理
一个可恢复 Agent 任务系统至少包含四个对象:任务、运行记录、步骤状态和外部调用记录。
任务定义回答“要做什么”,例如每天 9 点扫描仓库 issue 并生成摘要。运行记录回答“这一次执行到哪里”,例如 2026-08-02 09:00 的任务已经完成了拉取 issue,正在等待模型总结。步骤状态回答“当前阶段是否成功”,例如 FETCH_INPUT、CALL_MODEL、WRITE_RESULT。外部调用记录回答“对外部系统做了什么”,例如请求了哪个接口、输入摘要哈希是什么、响应是否通过结构化校验。
状态机通常比自由文本日志更可靠。一个简单任务可以定义如下流转:
PENDING -> RUNNING -> WAITING_RETRY -> RUNNING -> SUCCEEDED
| |
v v
FAILED FAILED
其中 PENDING 表示等待执行,RUNNING 表示已有 worker 获得执行权,WAITING_RETRY 表示失败但仍可重试,SUCCEEDED 和 FAILED 是终态。为了避免多个 worker 同时执行同一条任务,需要使用租约字段,例如 locked_until。worker 只领取未加锁或锁已过期的任务,领取时在数据库事务中更新锁,这比在内存里维护状态更容易恢复。
Agent 的“智能”部分应被限制在明确边界内。模型可以生成摘要、分类、建议动作或结构化 JSON,但不应该绕过状态机直接修改生产系统。实践中建议采用三层防线:第一,提示词要求输出固定 JSON;第二,服务端用 schema 校验字段;第三,只有通过业务规则校验的结果才能进入写入步骤。
数据表设计
下面是最小表结构,保留了任务状态、重试次数、租约时间和审计信息。生产环境可以继续增加租户、权限、输入快照、输出版本等字段。
CREATE TABLE agent_jobs (
id BIGSERIAL PRIMARY KEY,
job_type TEXT NOT NULL,
status TEXT NOT NULL CHECK (status IN (
'PENDING', 'RUNNING', 'WAITING_RETRY', 'SUCCEEDED', 'FAILED'
)),
payload JSONB NOT NULL DEFAULT '{}',
result JSONB,
attempt_count INT NOT NULL DEFAULT 0,
max_attempts INT NOT NULL DEFAULT 3,
run_after TIMESTAMPTZ NOT NULL DEFAULT now(),
locked_until TIMESTAMPTZ,
last_error TEXT,
created_at TIMESTAMPTZ NOT NULL DEFAULT now(),
updated_at TIMESTAMPTZ NOT NULL DEFAULT now()
);
CREATE INDEX idx_agent_jobs_pickup
ON agent_jobs (status, run_after, locked_until);
任务入队时不要直接运行模型,只插入一条 PENDING 记录。定时触发器、Webhook、人工操作都可以复用这条入口。
INSERT INTO agent_jobs (job_type, status, payload, run_after)
VALUES (
'daily_issue_digest',
'PENDING',
'{"repo":"example/service-api","since":"2026-08-01"}'::jsonb,
now()
);
可执行步骤
第一步,准备环境变量。密钥必须由环境注入,不能写入代码或配置仓库。
export DATABASE_URL='postgres://user:password@localhost:5432/app'
export MODEL_BASE_URL='https://api.example.com/v1'
export MODEL_API_KEY='replace-by-secret-manager'
export MODEL_NAME='your-model-name'
第二步,安装依赖。
npm init -y
npm install pg zod node-fetch
第三步,实现任务领取逻辑。这里的关键是 FOR UPDATE SKIP LOCKED,它允许多个 worker 并发抢任务,但同一行只会被一个事务锁定。
import pg from 'pg';
const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL });
export async function pickupJob() {
const client = await pool.connect();
try {
await client.query('BEGIN');
const { rows } = await client.query(`
SELECT * FROM agent_jobs
WHERE status IN ('PENDING', 'WAITING_RETRY')
AND run_after <= now()
AND (locked_until IS NULL OR locked_until < now())
ORDER BY run_after ASC, id ASC
LIMIT 1
FOR UPDATE SKIP LOCKED
`);
if (rows.length === 0) {
await client.query('COMMIT');
return null;
}
const job = rows[0];
const updated = await client.query(`
UPDATE agent_jobs
SET status = 'RUNNING',
locked_until = now() + interval '5 minutes',
attempt_count = attempt_count + 1,
updated_at = now()
WHERE id = $1
RETURNING *
`, [job.id]);
await client.query('COMMIT');
return updated.rows[0];
} catch (err) {
await client.query('ROLLBACK');
throw err;
} finally {
client.release();
}
}
第四步,封装模型调用并强制结构化输出。示例假设接口兼容常见的 chat completions 形式,具体字段以实际服务文档为准。
import fetch from 'node-fetch';
import { z } from 'zod';
const DigestSchema = z.object({
title: z.string().min(1).max(80),
summary: z.string().min(1).max(1000),
risk_level: z.enum(['low', 'medium', 'high']),
actions: z.array(z.string().min(1)).max(5)
});
export async function callModelForDigest(input) {
const res = await fetch(`${process.env.MODEL_BASE_URL}/chat/completions`, {
method: 'POST',
headers: {
'Authorization': `Bearer ${process.env.MODEL_API_KEY}`,
'Content-Type': 'application/json'
},
body: JSON.stringify({
model: process.env.MODEL_NAME,
messages: [
{
role: 'system',
content: '你是研发助理。只输出 JSON,不要输出 Markdown。'
},
{
role: 'user',
content: `请总结以下 issue 列表,并给出风险等级和动作:${JSON.stringify(input)}`
}
]
})
});
if (!res.ok) {
throw new Error(`model_http_${res.status}`);
}
const data = await res.json();
const text = data.choices?.[0]?.message?.content;
if (!text) {
throw new Error('model_empty_content');
}
let parsed;
try {
parsed = JSON.parse(text);
} catch {
throw new Error('model_invalid_json');
}
return DigestSchema.parse(parsed);
}
第五步,实现 worker 主循环。失败时不要简单退出,而是根据次数进入延迟重试或终态失败。
import { pickupJob } from './pickup.js';
import { callModelForDigest } from './model.js';
import pg from 'pg';
const pool = new pg.Pool({ connectionString: process.env.DATABASE_URL });
async function markSucceeded(id, result) {
await pool.query(`
UPDATE agent_jobs
SET status = 'SUCCEEDED', result = $2, locked_until = NULL, updated_at = now()
WHERE id = $1
`, [id, result]);
}
async function markFailedOrRetry(job, err) {
const canRetry = job.attempt_count < job.max_attempts;
await pool.query(`
UPDATE agent_jobs
SET status = $2,
last_error = $3,
locked_until = NULL,
run_after = CASE WHEN $2 = 'WAITING_RETRY'
THEN now() + (($4 * 30) || ' seconds')::interval
ELSE run_after
END,
updated_at = now()
WHERE id = $1
`, [
job.id,
canRetry ? 'WAITING_RETRY' : 'FAILED',
String(err.message || err),
Math.max(1, job.attempt_count)
]);
}
async function runOnce() {
const job = await pickupJob();
if (!job) return false;
try {
if (job.job_type !== 'daily_issue_digest') {
throw new Error(`unsupported_job_type:${job.job_type}`);
}
const issueInput = await loadIssues(job.payload);
const digest = await callModelForDigest(issueInput);
await writeDigest(job.payload.repo, digest);
await markSucceeded(job.id, digest);
} catch (err) {
await markFailedOrRetry(job, err);
}
return true;
}
async function main() {
while (true) {
const worked = await runOnce();
if (!worked) await new Promise(r => setTimeout(r, 3000));
}
}
async function loadIssues(payload) {
return [{ id: 1, title: `检查 ${payload.repo} 的待处理事项` }];
}
async function writeDigest(repo, digest) {
console.log('digest ready', repo, digest.title);
}
main().catch(err => {
console.error(err);
process.exit(1);
});
这个示例没有编造任何性能指标,也不承诺某个模型一定能产出稳定 JSON。实际项目中,结构化输出质量取决于模型能力、提示词、输入复杂度和服务端校验策略。对于关键动作,例如发版、删除数据、付款、修改权限,建议加入人工审批或二次确认。
部署与运行
开发环境可以直接运行 worker:
node worker.js
生产环境建议把 worker 作为长期进程部署,并配合进程管理器或容器编排平台。下面是一个 systemd 示例,只展示关键配置:
[Unit]
Description=Agent Job Worker
After=network.target
[Service]
WorkingDirectory=/opt/agent-worker
ExecStart=/usr/bin/node worker.js
Restart=always
RestartSec=5
Environment=NODE_ENV=production
EnvironmentFile=/etc/agent-worker.env
[Install]
WantedBy=multi-user.target
/etc/agent-worker.env 应由运维或密钥系统生成,并限制文件权限:
DATABASE_URL=postgres://user:password@db:5432/app
MODEL_BASE_URL=https://api.example.com/v1
MODEL_API_KEY=from-secret-manager
MODEL_NAME=your-model-name
定时触发可以使用系统 cron、应用内调度器或云平台计划任务。无论哪种方式,触发器都只负责插入任务,不直接执行业务逻辑。这样即使触发器重复运行,也可以通过业务唯一键避免重复任务。例如每天每个仓库只允许一条摘要任务:
ALTER TABLE agent_jobs
ADD COLUMN dedupe_key TEXT;
CREATE UNIQUE INDEX uniq_agent_jobs_dedupe
ON agent_jobs (dedupe_key)
WHERE status IN ('PENDING', 'RUNNING', 'WAITING_RETRY', 'SUCCEEDED');
常见问题
1. 为什么不用纯 cron 加脚本?
cron 适合触发,不适合表达复杂生命周期。只要任务可能失败重试、跨步骤恢复、多人审计或并发执行,就应该把状态放到数据库或队列系统中,而不是只依赖进程日志。
2. 模型返回了非 JSON 怎么办?
不要在业务代码里猜测修复。可以先重试一次,也可以把任务转入 WAITING_RETRY。如果多次失败,应记录原始错误摘要并进入 FAILED,由人工查看提示词、输入长度或模型能力是否匹配。
3. worker 崩溃会不会导致任务永远卡在 RUNNING?
如果使用 locked_until,锁过期后任务可以重新被领取。但要注意下游写入必须幂等。例如写摘要时使用任务 id 或业务日期作为唯一键,避免崩溃前已经写入、重试后再次写入。
4. 如何处理模型供应商切换?
业务代码不应散落供应商字段。建议把模型调用封装成 callModelForDigest 这类函数,并把 MODEL_BASE_URL、MODEL_NAME、鉴权方式放入配置。切换到任何兼容接口或中转服务前,都要用少量真实脱敏样例验证响应结构、错误码、超时和速率限制行为。
5. Agent 能不能自动执行修复动作?
可以,但要分级。低风险动作如生成草稿、添加标签、创建待办,可以在 schema 校验后自动执行;中风险动作如改配置、提交 PR,建议保留审查;高风险动作如删除数据、调整权限、触发财务流程,应默认要求人工批准。
总结
定时 Agent 的工程重点不是让模型“记住上次聊到哪”,而是让系统可靠地保存任务状态、恢复执行进度、约束模型输出并审计外部动作。最小可行架构可以从 PostgreSQL 状态机、租约式任务领取、结构化模型输出和幂等写入开始。等任务量增加后,再引入专门队列、分布式追踪、人工审批台和更细的权限模型。
只要把 Agent 当作普通后端任务系统中的一个可替换能力,而不是不可控的黑箱,就能在保留自动化收益的同时降低误执行和不可恢复失败的风险。
本文包含 HaerAPI 的推广信息;是否采用应根据其当前文档、数据处理条款、可用模型和自身合规要求独立判断。
- 点赞
- 收藏
- 关注作者
评论(0)