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());
        }
    }
}

关键设计点:

  1. 非阻塞领取tryAcquire() 不阻塞,没有槽位立即返回
  2. 批量领取while 循环持续领取,直到没有任务或槽位用尽
  3. 类型过滤:只领取已注册 Handler 的任务类型
  4. 异常处理:线程池拒绝时不破坏任务状态,由租约恢复兜底

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 &lt;= #{now})
          AND (deadline_at IS NULL OR deadline_at &gt; #{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>

排序策略:

  1. priority DESC:优先级高的先执行
  2. available_at ASC:可用时间早的先执行(支持延迟任务)
  3. 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):

尝试次数计算延迟实际延迟(加抖动)
130s24s ~ 36s
260s48s ~ 72s
3120s96s ~ 144s
4240s192s ~ 288s
5480s(达到上限)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-intervalheartbeat-intervallease-timeoutlocal-max-concurrency
开发环境1s5s30s2
生产环境1s10s60s4-8
长任务2s20s120s2-4
高吞吐500ms5s60s10-20

3. 故障排查

问题1: 任务一直排队不执行

排查步骤:

  1. 检查 Worker 是否启动(worker.enabled=true
  2. 检查是否注册了 Handler(handlerRegistry.registeredTaskTypes()
  3. 检查全局/类型并发是否已满
  4. 检查任务的 available_at 是否未到达
  5. 检查任务的 deadline_at 是否已过期

问题2: 任务反复进入 RECOVERING

排查步骤:

  1. 检查 lease-timeout 是否过短
  2. 检查 heartbeat-interval 是否过长
  3. 检查数据库连接是否稳定
  4. 检查 Worker 线程池是否过载
  5. 检查系统时钟是否同步(NTP)

问题3: 任务被重复执行

排查步骤:

  1. 检查 Handler 是否实现幂等
  2. 检查租约续期成功率
  3. 检查 Worker 是否频繁重启
  4. 检查数据库事务是否正确提交

总结

AsyncTaskWorker 是异步任务框架的心脏,通过精心设计的调度策略和容错机制,实现了:

  • 高效调度:轮询 + 批量领取,减少数据库往返
  • 可靠执行:租约机制 + 心跳续约,防止任务丢失
  • 优雅容错:失败分类 + 指数退避 + 死信队列
  • 平滑扩缩容:多实例自动负载均衡,无需协调
  • 监控完善:指标暴露 + 健康检查,运维友好

在下一篇文章中,我们将通过实际案例,展示如何接入异步任务框架,从需求分析到上线的完整过程。


作者: 无声源语架构团队
发布日期: 2026-08-27
相关文章:

最后修改:2026 年 08 月 25 日
如果觉得我的文章对你有用,请随意赞赏