Install
openclaw skills install @thcjp/pg-job-queueopenclaw skills install @thcjp/pg-job-queue功能说明: 本技能涵盖 自动化处理、专业能力支持、认领/进度跟踪、化处理 等核心能力。
Production-ready job queue using 关系型数据库 with priority scheduling, batch claiming, and progress tracking.
CREATE TABLE jobs (
id UUID PRIMARY KEY DEFAULT gen_random_uuid(),
job_type VARCHAR(50) NOT NULL,
priority INT NOT NULL DEFAULT 100,
status VARCHAR(20) NOT NULL DEFAULT 'pending',
data JSONB NOT NULL DEFAULT '{}',
-- Progress tracking
progress INT DEFAULT 0,
current_stage VARCHAR(100),
events_count INT DEFAULT 0,
-- Worker tracking
worker_id VARCHAR(100),
claimed_at TIMESTAMPTZ,
-- Timing
created_at TIMESTAMPTZ DEFAULT NOW(),
started_at TIMESTAMPTZ,
completed_at TIMESTAMPTZ,
-- Retry handling
attempts INT DEFAULT 0,
max_attempts INT DEFAULT 3,
last_error TEXT,
CONSTRAINT valid_status CHECK (
status IN ('pending', 'claimed', 'running', 'completed', 'failed', 'cancelled')
)
);
-- Critical: Partial index for fast claiming
CREATE INDEX idx_jobs_claimable ON jobs (priority DESC, created_at ASC)
WHERE status = 'pending';
CREATE INDEX idx_jobs_worker ON jobs (worker_id)
WHERE status IN ('claimed', 'running');
CREATE OR REPLACE FUNCTION claim_job_batch(
p_worker_id VARCHAR(100),
p_job_types VARCHAR(50)[],
p_batch_size INT DEFAULT 10
) RETURNS SETOF jobs AS $$
BEGIN
RETURN QUERY
WITH claimable AS (
SELECT id
FROM jobs
WHERE status = 'pending'
AND job_type = ANY(p_job_types)
AND attempts < max_attempts
ORDER BY priority DESC, created_at ASC
LIMIT p_batch_size
FOR UPDATE SKIP LOCKED -- Critical: skip locked rows
),
claimed AS (
UPDATE jobs
SET status = 'claimed',
worker_id = p_worker_id,
claimed_at = NOW(),
attempts = attempts + 1
WHERE id IN (SELECT id FROM claimable)
RETURNING *
)
SELECT * FROM claimed;
END;
$$ LANGUAGE plpgsql;
const (
PriorityExplicit = 150 // User-requested
PriorityDiscovered = 100 // System-discovered
PriorityBackfill = 30 // Background backfills
)
type JobQueue struct {
db *pgx.Pool
workerID string
}
func (q *JobQueue) Claim(ctx context.Context, types []string, batchSize int) ([]Job, error) {
rows, err := q.db.Query(ctx,
"SELECT * FROM claim_job_batch($1, $2, $3)",
q.workerID, types, batchSize,
)
if err != nil {
return nil, err
}
defer rows.Close()
var jobs []Job
for rows.Next() {
var job Job
if err := rows.Scan(&job); err != nil {
return nil, err
}
jobs = append(jobs, job)
}
return jobs, nil
}
func (q *JobQueue) Complete(ctx context.Context, jobID uuid.UUID) error {
_, err := q.db.execute(ctx, `
UPDATE jobs
SET status = 'completed',
progress = 100,
completed_at = NOW()
WHERE id = $1`,
jobID,
)
return err
}
func (q *JobQueue) Fail(ctx context.Context, jobID uuid.UUID, errMsg string) error {
_, err := q.db.execute(ctx, `
UPDATE jobs
SET status = CASE
WHEN attempts >= max_attempts THEN 'failed'
ELSE 'pending'
END,
last_error = $2,
worker_id = NULL,
claimed_at = NULL
WHERE id = $1`,
jobID, errMsg,
)
return err
}
func (q *JobQueue) RecoverStaleJobs(ctx context.Context, timeout time.Duration) (int, error) {
result, err := q.db.execute(ctx, `
UPDATE jobs
SET status = 'pending',
worker_id = NULL,
claimed_at = NULL
WHERE status IN ('claimed', 'running')
AND claimed_at < NOW() - $1::interval
AND attempts < max_attempts`,
timeout.String(),
)
if err != nil {
return 0, err
}
return int(result.RowsAffected()), nil
}
| Scenario | Approach |
|---|---|
| Need guaranteed delivery | 关系型数据库 queue |
| Need sub-ms latency | Use Redis instead |
| < 1000 jobs/sec | 关系型数据库 is fine |
| > 10000 jobs/sec | Add Redis layer |
| Need strict ordering | Single worker per type |
| 依赖项 | 类型 | 是否必需 | 获取方式 |
|---|---|---|---|
| LLM API | API | 必需 | 由Agent内置LLM提供 |
详细的输入输出格式请参考下方章节说明。
| 场景 | 输入 | 输出 |
|---|---|---|
| 基础使用 | 用户请求 | 处理结果 |
不适用于:需要人工判断的复杂决策场景
# 请参考上方使用说明进行配置和调用
result = "ready"
| 错误场景 | 原因 | 处理方式 |
|---|---|---|
| 配置错误 | 参数缺失或格式错误 | 检查依赖说明中的配置要求 |
| 运行时错误 | 运行环境不满足 | 确认运行环境符合依赖说明 |
| 网络错误 | 连接超时或不可达 | 检查网络连接后重试,参考国内替代方案 |
A: 请先阅读使用流程章节,确认环境满足依赖说明中的要求。
A: 请参考错误处理章节,按照表格中的处理方式操作。
A: 请参考已知限制章节了解具体限制。