切换主题
并发、锁与租约
问题:两个后台进程同时拿到同一任务
普通 SELECT status='PENDING' 后执行,多个处理器都可能看到同一条记录。进程内 mutex 只能保护一个进程;增加服务副本后,它们有各自的锁。
持久队列需要原子认领:判断可处理与写入所有权必须处于同一受保护步骤。处理器死亡后还要有办法释放工作,因此需要带截止时间的租约。
原理:租约提供回收,token 提供所有权
记录中保存 claim_id 与 expires_at。认领后提交短事务,再执行外部动作。更新成功必须匹配自己的 claim_id;过期被重新认领后,旧处理器的确认应失败。
但数据库中的 claim_id 不能阻止旧处理器继续执行外部动作。若外部系统能接收递增 fencing token 并拒绝旧值,可以保护其资源。否则还要使用稳定业务键、远端幂等和对账。
最小例子:认领与确认分离
sql
-- 教学 SQL:在短事务内认领,外部调用不持有行锁。
WITH next AS (
SELECT id FROM jobs
WHERE status = 'PENDING'
ORDER BY available_at
FOR UPDATE SKIP LOCKED LIMIT 10
)
UPDATE jobs SET status = 'RUNNING', claim_id = $1, expires_at = $2
FROM next WHERE jobs.id = next.id RETURNING jobs.*;
-- 完成时只允许当前所有者提交确认。
UPDATE jobs SET status = 'DONE'
WHERE id = $3 AND claim_id = $1 AND status = 'RUNNING';完整系统还需认领过期 RUNNING、续租、失败时间和受影响行数校验。示例只说明原子认领与所有权原则。
方案比较
| 协调手段 | 适用范围 | 关键限制 |
|---|---|---|
| mutex | 单进程内存状态 | 不协调多个服务副本 |
| 行锁 | 数据库短事务 | 不能长时间锁住网络调用 |
| advisory lock | 按业务键协调事务 | 所有协作者必须遵守同一规则 |
| 租约与认领 | 跨重启后台处理 | 过期不等于旧处理器物理停止 |
| fencing token | 下游支持的资源操作 | 下游必须实际校验单调 token |
SKIP LOCKED 适合工作队列减少等待,不适合一般一致查询。跳过被锁工作可能影响公平性,因此还需观察最老等待时间。
真实案例:Outbox 与 Attempt 调度
Outbox SQL 同时选候选和更新认领;标记成功按 claim_id 确认。调度事务按 StepID+AttemptNumber 获取 transaction advisory lock,验证 scheduling claim 与 Execution 状态,复用已存在 Attempt。
已核对的实现 · 本地代码快照
Outbox 认领与所有权确认 backend · c6dbe05c
server/db/queries/runtime.sql · 第 4–50 行
符号:ClaimOutboxEvents / MarkOutboxEventPublished · 核对日期 2026-10-02
来源与提交版本一致
-- name: ClaimOutboxEvents :many
WITH candidates AS (
SELECT id
FROM outbox_events
WHERE (outbox_events.status IN ('PENDING', 'FAILED') AND outbox_events.available_at <= sqlc.arg(now))
OR (outbox_events.status = 'PUBLISHING' AND outbox_events.claim_expires_at <= sqlc.arg(now))
ORDER BY outbox_events.available_at, outbox_events.created_at
FOR UPDATE SKIP LOCKED
LIMIT sqlc.arg(batch_limit)
)
UPDATE outbox_events AS event
SET status = 'PUBLISHING',
attempts = event.attempts + 1,
claim_id = sqlc.arg(claim_id),
claim_expires_at = sqlc.arg(claim_expires_at)
FROM candidates
WHERE event.id = candidates.id
RETURNING event.id, event.aggregate_id, event.event_type, event.sequence, event.payload,
event.created_at, event.attempts, event.claim_id;
-- name: GetNextOutboxAttemptAt :one
SELECT MIN(candidate.due_at)::timestamptz AS due_at
FROM (
(SELECT available_at AS due_at
FROM outbox_events
WHERE status IN ('PENDING', 'FAILED')
ORDER BY available_at
LIMIT 1)
UNION ALL
(SELECT claim_expires_at AS due_at
FROM outbox_events
WHERE status = 'PUBLISHING' AND claim_expires_at IS NOT NULL
ORDER BY claim_expires_at
LIMIT 1)
) AS candidate;
-- name: MarkOutboxEventPublished :execrows
UPDATE outbox_events
SET status = 'PUBLISHED', published_at = sqlc.arg(published_at), last_error = NULL,
claim_id = NULL, claim_expires_at = NULL
WHERE id = sqlc.arg(id) AND status = 'PUBLISHING' AND claim_id = sqlc.arg(claim_id);
-- name: MarkOutboxEventFailed :execrows
UPDATE outbox_events
SET status = sqlc.arg(status), available_at = sqlc.arg(available_at),
last_error = sqlc.arg(last_error), claim_id = NULL, claim_expires_at = NULL
WHERE id = sqlc.arg(id) AND status = 'PUBLISHING' AND claim_id = sqlc.arg(claim_id);片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
调度事务与幂等 Attempt backend · c6dbe05c
server/internal/adapters/outbound/postgres/attempt_dispatch.go · 第 45–133 行
符号:dispatchAttempt · 核对日期 2026-10-02
来源与提交版本一致
func (dispatcher *AttemptDispatcher) dispatchAttempt(ctx context.Context, request executionports.CreateAttemptRequest, claim *pgqueries.SchedulingRequest) (*executionports.AttemptRef, error) {
stepID, err := uuid.Parse(request.StepID)
if err != nil {
return nil, fmt.Errorf("parse step ID: %w", err)
}
tx, err := dispatcher.database.BeginTx(ctx, pgx.TxOptions{
IsoLevel: pgx.ReadCommitted,
})
if err != nil {
return nil, err
}
defer func() { _ = tx.Rollback(ctx) }()
txQueries := dispatcher.queries.WithTx(tx)
lockKey := request.StepID + ":" + fmt.Sprint(request.AttemptNumber)
if err := txQueries.AcquireTransactionAdvisoryLock(ctx, lockKey); err != nil {
return nil, err
}
if _, err := txQueries.LockSchedulingRequest(ctx, pgqueries.LockSchedulingRequestParams{
ID: claim.ID,
ClaimID: claim.ClaimID,
}); err != nil {
return nil, err
}
executionIDValue, err := uuid.Parse(request.ExecutionID)
if err != nil {
return nil, fmt.Errorf("%w: execution ID", executiondomain.ErrDispatchInvalid)
}
execution, err := txQueries.LockExecutionByID(ctx, executionIDValue)
if err != nil {
return nil, err
}
if execution.Status != "CREATED" && execution.Status != "QUEUED" && execution.Status != "RUNNING" {
return nil, fmt.Errorf("%w: execution closed", executiondomain.ErrDispatchInvalid)
}
existing, err := txQueries.FindAttemptByStepAndNumber(ctx, pgqueries.FindAttemptByStepAndNumberParams{
StepID: stepID,
AttemptNumber: int32(request.AttemptNumber),
})
if err == nil {
if !existing.WorkerID.Valid {
return nil, errors.New("existing attempt has no worker")
}
if err := txQueries.CompleteSchedulingRequest(ctx, pgqueries.CompleteSchedulingRequestParams{
ID: claim.ID,
ClaimID: claim.ClaimID,
State: "DISPATCHED",
AvailableAt: timestamp(time.Now().UTC()),
}); err != nil {
return nil, err
}
if err := tx.Commit(ctx); err != nil {
return nil, err
}
existingWorker := uuid.UUID(existing.WorkerID.Bytes)
_ = dispatcher.waker.WakeWorker(existingWorker.String())
return &executionports.AttemptRef{
AttemptID: existing.ID.String(),
LeaseVersion: existing.LeaseVersion,
}, nil
}
if !errors.Is(err, pgx.ErrNoRows) {
return nil, err
}
now := time.Now().UTC()
workerID, generation, err := dispatcher.reserveWorker(ctx, txQueries, stepID, request, now)
if err != nil {
if errors.Is(err, errNoEligibleWorker) {
if commitErr := tx.Commit(ctx); commitErr != nil {
return nil, commitErr
}
}
return nil, err
}
contextRow, err := txQueries.FindStepExecutionAndWorkerNode(ctx, pgqueries.FindStepExecutionAndWorkerNodeParams{
StepID: stepID,
WorkerID: workerID,
})
if err != nil {
return nil, err
}
executionID, nodeID := contextRow.ExecutionID, contextRow.NodeID
if err := appendTimelineTx(ctx, txQueries, timelineRecord{
ExecutionID: executionID,
MilestoneKey: "step.worker.selected:" + stepID.String() + ":" + fmt.Sprint(request.AttemptNumber),
Stage: "SCHEDULING",
AttemptNumber: request.AttemptNumber,
State: "SUCCEEDED",
Component: "Scheduler",
Title: "已选定 Node 与 Capability 唯一 Worker",片段展示核对时的源码;完整文件指纹用于检测后续变化。这里的路径用于定位,不要求手机访问源码仓库。
两层检查处理不同竞争:调度认领保护处理请求,Attempt 稳定键保护重复创建,Worker slot 复核保护资源占用。仅在纯函数里选了“最佳节点”还没有预留容量。
失败与边界:迟到的成功
A 认领任务后停顿;租约到期 B 接手;A 恢复并发送成功。确认条件必须拒绝 A 的陈旧 token。若 A 已在外部产生副作用,拒绝数据库确认也不能撤回它,需下游策略。
租约时间不是随意数字:太短造成重复,太长增加故障恢复时延。应结合执行延迟、续租周期和停顿风险。数据库时间与应用时间的选择也应统一,不混用多个不同步的时钟。
迁移练习与参考答案
练习:定时发送优惠券,每张券只能发放一次,服务跑三个副本。只在 Go 中加 mutex 够吗?
参考答案:不够。使用数据库认领和租约回收,以用户+活动建立唯一发放事实;远端短信以稳定消息键幂等。迟到处理器确认匹配 claim_id,旧任务不能覆盖新所有者。能重复发送短信与能重复发券是不同要求,分别定义保证。