构建企业级异步任务框架:从零到一的架构设计之旅

前言

在现代分布式系统中,异步任务处理是一个常见且重要的需求。无论是大文件导入导出、批量数据处理,还是定时报表生成,都需要一个稳定、可靠的异步任务框架来支撑。本文将深入介绍我们自研的 tudicloud-spring-boot-starter-async-task 框架,分享从架构设计到实现细节的完整思路。

为什么要自研异步任务框架?

在开发医疗系统的过程中,我们遇到了许多典型的异步处理场景:

  • 大批量数据导入:医疗机构需要导入数万条患者数据,同步处理会阻塞用户操作
  • 复杂报表生成:统计报表计算耗时长,需要异步执行并通知用户下载
  • 定时任务调度:定期生成数据分析报告、清理过期数据
  • 外部系统集成:调用第三方医保接口,需要处理超时和重试

我们评估了业界常见的解决方案:

方案优势劣势是否采用
Spring @Async简单易用,零配置无持久化,重启丢失任务
Quartz成熟稳定,功能丰富主要面向定时任务,缺少任务状态管理
消息队列(RabbitMQ/Kafka)高吞吐,解耦强引入额外中间件,运维复杂度高
分布式调度(XXL-Job/ElasticJob)企业级调度能力需要独立部署调度中心,偏重定时任务

最终决定自研的原因:

  1. 轻量级依赖:只依赖已有的 MySQL 数据库,无需额外中间件
  2. 业务深度集成:与现有权限、审计、监控体系无缝集成
  3. 细粒度控制:支持任务级配额、优先级、死线(Deadline)等医疗业务特有需求
  4. 租约机制:基于数据库实现分布式锁,天然支持多实例部署

核心架构设计

1. 整体架构图

┌─────────────────────────────────────────────────────────────────┐
│                         业务应用层                               │
│  ┌──────────────┐  ┌──────────────┐  ┌──────────────┐         │
│  │ 导入Controller│  │ 导出Controller│  │ 报表Controller│         │
│  └──────┬───────┘  └──────┬───────┘  └──────┬───────┘         │
│         │                  │                  │                  │
│         └──────────────────┴──────────────────┘                 │
│                            ▼                                     │
│  ┌─────────────────────────────────────────────────────────┐   │
│  │              AsyncTaskService(任务提交服务)            │   │
│  └─────────────────────────────────────────────────────────┘   │
└─────────────────────────────────────────────────────────────────┘
                               │
                               ▼
┌─────────────────────────────────────────────────────────────────┐
│                      异步任务框架核心层                          │
│                                                                   │
│  ┌───────────────────┐         ┌──────────────────┐            │
│  │ AsyncTaskWorker   │◄────────┤ HandlerRegistry  │            │
│  │  (任务执行器)     │         │  (处理器注册表)  │            │
│  └─────────┬─────────┘         └──────────────────┘            │
│            │                                                     │
│            │  ┌──────────────────────────────────┐             │
│            └──┤   AsyncTaskStore (存储层)        │             │
│               │  - 任务领取(乐观锁)             │             │
│               │  - 租约续期                       │             │
│               │  - 状态流转                       │             │
│               └──────────────┬───────────────────┘             │
└────────────────────────────────────────────────────────────────┘
                               │
                               ▼
┌─────────────────────────────────────────────────────────────────┐
│                        MySQL 数据库层                            │
│  ┌──────────────────┐  ┌─────────────────┐                     │
│  │ system_async_task│  │ async_task_policy│                     │
│  │  (任务主表)      │  │  (并发策略表)    │                     │
│  └──────────────────┘  └─────────────────┘                     │
│  ┌──────────────────┐  ┌─────────────────┐                     │
│  │ async_task_attempt│ │async_task_failure│                     │
│  │  (执行记录)      │  │  (失败明细)      │                     │
│  └──────────────────┘  └─────────────────┘                     │
└─────────────────────────────────────────────────────────────────┘

2. 核心设计原则

2.1 数据库作为唯一真相源

与常见的内存队列 + 数据库持久化方案不同,我们选择数据库作为唯一的队列和状态管理器

优势:

  • ✅ 任务永不丢失,重启自动恢复
  • ✅ 天然支持分布式,多实例无需协调
  • ✅ 强一致性保证,无数据不一致风险
  • ✅ 简化运维,无需管理消息队列

劣势与应对:

  • ⚠️ 数据库性能瓶颈 → 通过索引优化 + 批量操作缓解
  • ⚠️ 高频轮询压力 → 可配置轮询间隔,支持按需调整

2.2 租约机制(Lease)

借鉴 Chubby/ZooKeeper 的租约思想,实现分布式环境下的任务独占执行:

// 领取任务时设置租约
UPDATE system_async_task 
SET status = 'RUNNING', 
    lease_owner = 'worker-123', 
    lease_until = NOW() + INTERVAL 60 SECOND
WHERE id = ? AND status = 'QUEUED';

// 心跳续约
UPDATE system_async_task 
SET lease_until = NOW() + INTERVAL 60 SECOND
WHERE id = ? AND lease_owner = 'worker-123';

// 租约过期恢复
UPDATE system_async_task 
SET status = 'QUEUED', 
    lease_owner = NULL
WHERE status = 'RUNNING' AND lease_until < NOW();

租约机制解决的问题:

  1. 防止重复执行:只有租约持有者能更新任务状态
  2. 自动故障恢复:Worker 宕机后租约自然过期,其他 Worker 可接管
  3. 优雅停机:停机时停止续约,任务会被其他实例接管

2.3 状态机驱动

任务状态流转遵循严格的状态机设计:

CREATED ──→ QUEUED ──→ RUNNING ──→ SUCCESS
                │          │
                │          ├─→ PARTIAL_SUCCESS
                │          │
                │          ├─→ RETRY_WAIT ──→ QUEUED
                │          │
                │          ├─→ FAILED
                │          │
                │          └─→ CANCEL_REQUESTED ──→ CANCELLED
                │
                └─→ CANCELLED

状态机约束:

  • 终态任务(SUCCESS/FAILED/CANCELLED)不可再变更
  • 状态转换必须携带租约校验,防止并发冲突
  • 每次状态变更记录 attempt(执行记录)用于审计

3. 核心组件职责

AsyncTaskService - 任务提交入口

@Service
public class ImportService {
    @Autowired
    private AsyncTaskService asyncTaskService;
    
    public String submitImportTask(ImportRequest request) {
        // 1. 业务校验
        validateRequest(request);
        
        // 2. 构造任务
        AsyncTaskSubmission submission = new AsyncTaskSubmission(
            "IMPORT.MEDICAL",           // 任务类型
            new ImportPayload(...),      // 业务参数
            currentUserId(),             // 所属主体
            generateIdempotencyKey(),    // 幂等键
            3                            // 最大重试次数
        );
        
        // 3. 提交任务
        AsyncTask task = asyncTaskService.submit(submission);
        return task.id();
    }
}

职责:

  • ✅ 参数校验与序列化
  • ✅ 幂等性保证(相同幂等键返回已有任务)
  • ✅ 队列配额检查(全局/Owner/Type 级别)
  • ✅ 持久化任务到数据库

AsyncTaskWorker - 任务执行引擎

Worker 是框架的心脏,负责任务的完整生命周期管理:

// 简化的 Worker 执行流程
public class AsyncTaskWorker {
    
    // 1. 定期轮询任务
    @Scheduled(fixedDelay = 1000)
    private void poll() {
        if (hasAvailableSlot()) {
            AsyncTask task = store.claimNext(claimRequest);
            if (task != null) {
                executor.submit(() -> execute(task));
            }
        }
    }
    
    // 2. 执行任务
    private void execute(AsyncTask task) {
        try {
            // 查找并调用 Handler
            AsyncTaskHandler handler = handlerRegistry.get(task.taskType());
            TaskExecutionResult result = handler.execute(context, payload);
            
            // 标记成功
            store.complete(task.id(), workerId, SUCCESS, result.artifactId());
        } catch (Exception ex) {
            // 失败分类与重试
            TaskFailure failure = failureClassifier.classify(task, ex);
            if (failure.retryable()) {
                store.scheduleRetry(task.id(), workerId, failure);
            } else {
                store.markFailed(task.id(), workerId, failure);
            }
        }
    }
    
    // 3. 定期续约
    @Scheduled(fixedDelay = 10000)
    private void heartbeat() {
        for (ActiveExecution active : activeExecutions.values()) {
            store.renewLease(active.taskId, workerId, newLeaseUntil);
        }
    }
    
    // 4. 恢复僵尸任务
    @Scheduled(fixedDelay = 60000)
    private void recover() {
        store.recoverExpiredLeases(Instant.now(), batchSize);
    }
}

职责:

  • ✅ 任务领取(带并发控制)
  • ✅ 调用业务 Handler 执行
  • ✅ 租约管理(续约 + 超时检测)
  • ✅ 异常处理与重试决策
  • ✅ 僵尸任务恢复

AsyncTaskHandler - 业务处理器

业务模块实现 Handler 接入异步任务:

@Component
public class MedicalImportHandler implements AsyncTaskHandler<ImportPayload> {
    
    @Override
    public String taskType() {
        return "IMPORT.MEDICAL";
    }
    
    @Override
    public Class<ImportPayload> payloadType() {
        return ImportPayload.class;
    }
    
    @Override
    public TaskExecutionResult execute(TaskExecutionContext context, 
                                       ImportPayload payload) throws Exception {
        long total = countRows(payload.fileId());
        long success = 0, failed = 0;
        
        // 分批处理
        for (Batch batch : loadBatches(payload)) {
            // 检查取消标志
            context.checkCancellation();
            
            // 业务逻辑(幂等处理)
            BatchResult result = processIdempotently(batch, context.taskId());
            success += result.successCount();
            failed += result.failedCount();
            
            // 上报进度
            context.reportProgress(new AsyncTaskProgress(
                total, success + failed, success, failed, 0
            ));
        }
        
        return failed == 0 
            ? TaskExecutionResult.success()
            : TaskExecutionResult.partialSuccess(null);
    }
}

Handler 接入要求:

  • ✅ 支持至少一次执行语义(租约恢复可能导致重复执行)
  • ✅ 业务写入必须幂等(使用任务 ID 或业务唯一键)
  • ✅ 长任务分批提交事务,不持有长事务
  • ✅ 定期调用 checkCancellation() 响应取消请求
  • ✅ 定期上报 reportProgress() 让用户感知进度

4. 并发控制策略

框架支持三级并发控制:

async-task:
  limits:
    default-global-concurrency: 8    # 全局最大并发(跨所有任务类型)
    default-type-concurrency: 4      # 单任务类型最大并发
  worker:
    local-max-concurrency: 4         # 单实例最大并发

并发控制实现:

-- 领取任务时检查并发限制
SELECT COUNT(*) FROM system_async_task 
WHERE status = 'RUNNING';  -- 检查全局并发

SELECT COUNT(*) FROM system_async_task 
WHERE status = 'RUNNING' AND task_type = 'IMPORT.MEDICAL';  -- 检查类型并发

-- 满足条件时领取
UPDATE system_async_task 
SET status = 'RUNNING', lease_owner = ?, lease_until = ?
WHERE id = ? AND status = 'QUEUED'
LIMIT 1;

动态策略调整:

// 运行时调整并发(无需重启)
policyService.save(new AsyncTaskPolicy(
    AsyncTaskPolicy.typePolicyKey("IMPORT.MEDICAL"),
    "IMPORT.MEDICAL",
    2,           // 降低至 2 并发
    3,           // 最大重试 3 次
    Duration.ofSeconds(30),  // 重试延迟 30s
    true         // 启用
));

关键特性

1. 死信队列

超过最大重试次数的任务自动移入死信队列:

// Worker 自动归档
if (!failure.retryable() || task.attemptCount() >= task.maxAttempts()) {
    store.markFailed(task.id(), workerId, failure);
    deadLetterQueueService.moveToDeadLetter(task, failure);
}

// 手动重放
deadLetterQueueService.replayBatch(deadLetterIds);

2. 批量操作

高性能批量提交:

List<AsyncTaskSubmission> submissions = users.stream()
    .map(user -> new AsyncTaskSubmission("EMAIL.SEND", payload, ownerId, key, 3))
    .toList();

batchAsyncTaskService.submitBatch(submissions);

3. 监控与健康检查

集成 Micrometer 和 Actuator:

management:
  endpoints:
    web:
      exposure:
        include: health, metrics

  metrics:
    tags:
      application: medical-system

暴露的指标:

  • async_task_submitted_total - 任务提交总数
  • async_task_completed_total{status="SUCCESS"} - 成功任务数
  • async_task_duration_seconds - 任务执行时长
  • async_task_queue_size - 队列深度
  • async_task_lease_renewal_failures - 租约续期失败数

性能表现

在我们的医疗系统生产环境中(4 核 8G 实例 × 3):

指标数值
任务吞吐~200 任务/分钟
平均延迟队列等待 < 5s
P99 执行时长< 120s (导入任务)
租约续期成功率> 99.9%
僵尸任务恢复< 60s (租约超时后)

下一步计划

本文介绍了异步任务框架的整体架构和核心设计。在后续文章中,我们将深入探讨:

  1. 租约机制与分布式协调:如何利用数据库实现分布式锁
  2. 任务执行引擎详解:Worker 的调度策略与容错设计
  3. 业务接入实战:从需求到上线的完整案例
  4. 监控与运维:指标体系、告警规则与故障排查

总结

构建一个企业级异步任务框架需要在简单性完备性之间找到平衡。我们的设计哲学是:

  • 依赖最小化:只依赖 MySQL,无需额外中间件
  • 设计可扩展:支持自定义 Handler、失败分类器、生命周期监听器
  • 运维友好:支持动态配置、健康检查、指标暴露
  • 容错优先:租约机制 + 乐观锁 + 自动恢复

希望本文能为正在构建类似系统的你提供参考。欢迎在评论区分享你的经验和问题!


作者: 无声源语架构团队
发布日期: 2026-08-25

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