一个任务投送系统出现过两种症状。待执行任务进入 Kafka 的速度不够,数据库也频繁出现锁等待。两种症状来自同一段调度代码。
旧流程用一个线程逐条发送任务。整批投送又共用一个数据库事务。前面的任务已经改过状态,却要等批次尾部的任务发送结束。批次越慢,数据库锁就占得越久。
本次修复没有重写消息架构。它只做两件事。多线程负责投送,每个任务使用独立事务。任务进入队列后标记为 RUNNING。周期调度只选择 RUNABLE。队列里的任务不会重复投送。
本文沿用系统中的拼写 RUNABLE。它表示“当前可以投送”的任务状态。
单线程和长事务叠在了一起

旧代码可以简化成下面的样子。
@Transactional
public void dispatchBatch(List<Task> tasks) {
for (Task task : tasks) {
taskRepository.markRunning(task.getId());
kafkaPublisher.send(task.toMessage());
}
}
单线程让任务只能排队发送。一个任务没有完成投送,下一个任务就不能开始。Kafka 发送速度因此受限于整条串行链路。
大事务又把批次耗时变成了持锁时间。任务在循环前部改成 RUNNING。相关行锁却不会马上释放。数据库要等整个方法结束,才能提交这批修改。
可以把旧流程理解成一个仓库只有一名发货员。他每拿一件货,还把前面所有订单都压在桌上。整批货发完以前,其他人不能修改这些订单。
一个线程:任务 A ──▶ 任务 B ──▶ 任务 C ──▶ …… ──▶ 批次结束
一个事务:└──────── 所有已修改记录持续持锁 ────────┘
只增加线程,数据库事务仍可能很长。只拆事务,任务投送仍然串行。两个问题要一起处理。
用有界线程池并行投送
批量调度现在只负责找到待执行任务,再把每个任务交给线程池。多个任务可以并行准备和发送,不再排成一条长队。
public class BatchDispatcher {
private final Executor dispatchExecutor;
private final SingleTaskDispatcher singleTaskDispatcher;
public void dispatchRunnableTasks(List<String> taskIds) {
for (String taskId : taskIds) {
dispatchExecutor.execute(
() -> singleTaskDispatcher.dispatch(taskId));
}
}
}
线程池必须有上限。无限创建线程只会改变排队位置。任务最终会堵在数据库连接池或 Kafka 客户端。并发数、等待队列和拒绝策略都要能配置,也要能监控。
这里不建议用一次性的无限并行写法。系统需要知道线程池何时已满。新任务是等待、拒绝,还是留到下个周期,也要有明确规则。
多线程解决的是投送吞吐。数据库锁能否及时释放,还要看事务边界。
每个任务使用独立事务

每个工作线程调用一次单任务服务。单任务服务读取状态、修改状态并向 Kafka 发送任务。一次调用只处理一个任务,也只占用一笔事务。
public class SingleTaskDispatcher {
private final TaskRepository taskRepository;
private final KafkaPublisher kafkaPublisher;
@Transactional
public void dispatch(String taskId) {
Task task = taskRepository.findById(taskId);
if (task == null || task.getState() != TaskState.RUNABLE) {
return;
}
taskRepository.markRunning(taskId);
kafkaPublisher.send(task.toMessage());
}
}
单任务事务仍然可能包含一次 Kafka 发送。区别在于,任何网络等待只影响当前任务。它不会让整批任务修改过的数据库记录一起持锁。
任务发送完成后,当前事务就可以提交并释放行锁。某个任务失败时,回滚范围也限于当前任务,不会拖住整个批次。
Spring 项目还有一个容易漏掉的实现细节。@Transactional 依赖 Spring 代理。this.dispatch(taskId) 会绕过代理。独立事务可能不生效。
可以把批量调度和单任务投送放在两个 Spring Bean 中。线程池调用单任务 Bean。事务代理会包住每次任务调用。
RUNNING 是周期投送的状态门
周期调度只查询 RUNABLE 任务。任务发往 Kafka 时,数据库状态改成 RUNNING。任务还在队列或执行中,下一轮调度就不会选择它。
周期调度查询:WHERE state = 'RUNABLE'
RUNABLE ──投送 Kafka──▶ RUNNING
│
├─仍在队列或执行中:不再投送
└─执行失败或超时:回到 RUNABLE
下一轮周期调度只会重新选择已经回到 RUNABLE 的任务
这里的重试不是“每个周期都再发一次”。系统先用 RUNNING 把任务挡在周期查询之外。执行失败或超时后,状态才改回 RUNABLE。后续周期可以再次投送它。
这个状态回路承担了任务级恢复。Kafka 队列暂时积压时,任务仍是 RUNNING。下一轮调度不会再发一条相同消息。
失败回退也要保留尝试次数、最近错误和更新时间。持续失败的任务可能在两个状态之间循环。没有这些记录,运维人员看不到原因。
为什么锁冲突会减少
旧流程的锁范围是整个批次。一个任务发送变慢,已经处理过的许多记录都会延迟提交。
新流程把锁范围缩到单个任务。工作线程处理完一条任务,就提交一条任务。其他线程写不同任务时,不需要等待整个批次结束。
旧流程
事务 T:任务 A ── 任务 B ── 任务 C ── 批次提交
└──────── 多行锁一起保留 ────────┘
新流程
线程 A:任务 A ── 提交
线程 B:任务 B ── 提交
线程 C:任务 C ── 提交
每个任务只保留自己的短事务
多线程本身不会消灭数据库竞争。如果不同线程处理同一个任务,它们仍可能争用同一行。当前设计要求任务清单不重复。RUNNING 状态负责阻止后续周期再次投送。
如果将来出现多个调度实例,多个入口可能并发领取同一任务。届时可以改用条件更新。CAS 只允许一个调用成功。它负责把 RUNABLE 改成 RUNNING。
UPDATE task_execution
SET state = 'RUNNING'
WHERE id = :task_id
AND state = 'RUNABLE';
CAS 是并发加固,不是本次修复的主体。只有存在多入口竞争时,它才解决实际问题。
Outbox 不是本次修复的一部分
每任务事务缩短了锁时间。数据库和 Kafka 仍是两个事务系统。进程退出或网络异常时,系统仍要恢复任务。恢复依靠失败检测、超时和状态回退。
当前机制使用业务级恢复。执行失败或超时后,任务回到 RUNABLE。后续周期再投送它。这个做法符合已有状态机。
业务以后可能要求状态和待发消息严格绑定。那时可以评估 Transactional Outbox。Outbox 把待发事件保存在数据库中,再由独立投递器发送。它会增加表结构、重试和清理。消费端也要处理幂等。本次修改没有引入这些机制。
上线前怎样验证
修复不能只看线程池启动,也不能只看任务最终进入 Kafka。下面几条要一起验证。
- 用同一批任务对比串行和并行投送耗时,确认投送吞吐确实提高。
- 观察数据库事务持续时间,确认不再出现覆盖整个批次的长事务。
- 连续运行周期调度。确认
RUNNING任务不会重新进入 Kafka。 - 制造执行失败和超时。确认状态回到
RUNABLE,后续周期可以重新投送。 - 让线程池达到容量上限。拒绝或延迟行为必须可观察,任务不能静默丢失。
- 在单个任务发送时制造异常,确认其他任务的事务仍能独立提交。
监控要覆盖线程池活动数、等待队列和拒绝次数。还要观察单任务事务耗时、行锁等待和各状态任务数量。超时回退次数也不能漏。吞吐提高后,瓶颈可能转移。检查数据库连接池和 Kafka 客户端。
如果你正在修类似问题,先分开画出批次并发方式和单条事务边界。并行负责加快投送。短事务及时释放锁。状态机避免重复,也负责恢复失败。三件事不要混成模糊的“异步优化”。






