AsyncTaskWorker 深度解析:任务执行引擎的调度与容错设计
引言
如果说租约机制是异步任务框架的"锁",那么 AsyncTaskWorker 就是这把锁的"持有者"和"使用者"。Worker 负责任务的完整生命周期管理:从数据库领取任务、调用业务 Handler 执行、定期续约、处理异常、直到任务完成或失败。
本文将深入剖析 Worker 的内部实现,揭秘它是如何在分布式环境下实现高效、可靠的任务调度。
Worker 的核心职责
/**
* AsyncTaskWorker - 任务调度与执行的核心组件
*
* 核心职责:
* 1. 任务领取:周期性从数据库领取待执行任务(基于乐观锁)
* 2. 任务执行:调用对应的任务处理器执行业务逻辑
* 3. 租约管理:定期续约(心跳),防止任务被其他 Worker 抢占
* 4. 状态管理:更新任务状态(RUNNING → SUCCESS/FAILED)
* 5. 异常处理:处理超时、取消、执行失败等异常情况
* 6. 任务恢复:恢复租约过期的僵尸任务,重新排队
*/
public class AsyncTaskWorker implements SmartLifecycle, DisposableBean {
// ...
}整体架构
1. 线程模型
Worker 内部使用两个线程池:
public AsyncTaskWorker(...) {
// 调度线程池:2 个线程,执行定时任务
this.scheduler = Executors.newScheduledThreadPool(2,
namedThreadFactory("async-task-scheduler-"));
// 执行线程池:大小 = localMaxConcurrency
this.executor = Executors.newFixedThreadPool(
properties.getWorker().getLocalMaxConcurrency(),
namedThreadFactory("async-task-worker-"));
// 本地并发控制信号量
this.localSlots = new Semaphore(
properties.getWorker().getLocalMaxConcurrency()
);
}设计理念:
- 调度线程池:负责定时任务(轮询、心跳、恢复),线程数固定为 2
- 执行线程池:负责执行业务 Handler,线程数等于本地最大并发数
- 信号量控制:在应用层控制并发,避免线程池过载
2. 三大定时任务
@Override
public void start() {
// 任务1: 轮询任务队列(立即开始,周期执行)
scheduler.scheduleWithFixedDelay(
this::pollSafely,
0, // 立即开始
pollInterval.toMillis(),
TimeUnit.MILLISECONDS
);
// 任务2: 心跳续约(延迟启动,周期执行)
scheduler.scheduleWithFixedDelay(
this::heartbeatSafely,
heartbeatInterval.toMillis(), // 延迟启动
heartbeatInterval.toMillis(),
TimeUnit.MILLISECONDS
);
// 任务3: 租约恢复(延迟启动,周期执行)
scheduler.scheduleWithFixedDelay(
this::recoverSafely,
recoveryInterval.toMillis(),
recoveryInterval.toMillis(),
TimeUnit.MILLISECONDS
);
}时序图:
时间轴(秒) 0 1 2 ... 10 11 ... 60 61
↓ ↓ ↓ ↓ ↓ ↓ ↓
轮询任务 ●────●────●─────●────●────●─────●────●──── (每 1s)
心跳任务 ▲─────────────▲─────────────▲──────── (每 10s)
恢复任务 △────────────────── (每 60s)3. 活跃任务追踪
// 活跃任务映射表:taskId -> ActiveExecution
private final Map<String, ActiveExecution> activeExecutions = new ConcurrentHashMap<>();
private static final class ActiveExecution {
private final AsyncTask task; // 任务快照
private final Instant startedAt; // 开始时间
private final AtomicBoolean timeoutRequested; // 超时标志
private volatile Future<?> future; // 执行 Future
}用途:
- 心跳续约时批量续约所有活跃任务
- 超时检测时中断执行线程
- 优雅停机时等待活跃任务完成
核心流程详解
1. 任务轮询(Poll)
1.1 轮询逻辑
private void poll() {
// 持续领取,直到没有任务或槽位用尽
while (running.get() && localSlots.tryAcquire()) {
// 1. 构造领取请求
TaskClaimRequest request = new TaskClaimRequest(
workerId, // 当前 Worker ID
handlerRegistry.registeredTaskTypes(), // 已注册的任务类型
clock.instant(), // 当前时间
properties.getWorker().getLeaseTimeout(), // 租约时长
properties.getLimits().getDefaultGlobalConcurrency(),
properties.getLimits().getDefaultTypeConcurrency(),
properties.getWorker().getCandidateLimit()
);
// 2. 领取任务
AsyncTask claimed = store.claimNext(request).orElse(null);
if (claimed == null) {
// 没有任务,释放槽位并退出
localSlots.release();
return;
}
// 3. 记录活跃任务
ActiveExecution active = new ActiveExecution(claimed, clock.instant());
activeExecutions.put(claimed.id(), active);
if (metrics != null) metrics.activeIncrement();
// 4. 提交到执行线程池
try {
Future<?> future = executor.submit(() -> execute(active));
active.future = future;
} catch (RejectedExecutionException ex) {
// 线程池拒绝,保留 RUNNING 状态,等待租约恢复
activeExecutions.remove(claimed.id());
localSlots.release();
log.error("线程池拒绝任务, taskId={}, 等待租约恢复", claimed.id());
}
}
}关键设计点:
- 非阻塞领取:
tryAcquire()不阻塞,没有槽位立即返回 - 批量领取:
while循环持续领取,直到没有任务或槽位用尽 - 类型过滤:只领取已注册 Handler 的任务类型
- 异常处理:线程池拒绝时不破坏任务状态,由租约恢复兜底
1.2 领取策略
数据库端实现(MyBatis):
<select id="claimNext" resultMap="asyncTaskMap">
-- 子查询 + 更新
UPDATE system_async_task
SET status = 'RUNNING',
lease_owner = #{workerId},
lease_until = #{leaseUntil},
attempt_count = attempt_count + 1,
update_time = NOW()
WHERE id = (
SELECT id FROM system_async_task
WHERE status = 'QUEUED'
AND task_type IN
<foreach collection="taskTypes" item="type" open="(" separator="," close=")">
#{type}
</foreach>
AND (available_at IS NULL OR available_at <= #{now})
AND (deadline_at IS NULL OR deadline_at > #{now})
ORDER BY priority DESC, available_at ASC, create_time ASC
LIMIT 1
FOR UPDATE SKIP LOCKED -- 跳过已被其他事务锁定的行
);
-- 返回领取的任务
SELECT * FROM system_async_task
WHERE lease_owner = #{workerId}
AND status = 'RUNNING'
ORDER BY update_time DESC
LIMIT 1;
</select>排序策略:
- priority DESC:优先级高的先执行
- available_at ASC:可用时间早的先执行(支持延迟任务)
- create_time ASC:创建时间早的先执行(FIFO)
并发控制:
// 在领取前检查并发限制
public Optional<AsyncTask> claimNext(TaskClaimRequest request) {
return transactionTemplate.execute(status -> {
// 1. 检查全局并发
long runningCount = countRunningTasks();
if (runningCount >= globalMaxConcurrency) {
return Optional.empty();
}
// 2. 检查类型并发
for (String taskType : request.taskTypes()) {
long typeRunning = countRunningTasksByType(taskType);
if (typeRunning >= typeMaxConcurrency) {
continue; // 跳过此类型
}
// 3. 领取该类型的任务
return claimNextOfType(taskType, request);
}
return Optional.empty();
});
}2. 任务执行(Execute)
2.1 执行流程
private void execute(ActiveExecution active) {
AsyncTask task = active.task;
long startedNanos = System.nanoTime();
try {
// 1. 记录队列等待时长
if (metrics != null && task.createdAt() != null) {
Duration queueWait = Duration.between(task.createdAt(), active.startedAt);
metrics.recordQueueWait(queueWait);
}
// 2. 创建执行上下文
DefaultTaskExecutionContext context = new DefaultTaskExecutionContext(
task, workerId, store, clock,
flushInterval, flushThreshold
);
// 3. 通知监听器:任务开始
notifyStarted(task);
// 4. 查找 Handler 并执行
AsyncTaskHandler<?> handler = handlerRegistry.requireHandler(task.taskType());
TaskExecutionResult result = invokeHandler(handler, context, task.payloadJson());
// 5. 刷新进度
context.flush();
// 6. 校验完成状态
context.validateCompletion();
context.checkCancellation();
// 7. 再次检查超时标志
if (active.timeoutRequested.get()) {
throw new TaskTimeoutException("任务执行超过最长时间");
}
// 8. 检查任务截止时间
Instant deadline = task.deadlineAt();
if (deadline == null && !maxExecutionTime.isZero()) {
deadline = active.startedAt.plus(maxExecutionTime);
}
if (deadline != null && !deadline.isAfter(clock.instant())) {
throw new TaskTimeoutException("任务已超过截止时间: " + deadline);
}
// 9. 写入成功状态
AsyncTaskStatus target = result.status() == SUCCESS
? AsyncTaskStatus.SUCCESS
: AsyncTaskStatus.PARTIAL_SUCCESS;
if (!store.complete(task.id(), workerId, target, result.artifactId(), clock.instant())) {
// 失去租约,检查是否被取消
if (store.isCancellationRequested(task.id())) {
handleCancellation(task);
return;
}
log.error("失去租约,无法写入终态, taskId={}", task.id());
return;
}
// 10. 通知监听器:任务完成
AsyncTask completed = store.findById(task.id()).orElse(task);
notifyCompleted(completed);
if (metrics != null) metrics.completed();
} catch (TaskCancellationException ex) {
// 任务被取消
handleCancellation(task);
} catch (Throwable ex) {
// 任务执行异常
if (!running.get()) {
// Worker 停机中,不记录失败
log.warn("Worker 停机导致任务中断, taskId={}", task.id());
return;
}
if (store.isCancellationRequested(task.id())) {
handleCancellation(task);
return;
}
if (active.timeoutRequested.get()) {
ex = new TaskTimeoutException("任务执行超过最长时间");
}
handleFailure(task, ex);
} finally {
// 11. 记录执行时长
Duration executionDuration = Duration.ofNanos(System.nanoTime() - startedNanos);
if (metrics != null) {
metrics.recordExecution(task.taskType(), executionDuration);
}
// 12. 慢任务告警
Duration slowThreshold = properties.getWorker().getSlowTaskThreshold();
if (!slowThreshold.isZero() && executionDuration.compareTo(slowThreshold) > 0) {
log.warn("慢任务告警, taskId={}, duration={}ms",
task.id(), executionDuration.toMillis());
notifySlowTask(task, executionDuration);
}
// 13. 清理活跃任务
activeExecutions.remove(task.id());
localSlots.release();
if (metrics != null) metrics.activeDecrement();
}
}流程图:
graph TD
A[开始执行] --> B[创建执行上下文]
B --> C[通知监听器: onStarted]
C --> D[调用 Handler.execute]
D --> E{执行结果}
E -->|成功| F[刷新进度]
F --> G[校验完成状态]
G --> H{检查超时/取消}
H -->|正常| I[写入成功状态]
I -->|写入成功| J[通知监听器: onCompleted]
I -->|写入失败| K[检查取消请求]
K -->|是| L[处理取消]
K -->|否| M[等待租约恢复]
E -->|失败| N[失败分类]
N --> O{是否可重试}
O -->|是| P[调度重试]
O -->|否| Q[标记失败]
Q --> R[移入死信队列]
H -->|超时| S[抛出超时异常]
H -->|取消| T[抛出取消异常]
J --> U[清理资源]
L --> U
P --> U
R --> U
M --> U
U --> V[结束]2.2 Handler 调用
@SuppressWarnings("unchecked")
private <P> TaskExecutionResult invokeHandler(
AsyncTaskHandler<?> rawHandler,
DefaultTaskExecutionContext context,
String payloadJson
) throws Exception {
// 类型转换
AsyncTaskHandler<P> handler = (AsyncTaskHandler<P>) rawHandler;
// 反序列化 Payload
P payload = objectMapper.readValue(payloadJson, handler.payloadType());
// 执行前校验(动态业务前提)
handler.validate(payload);
// 再次检查取消标志
context.checkCancellation();
// 执行业务逻辑
return handler.execute(context, payload);
}设计要点:
- ✅ 泛型安全:类型转换由框架保证
- ✅ 执行前校验:防止权限或业务状态变化后继续处理
- ✅ 取消检测:在执行前最后一次检查取消标志
2.3 执行上下文
public class DefaultTaskExecutionContext implements TaskExecutionContext {
@Override
public String taskId() {
return task.id();
}
@Override
public void reportProgress(AsyncTaskProgress progress) {
// 节流刷新:累积到阈值或定时刷新
this.latestProgress = progress;
long now = System.currentTimeMillis();
if (now - lastFlushTime >= flushInterval.toMillis() ||
progress.processedCount() - lastFlushCount >= flushThreshold) {
flush();
}
}
@Override
public void checkCancellation() throws TaskCancellationException {
if (store.isCancellationRequested(task.id())) {
throw new TaskCancellationException("任务已被取消");
}
// 同时检查线程中断标志
if (Thread.interrupted()) {
throw new TaskCancellationException("任务执行线程被中断");
}
}
public void flush() {
if (latestProgress != null) {
store.updateProgress(task.id(), workerId, latestProgress, clock.instant());
lastFlushTime = System.currentTimeMillis();
lastFlushCount = latestProgress.processedCount();
}
}
}进度刷新策略:
- 时间阈值:每隔 1 秒刷新一次(可配置)
- 数量阈值:每处理 500 条记录刷新一次(可配置)
- 强制刷新:任务完成前强制刷新最终进度
3. 心跳续约(Heartbeat)
@Scheduled(fixedDelay = 10000)
private void heartbeat() {
Instant newLeaseUntil = clock.instant().plus(leaseTimeout);
// 1. 获取活跃任务快照
Map<String, ActiveExecution> snapshot = Map.copyOf(activeExecutions);
List<String> renewableTaskIds = new ArrayList<>();
// 2. 检查每个任务的超时状态
for (ActiveExecution active : snapshot.values()) {
// 计算截止时间
Instant deadline = active.task.deadlineAt();
if (deadline == null && !maxExecutionTime.isZero()) {
deadline = active.startedAt.plus(maxExecutionTime);
}
// 检查是否超时
if (deadline != null && !deadline.isAfter(clock.instant())) {
// 设置超时标志
active.timeoutRequested.set(true);
if (metrics != null) metrics.recordTimeout();
log.error("任务执行超时,正在中断, taskId={}, deadline={}",
active.task.id(), deadline);
// 中断执行线程
if (active.future != null) {
active.future.cancel(true);
}
continue;
}
renewableTaskIds.add(active.task.id());
}
// 3. 批量续约
if (!renewableTaskIds.isEmpty()) {
Set<String> renewed = store.renewLeases(
renewableTaskIds,
workerId,
newLeaseUntil
);
// 4. 检查续约失败的任务
for (String taskId : renewableTaskIds) {
if (!renewed.contains(taskId)) {
// 续约失败,检查任务状态
AsyncTask current = store.findById(taskId).orElse(null);
if (current != null && current.status().isTerminal()) {
// 任务已完成,无需处理
log.info("任务已完成,无需续约, taskId={}", taskId);
continue;
}
// 续约失败,中断执行
if (metrics != null) {
metrics.heartbeatFailure();
metrics.recordLeaseRenewalFailure();
}
log.error("租约续期失败,中断执行, taskId={}", taskId);
ActiveExecution active = activeExecutions.get(taskId);
if (active != null && active.future != null) {
active.future.cancel(true);
}
}
}
}
}心跳流程图:
graph TD
A[开始心跳] --> B[遍历活跃任务]
B --> C{检查超时}
C -->|超时| D[设置超时标志]
D --> E[中断线程]
C -->|正常| F[加入续约列表]
F --> G[批量续约]
G --> H{续约结果}
H -->|成功| I[继续执行]
H -->|失败| J[检查任务状态]
J -->|已完成| K[忽略]
J -->|仍运行| L[中断线程]4. 租约恢复(Recovery)
@Scheduled(fixedDelay = 60000)
private void recover() {
Instant now = clock.instant();
// 1. 恢复等待重试的任务
int dueRetries = store.enqueueDueRetries(
now,
recoveryBatchSize
);
// 2. 恢复租约过期的任务
RecoveryResult result = store.recoverExpiredLeases(
now,
recoveryBatchSize
);
// 3. 记录指标
if (metrics != null) {
metrics.recordRecovery(
result.requeuedCount(),
result.failedCount(),
result.cancelledCount()
);
}
// 4. 通知监听器
invokeListeners(listener -> listener.onRecoveryBatch(result),
workerId, "onRecoveryBatch");
// 5. 日志记录
if (dueRetries > 0 || result.totalCount() > 0) {
log.warn("租约恢复完成, dueRetries={}, requeued={}, failed={}, cancelled={}",
dueRetries, result.requeuedCount(),
result.failedCount(), result.cancelledCount());
}
}恢复策略:
public RecoveryResult recoverExpiredLeases(Instant now, int limit) {
return transactionTemplate.execute(status -> {
// 查询租约过期的任务
List<AsyncTask> expired = jdbcTemplate.query(
"SELECT * FROM system_async_task " +
"WHERE status IN ('RUNNING', 'CANCEL_REQUESTED') " +
" AND lease_until < ? " +
"ORDER BY lease_until ASC " +
"LIMIT ? " +
"FOR UPDATE SKIP LOCKED",
now, limit
);
int requeued = 0, failed = 0, cancelled = 0;
for (AsyncTask task : expired) {
if (task.status() == CANCEL_REQUESTED) {
// 取消请求 → 直接标记为已取消
markCancelled(task.id());
cancelled++;
} else if (task.attemptCount() < task.maxAttempts()) {
// 未达重试上限 → 重新排队
requeueForRetry(task.id());
requeued++;
} else {
// 达到重试上限 → 标记为失败
markFailed(task.id(), "租约过期且达到最大重试次数");
failed++;
}
}
return new RecoveryResult(requeued, failed, cancelled);
});
}异常处理与重试
1. 失败分类
private void handleFailure(AsyncTask task, Throwable throwable) {
// 1. 使用失败分类器分析异常
TaskFailure failure = failureClassifier.classify(task, throwable);
// 2. 记录失败日志(不暴露敏感信息)
log.error("任务执行失败, taskId={}, taskType={}, attempt={}, " +
"errorCode={}, traceId={}, retryable={}, exceptionType={}",
task.id(), task.taskType(), task.attemptCount(),
failure.errorCode(), failure.traceId(), failure.retryable(),
throwable.getClass().getName()
);
// 3. 通知监听器
notifyFailure(task, failure);
if (metrics != null) metrics.recordFailure(task.taskType());
// 4. 决定重试或失败
boolean retry = failure.retryable() &&
task.attemptCount() < task.maxAttempts();
if (retry) {
// 调度重试
Duration retryDelay = calculateRetryDelay(baseDelay, task.attemptCount());
store.scheduleRetry(
task.id(),
workerId,
failure,
clock.instant().plus(retryDelay)
);
} else {
// 标记失败
store.markFailed(task.id(), workerId, failure, clock.instant());
// 移入死信队列
if (deadLetterQueueService != null) {
try {
AsyncTask failedTask = store.findById(task.id()).orElse(task);
String deadLetterId = deadLetterQueueService.moveToDeadLetter(
failedTask, failure
);
log.info("任务已移入死信队列, taskId={}, deadLetterId={}",
task.id(), deadLetterId);
} catch (Exception ex) {
log.error("移入死信队列失败, taskId={}", task.id(), ex);
}
}
}
}2. 重试策略
指数退避 + 随机抖动
private Duration calculateRetryDelay(Duration baseDelay, int attempt) {
// 1. 指数退避
double multiplier = Math.pow(retryBackoffMultiplier, Math.max(0, attempt - 1));
long baseMillis = (long) Math.min(
Long.MAX_VALUE,
baseDelay.toMillis() * multiplier
);
// 2. 限制最大延迟
baseMillis = Math.min(maxRetryDelay.toMillis(), baseMillis);
// 3. 随机抖动(±20%)
double jitter = retryJitterRatio;
long jittered = (long) (baseMillis * (1.0 +
ThreadLocalRandom.current().nextDouble(-jitter, jitter)));
// 4. 确保非负且不超过上限
return Duration.ofMillis(
Math.max(0, Math.min(maxRetryDelay.toMillis(), jittered))
);
}示例(baseDelay=30s, multiplier=2.0, jitter=0.2):
| 尝试次数 | 计算延迟 | 实际延迟(加抖动) |
|---|---|---|
| 1 | 30s | 24s ~ 36s |
| 2 | 60s | 48s ~ 72s |
| 3 | 120s | 96s ~ 144s |
| 4 | 240s | 192s ~ 288s |
| 5 | 480s(达到上限) | 384s ~ 576s |
为什么需要抖动?
- 避免大量失败任务在同一时刻重试(雪崩效应)
- 分散数据库压力,提高系统稳定性
并发控制
1. 三级并发限制
async-task:
limits:
default-global-concurrency: 20 # 全局最大并发
default-type-concurrency: 5 # 单类型最大并发
worker:
local-max-concurrency: 4 # 单实例最大并发并发检查流程:
public Optional<AsyncTask> claimNext(TaskClaimRequest request) {
// 第1层:检查全局并发
long globalRunning = countRunningTasks();
if (globalRunning >= globalMaxConcurrency) {
return Optional.empty(); // 全局并发已满
}
// 第2层:检查类型并发
for (String taskType : request.taskTypes()) {
long typeRunning = countRunningTasksByType(taskType);
if (typeRunning >= typeMaxConcurrency) {
continue; // 此类型并发已满,尝试下一个
}
// 第3层:本地并发(Semaphore 控制)
// 在 poll() 方法中已通过 tryAcquire() 检查
// 领取该类型的任务
return claimNextOfType(taskType, request);
}
return Optional.empty();
}2. 动态调整并发
// 运行时调整全局并发(无需重启)
policyService.save(new AsyncTaskPolicy(
AsyncTaskPolicy.GLOBAL_POLICY_KEY,
null,
30, // 提升至 30 并发
3,
Duration.ofSeconds(30),
true
));
// 调整特定类型并发
policyService.save(new AsyncTaskPolicy(
AsyncTaskPolicy.typePolicyKey("IMPORT.MEDICAL"),
"IMPORT.MEDICAL",
2, // 降低至 2 并发
3,
Duration.ofSeconds(30),
true
));生效时机: 下一次 claimNext() 调用时立即生效。
优雅停机
@Override
public void stop() {
if (!running.compareAndSet(true, false)) {
return;
}
// 1. 停止接受新任务
scheduler.shutdown();
log.info("Worker 正在停止, workerId={}, activeTaskCount={}",
workerId, activeExecutions.size());
// 2. 停止执行线程池(不接受新任务)
executor.shutdown();
try {
// 3. 等待活跃任务完成(最多等待 leaseTimeout)
if (!executor.awaitTermination(
leaseTimeout.toMillis(),
TimeUnit.MILLISECONDS
)) {
// 4. 超时后强制停止
executor.shutdownNow();
log.warn("部分任务未完成即停止,将由其他 Worker 接管");
} else {
log.info("所有任务已完成,Worker 优雅停止");
}
} catch (InterruptedException ex) {
executor.shutdownNow();
Thread.currentThread().interrupt();
}
}停机流程图:
graph TD
A[收到停止信号] --> B[停止接受新任务]
B --> C[停止调度线程池]
C --> D[停止执行线程池]
D --> E[等待活跃任务完成]
E --> F{等待结果}
F -->|全部完成| G[优雅停止成功]
F -->|超时| H[强制停止]
H --> I[任务由其他 Worker 接管]监控指标
核心指标
// 1. 活跃任务数
metrics.activeIncrement();
metrics.activeDecrement();
// 2. 队列等待时长
metrics.recordQueueWait(Duration.between(createdAt, claimedAt));
// 3. 执行时长
metrics.recordExecution(taskType, executionDuration);
// 4. 成功/失败计数
metrics.completed();
metrics.recordFailure(taskType);
// 5. 超时计数
metrics.recordTimeout();
// 6. 租约续期失败
metrics.recordLeaseRenewalFailure();
// 7. 恢复统计
metrics.recordRecovery(requeuedCount, failedCount, cancelledCount);Prometheus 查询
# 活跃任务数
async_task_active_count
# 队列等待时长 P99
histogram_quantile(0.99,
rate(async_task_queue_wait_seconds_bucket[5m])
)
# 任务执行时长 P99(按类型)
histogram_quantile(0.99,
rate(async_task_execution_seconds_bucket{task_type="IMPORT.MEDICAL"}[5m])
)
# 任务成功率
sum(rate(async_task_completed_total[5m])) /
sum(rate(async_task_submitted_total[5m]))
# 租约续期失败率
rate(async_task_lease_renewal_failures_total[5m]) /
rate(async_task_lease_renewal_total[5m])最佳实践
1. Handler 实现规范
✅ 正确示例:
@Override
public TaskExecutionResult execute(TaskExecutionContext context,
ImportPayload payload) {
long total = countRows(payload);
long processed = 0;
// 分批处理
for (Batch batch : loadBatches(payload)) {
// ✅ 定期检查取消
context.checkCancellation();
// ✅ 幂等处理
if (!isBatchProcessed(context.taskId(), batch.id())) {
processBatch(batch);
markBatchProcessed(context.taskId(), batch.id());
}
processed += batch.size();
// ✅ 及时上报进度
context.reportProgress(new AsyncTaskProgress(
total, processed, processed, 0, 0
));
}
return TaskExecutionResult.success();
}❌ 错误示例:
@Override
public TaskExecutionResult execute(TaskExecutionContext context,
ImportPayload payload) {
// ❌ 长事务持有数据库连接
transactionTemplate.execute(status -> {
for (int i = 0; i < 10000; i++) {
processRow(i);
}
return null;
});
// ❌ 不检查取消标志
// ❌ 不上报进度
// ❌ 不做幂等处理
return TaskExecutionResult.success();
}2. 参数调优建议
| 场景 | poll-interval | heartbeat-interval | lease-timeout | local-max-concurrency |
|---|---|---|---|---|
| 开发环境 | 1s | 5s | 30s | 2 |
| 生产环境 | 1s | 10s | 60s | 4-8 |
| 长任务 | 2s | 20s | 120s | 2-4 |
| 高吞吐 | 500ms | 5s | 60s | 10-20 |
3. 故障排查
问题1: 任务一直排队不执行
排查步骤:
- 检查 Worker 是否启动(
worker.enabled=true) - 检查是否注册了 Handler(
handlerRegistry.registeredTaskTypes()) - 检查全局/类型并发是否已满
- 检查任务的
available_at是否未到达 - 检查任务的
deadline_at是否已过期
问题2: 任务反复进入 RECOVERING
排查步骤:
- 检查
lease-timeout是否过短 - 检查
heartbeat-interval是否过长 - 检查数据库连接是否稳定
- 检查 Worker 线程池是否过载
- 检查系统时钟是否同步(NTP)
问题3: 任务被重复执行
排查步骤:
- 检查 Handler 是否实现幂等
- 检查租约续期成功率
- 检查 Worker 是否频繁重启
- 检查数据库事务是否正确提交
总结
AsyncTaskWorker 是异步任务框架的心脏,通过精心设计的调度策略和容错机制,实现了:
- ✅ 高效调度:轮询 + 批量领取,减少数据库往返
- ✅ 可靠执行:租约机制 + 心跳续约,防止任务丢失
- ✅ 优雅容错:失败分类 + 指数退避 + 死信队列
- ✅ 平滑扩缩容:多实例自动负载均衡,无需协调
- ✅ 监控完善:指标暴露 + 健康检查,运维友好
在下一篇文章中,我们将通过实际案例,展示如何接入异步任务框架,从需求分析到上线的完整过程。
作者: 无声源语架构团队
发布日期: 2026-08-27
相关文章: