跳到主要内容
赞助推荐 Claude Team 合租,少折腾账号
>80aj_
架构设计

从长事务到可靠投递:任务调度链路的修复方法

15 分钟阅读阅读(5)
赞助推荐 团队协作里的 AI 办公工作台

一个任务调度系统出现过一个怪现象。任务下发正常,结果回写却越来越慢。应用日志没有明显异常,数据库更新语句也不复杂。只要批量任务开始运行,回写请求就会长时间停在数据库里。

真正的故障点不在回写代码。调度程序打开数据库事务后,又去调用外部接口。它还要循环处理任务和发送消息。事务一直没有提交,前面改过的记录也一直没有释放行锁。回写程序想更新同一条记录,只能在后面等。

赞助推荐 一人公司 · 创业装备库
赞助推荐 一人公司 · 创业装备库

修复方向不该是继续压缩回写 SQL。应该先缩小事务边界。事务只保护必要的数据库修改,网络调用放到事务外。

慢 SQL 可能只是排队的人

长事务把网络等待变成行锁等待

数据库事务可以理解成在柜台办理一组不能拆开的手续。手续办完以前,相关资料不能交给别人修改。行锁就是柜台暂时扣住的那份资料。

下面是一种常见写法。

@Transactional
public void dispatchBatch(List<Task> tasks) {
  for (Task task : tasks) {
    if (!presenceClient.isOnline(task.getTargetId())) {
      continue;
    }

    taskRepository.markRunning(task.getId());
    messageBroker.send(task.toMessage());
  }
}

代码看起来很整齐。事务却包住了两类不可控工作:外部接口和消息系统。只要某次网络调用变慢,事务就跟着变长。循环前部修改过的记录,要等整批任务结束才释放锁。

此时查看数据库,容易把注意力放在等待中的 UPDATE 上。那条语句未必慢,它可能只是在等另一笔事务交还行锁。加索引或精简回写逻辑只能改善局部耗时。堵塞源头仍然存在。

排查这类故障,我会先画一条时间线,而不是先改 SQL:

事务开始
  ├─ 更新任务 A ────────────────┐
  ├─ 外部请求                    │ A 的行锁仍在
  ├─ 更新任务 B                  │
  ├─ 消息发送                    │
  └─ 整批结束并提交 ─────────────┘

结果回写 A ── 等待行锁 ─────────▶ 获得锁并更新

时间线能分开两种耗时。数据库可能真的执行得慢,也可能在执行前已经等了很久。

把事务缩到数据库内部

修复长事务的直接办法,是把外部工作移到事务外。在线检查只是一种过滤,目的是少发无效消息。它不应该决定数据库并发是否正确。

public void dispatch(Task task) {
  if (!presenceClient.isOnline(task.getTargetId())) {
    return;
  }

  boolean claimed = taskService.claim(task.getId());
  if (!claimed) {
    return;
  }

  messageBroker.send(task.toMessage());
}

@Transactional
public boolean claim(String taskId) {
  return taskRepository.markRunningIfRunnable(taskId) == 1;
}

事务现在只做一次状态修改,并尽快提交。外部接口变慢时,当前调用依然会慢。数据库行锁不再覆盖整段网络时间。

拆掉大事务后,还要检查一个竞争条件。代码是不是先查状态,再根据查询结果修改状态。

线程 A:查到 RUNNABLE ───────▶ 更新 RUNNING
线程 B:查到 RUNNABLE ───────▶ 更新 RUNNING

两个线程可能读到相同旧状态,然后各自发送消息。大事务没有了,重复领取的问题反而更容易暴露。

用条件更新领取任务

CAS 是 Compare-And-Swap 的缩写。在数据库中,它通常是一条带旧状态条件的更新语句。任务仍可领取,当前调用才能把它改成运行中。

UPDATE task_execution
SET state = 'RUNNING',
    attempts = attempts + 1,
    updated_at = CURRENT_TIMESTAMP
WHERE id = :task_id
  AND state = 'RUNNABLE';

应用程序只需要检查受影响的行数。

int affectedRows = taskRepository.claim(taskId);
if (affectedRows == 0) {
  return;
}

messageBroker.send(message);

受影响行数为一,说明当前线程拿到了任务。受影响行数为零,说明任务已被其他线程领取,当前线程直接结束。数据库成了并发裁判,不再依赖应用代码碰巧按顺序运行。

这里有一个小边界。在线检查完成后,目标可能立刻离线。目标也可能在检查时离线,随后马上恢复。在线检查只能减少无效发送,不能充当正确性保证。真正的并发保证来自带状态条件的更新。

CAS 防止重复领取,却管不了数据库与消息系统的一致性。

状态已经改变,消息却没有发出

考虑下面的故障窗口。

数据库:任务状态已经提交为 RUNNING
进程:准备发送消息时退出
消息系统:没有收到任务消息

数据库认为任务正在运行,消费者却从未收到指令。任务会停在一个看似合理、实际没有执行者的状态里。

把消息发送放回数据库事务,也不能真正消除这个窗口。数据库和消息系统各自提交。普通的本地事务不能让两边严格地一起成功。相反的情况也会出现。消息已经送达,数据库提交却失败,消费者随后收到重复任务。

工程上通常用 Transactional Outbox 处理。Outbox 就是一张数据库发件箱表。业务事务不承诺消息已经送达。它只负责记下一条可靠的“待发送事件”。

Outbox 把消息变成可恢复的待办

短事务、CAS 与 Outbox 可靠投递流程

领取任务时,系统要完成两次写入。状态更新和 Outbox 事件必须共用一个数据库事务。

BEGIN;

UPDATE task_execution
SET state = 'RUNNING',
    attempts = attempts + 1,
    updated_at = CURRENT_TIMESTAMP
WHERE id = :task_id
  AND state = 'RUNNABLE';

-- 仅在上面的 UPDATE 修改成功时写入
INSERT INTO event_outbox (
  event_id,
  aggregate_id,
  event_type,
  payload,
  status,
  available_at,
  created_at
) VALUES (
  :event_id,
  :task_id,
  'TASK_CLAIMED',
  :payload,
  'PENDING',
  CURRENT_TIMESTAMP,
  CURRENT_TIMESTAMP
);

COMMIT;

状态更新和待发消息写在同一笔数据库事务里。事务成功,两条记录都存在;事务失败,两条记录一起回滚。

后台投递器负责扫描 PENDING 事件。它领取一批记录,并把状态改成 SENDING。事务提交后,投递器再调用消息系统。发送成功就标记为 SENT,失败则记录错误并安排重试。

PENDING ──短事务领取──▶ SENDING ──发送成功──▶ SENT
                           │
                           ├─发送失败──▶ 延迟后重试
                           └─租约过期──▶ 重新领取

投递器也不能一边持有数据库锁,一边等待消息系统。比较稳妥的流程分为三个短步骤:

  • 用短事务领取事件,并写入投递租约。
  • 在事务外发送消息。
  • 用另一个短事务记录成功或失败。

多实例投递可以用 FOR UPDATE SKIP LOCKED。一个实例锁定一批记录。其他实例跳过它们,继续领取后面的事件。

SELECT id, event_id, payload
FROM event_outbox
WHERE status = 'PENDING'
  AND available_at <= CURRENT_TIMESTAMP
ORDER BY id
LIMIT :batch_size
FOR UPDATE SKIP LOCKED;

Outbox 表要保存事件编号、业务对象编号、消息类型和内容。它还要保存状态、可重试时间、尝试次数和最近错误。事件编号应有唯一约束。历史数据要定期归档或清理,否则发件箱会变成新的大表问题。

至少一次投递必须配幂等

Outbox 保证的是 At-Least-Once。消息至少有机会送达一次,但可能重复。它不保证消息永远只出现一次。

一个典型窗口是,投递器已经把消息交给 Broker。它却在写入 SENT 之前退出。服务恢复后,同一条事件会再次发送。这个重复不能靠运气避免,消费者必须能识别它。

常见做法有两种。

  • 消费者保存处理过的 event_id。唯一约束负责挡住重复事件。
  • 消费者用条件更新推进业务状态。当前状态不符合预期时,直接忽略消息。

消费动作可能继续调用支付、邮件或 HTTP 回调。此时,下游也要使用同一个幂等键。只在数据库入口幂等,挡不住下游副作用重复发生。

我更愿意把“可靠投递”理解成恢复能力。进程可以在任何步骤退出。恢复后,它仍能根据数据库记录继续工作。系统不必假设网络稳定,也不必假设每次调用只发生一次。

不要急着把调度器重写成任务平台

任务租约和独立 Dispatcher 都很有用。它们可以支持暂停、限速和灰度。任务规模较大或运维要求明确时,可以考虑它们。

但长事务故障出现时,直接重写完整调度架构通常不是最好的起点。短事务和条件更新能处理锁等待与重复领取。Outbox 负责消息丢失。结构更大的改造,应由真实需求推动。

如果系统确实需要租约,任务状态可以按下面的方向演进。

READY ──领取并设置到期时间──▶ LEASED
  ▲                              │
  └────租约过期,允许重试─────────┤
                                 └─收到结果──▶ SUCCESS / FAILED

租约的价值是把“任务是不是卡住了”从猜测变成明确字段。代价是新的状态机、迁移脚本、调度组件和运维规则。

上线前要故意制造故障

正常路径通过,只能说明系统在一切顺利时可用。可靠性改造要检查故障窗口。

  • 并发启动多个领取线程,确认只有一个 CAS 更新成功。
  • 让外部在线检查变慢,确认数据库事务和行锁没有覆盖等待时间。
  • 暂停消息系统,确认 Outbox 事件仍然保留。恢复后,它应该继续投递。
  • 消息发送成功后,在标记 SENT 前结束投递器。消费者应该忽略重复事件。
  • 重启业务进程,确认没有永久停留在 SENDING 的事件。
  • 达到重试上限后,事件应该进入失败状态。错误原因必须可以查询。

监控也要跟着语义走。接口耗时不够,我还会观察事务持续时间和行锁等待。Outbox 积压量、最老事件等待时间也很重要。重试次数和幂等冲突数,则能说明投递质量。

如果你正在处理类似故障,先把事务时间线画出来。标出每个网络调用、数据库写入和提交点。然后再决定缩短事务、增加条件更新,还是引入 Outbox。不要从等待最久的那条 SQL 开始猜。

赞助推荐 一键部署 AI 大模型
赞助推荐 一键部署 AI 大模型
赞(0)
未经允许不得转载:80aj » 从长事务到可靠投递:任务调度链路的修复方法
赞助推荐 低成本上手 Claude Code 的中转选择
赞助推荐 低成本上手 Claude Code 的中转选择
赞助推荐 一键部署 AI 大模型
赞助推荐 一键部署 AI 大模型