深入理解租约机制:用数据库实现分布式任务调度

引言

在分布式系统中,如何保证同一个任务只被执行一次?如何在 Worker 节点宕机后自动恢复任务?传统方案通常依赖 Redis 分布式锁或 ZooKeeper 协调服务。但在 tudicloud-async-task​ 框架中,我们选择了一条不同的路:基于数据库的租约机制

本文将深入剖析租约机制的设计思想、实现细节以及如何解决分布式环境下的各种边界问题。

什么是租约(Lease)?

租约是分布式系统中的经典概念,最早由 Gray 和 Cheriton 在 1989 年提出。核心思想是:

资源的独占权是有时间期限的。持有者必须定期续约,否则租约自然过期,其他节点可以接管。

租约 vs 传统锁

特性传统锁租约
持有时长无限期(需要主动释放)有期限(到期自动失效)
故障恢复持有者宕机后锁永久占用持有者宕机后自动释放
实现复杂度需要检测节点存活无需检测,时间自然流逝
适用场景短期临界区保护长时间资源独占

在异步任务场景中,一个任务可能执行几秒到几分钟,使用传统锁会面临:

  • ❌ Worker 宕机导致任务永久锁定
  • ❌ 需要复杂的心跳检测和故障转移逻辑
  • ❌ 锁超时难以设置(过短误释放,过长影响恢复速度)

租约机制完美解决这些问题:

  • ✅ 自动过期,无需检测节点存活
  • ✅ 定期续约,确保执行中的任务不被抢占
  • ✅ 租约时长可配置,平衡性能与恢复速度

数据库租约实现原理

1. 数据模型

任务表中增加租约相关字段:

CREATE TABLE system_async_task (
    id BIGINT PRIMARY KEY,
    status VARCHAR(20) NOT NULL,  -- QUEUED, RUNNING, SUCCESS, FAILED, ...
    task_type VARCHAR(100) NOT NULL,
    
    -- 租约字段
    lease_owner VARCHAR(255),      -- 租约持有者(Worker ID)
    lease_until DATETIME(3),       -- 租约到期时间(毫秒精度)
    
    -- 其他字段...
    create_time DATETIME NOT NULL,
    update_time DATETIME NOT NULL,
    INDEX idx_lease_expire (status, lease_until)  -- 恢复查询索引
) ENGINE=InnoDB;

2. 核心操作

2.1 任务领取(Claim)

Worker 领取任务时,原子地设置租约:

public Optional<AsyncTask> claimNext(TaskClaimRequest request) {
    return transactionTemplate.execute(status -> {
        // 1. 检查全局并发
        long runningCount = countRunningTasks();
        if (runningCount >= globalMaxConcurrency) {
            return Optional.empty();
        }
        
        // 2. 检查类型并发
        long typeRunningCount = countRunningTasks(taskType);
        if (typeRunningCount >= typeMaxConcurrency) {
            return Optional.empty();
        }
        
        // 3. 领取任务(乐观锁 + 租约)
        Instant leaseUntil = Instant.now().plus(leaseTimeout);
        int updated = jdbcTemplate.update(
            "UPDATE system_async_task " +
            "SET status = 'RUNNING', " +
            "    lease_owner = ?, " +
            "    lease_until = ?, " +
            "    update_time = NOW() " +
            "WHERE id = (SELECT id FROM system_async_task " +
            "            WHERE status = 'QUEUED' " +
            "              AND task_type IN (?) " +
            "              AND (available_at IS NULL OR available_at <= NOW()) " +
            "            ORDER BY priority DESC, available_at ASC, create_time ASC " +
            "            LIMIT 1 FOR UPDATE SKIP LOCKED)",
            workerId, leaseUntil, taskTypes
        );
        
        if (updated == 0) {
            return Optional.empty();
        }
        
        return findByLeaseOwner(workerId).stream().findFirst();
    });
}

关键设计点:

  1. FOR UPDATE SKIP LOCKED:跳过已被其他事务锁定的行,避免等待
  2. 子查询 + 更新:一次性完成查询和状态变更,避免竞争
  3. 优先级排序:支持紧急任务优先执行
  4. 并发检查:在同一事务中检查并发限制,保证一致性

2.2 租约续期(Renew)

Worker 定期续约,延长租约到期时间:

@Scheduled(fixedDelay = 10000)  // 每 10 秒续约一次
private void heartbeat() {
    Instant newLeaseUntil = clock.instant().plus(leaseTimeout);
    
    // 批量续约所有活跃任务
    Set<String> renewed = store.renewLeases(
        activeExecutions.keySet(),
        workerId,
        newLeaseUntil
    );
    
    // 检查续约失败的任务
    for (String taskId : activeExecutions.keySet()) {
        if (!renewed.contains(taskId)) {
            // 续约失败,可能已被其他 Worker 接管
            log.error("租约续期失败,中断本地执行, taskId={}", taskId);
            activeExecutions.get(taskId).cancel();
        }
    }
}

SQL 实现(批量优化):

UPDATE system_async_task
SET lease_until = ?,
    update_time = NOW()
WHERE id IN (?, ?, ?)            -- 批量更新
  AND lease_owner = ?             -- 只能续自己的租约
  AND status = 'RUNNING'          -- 只有运行中的任务才需要续约
  AND lease_until > NOW();        -- 租约仍有效

为什么需要定期续约?

  • 证明 Worker 仍然存活且正在处理任务
  • 防止任务被其他 Worker 抢占
  • 租约时长(60s)远大于续约间隔(10s),留有充足缓冲

2.3 租约过期恢复(Recovery)

定期扫描租约过期的任务,重新排队或标记失败:

@Scheduled(fixedDelay = 60000)  // 每 60 秒恢复一次
private void recover() {
    RecoveryResult result = store.recoverExpiredLeases(
        Instant.now(),
        100  // 批次大小
    );
    
    log.info("租约恢复完成, requeued={}, failed={}, cancelled={}",
        result.requeuedCount(), 
        result.failedCount(), 
        result.cancelledCount()
    );
}

恢复逻辑(状态机驱动):

public RecoveryResult recoverExpiredLeases(Instant now, int limit) {
    return transactionTemplate.execute(status -> {
        // 1. 查询租约过期的任务
        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 ?",
            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);
    });
}

恢复策略:

graph TD
    A[检测到租约过期] --> B{检查任务状态}
    B -->|CANCEL_REQUESTED| C[标记为 CANCELLED]
    B -->|RUNNING| D{检查重试次数}
    D -->|未达上限| E[重新排队 QUEUED]
    D -->|达到上限| F[标记为 FAILED]

3. Worker 生命周期中的租约管理

启动阶段

@Override
public void start() {
    if (!running.compareAndSet(false, true)) {
        return;
    }
    
    // 1. 启动任务轮询线程
    scheduler.scheduleWithFixedDelay(
        this::poll,
        0,  // 立即开始
        pollInterval.toMillis(),
        TimeUnit.MILLISECONDS
    );
    
    // 2. 启动心跳续约线程
    scheduler.scheduleWithFixedDelay(
        this::heartbeat,
        heartbeatInterval.toMillis(),  // 延迟启动
        heartbeatInterval.toMillis(),
        TimeUnit.MILLISECONDS
    );
    
    // 3. 启动租约恢复线程
    scheduler.scheduleWithFixedDelay(
        this::recover,
        recoveryInterval.toMillis(),
        recoveryInterval.toMillis(),
        TimeUnit.MILLISECONDS
    );
    
    log.info("Worker 启动, workerId={}, leaseTimeout={}", 
        workerId, leaseTimeout);
}

运行阶段

private void execute(ActiveExecution active) {
    AsyncTask task = active.task;
    
    try {
        // 执行业务逻辑
        TaskExecutionResult result = invokeHandler(handler, context, payload);
        
        // 写入成功状态(带租约校验)
        boolean updated = store.complete(
            task.id(),
            workerId,    // 必须是租约持有者
            SUCCESS,
            result.artifactId()
        );
        
        if (!updated) {
            // 可能已失去租约
            if (store.isCancellationRequested(task.id())) {
                handleCancellation(task);
            } else {
                log.error("失去租约,无法写入终态, taskId={}", task.id());
            }
        }
    } catch (Exception ex) {
        handleFailure(task, ex);
    } finally {
        // 释放本地并发槽位
        activeExecutions.remove(task.id());
        localSlots.release();
    }
}

关键约束:

  • ✅ 只有租约持有者能更新任务状态
  • ✅ 状态更新失败不抛异常,由租约恢复兜底
  • ✅ 本地执行被中断不影响数据库状态

停止阶段

@Override
public void stop() {
    if (!running.compareAndSet(true, false)) {
        return;
    }
    
    // 1. 停止接受新任务
    scheduler.shutdown();
    
    // 2. 等待活跃任务完成(最多等待租约超时时间)
    log.info("等待活跃任务完成, count={}", activeExecutions.size());
    
    try {
        executor.shutdown();
        if (!executor.awaitTermination(
            leaseTimeout.toMillis(), 
            TimeUnit.MILLISECONDS
        )) {
            // 超时后强制停止
            executor.shutdownNow();
            log.warn("部分任务未完成即停止,将由其他 Worker 接管");
        }
    } catch (InterruptedException ex) {
        executor.shutdownNow();
        Thread.currentThread().interrupt();
    }
}

优雅停机策略:

  1. 停止接受新任务(停止轮询线程)
  2. 停止续约(租约自然过期)
  3. 等待活跃任务完成(最多等待 leaseTimeout)
  4. 超时后强制停止,任务由其他 Worker 接管

租约参数配置

推荐配置

async-task:
  worker:
    lease-timeout: 60s          # 租约超时时间
    heartbeat-interval: 10s     # 心跳间隔(续约)
    recovery-interval: 30s      # 恢复扫描间隔
    max-execution-time: 5m      # 单任务最长执行时间

参数调优指南

1. lease-timeout(租约超时时间)

作用: 租约的有效期,超过此时间未续约则视为过期。

取值建议:

  • 开发环境:30s(快速发现问题)
  • 生产环境:60s(平衡性能与恢复速度)
  • 长任务场景:120s(减少续约频率)

约束:

  • ✅ 必须 >= 2 × heartbeat-interval(避免续约时机与过期时间重叠)
  • ✅ 应 < max-execution-time(避免任务执行完毕但租约未过期)

影响:

  • ⬆️ 过大:Worker 宕机后恢复慢
  • ⬇️ 过小:网络抖动导致误判过期

2. heartbeat-interval(心跳间隔)

作用: Worker 续约的频率。

取值建议:

  • 高频心跳:5s(实时性要求高)
  • 标准配置:10s(推荐)
  • 低频心跳:20s(减少数据库压力)

约束:

  • ✅ 必须 < lease-timeout / 2(确保至少续约一次)
  • ✅ 建议 = lease-timeout / 4 ~ lease-timeout / 6(留有充足缓冲)

影响:

  • ⬆️ 过大:Worker 宕机后任务继续执行,可能产生重复副作用
  • ⬇️ 过小:心跳频繁,增加数据库负载

3. recovery-interval(恢复扫描间隔)

作用: 扫描租约过期任务的频率。

取值建议:

  • 快速恢复:30s
  • 标准配置:60s(推荐)
  • 低优先级:120s

约束:

  • ✅ 建议 >= lease-timeout(避免与正常执行重叠)
  • ✅ 建议 <= 2 × lease-timeout(控制恢复延迟)

影响:

  • ⬆️ 过大:任务恢复延迟增加
  • ⬇️ 过小:频繁扫描,增加数据库负载

配置示例对比

场景lease-timeoutheartbeat-intervalrecovery-interval说明
快速恢复30s5s30sWorker 宕机后 30s 内恢复
标准生产60s10s60s平衡性能与恢复速度
长任务120s20s120s减少心跳频率,适合耗时任务
高可用60s5s30s心跳频繁 + 快速恢复

边界问题处理

1. 续约失败场景

问题: Worker 网络抖动,续约失败但任务仍在执行。

处理策略:

Set<String> renewed = store.renewLeases(activeTaskIds, workerId, leaseUntil);

for (String taskId : activeTaskIds) {
    if (!renewed.contains(taskId)) {
        // 续约失败,检查任务状态
        AsyncTask current = store.findById(taskId).orElse(null);
        
        if (current != null && current.status().isTerminal()) {
            // 任务已完成,无需处理
            log.info("任务已完成,无需继续续约, taskId={}", taskId);
            continue;
        }
        
        // 任务仍在运行但续约失败,中断执行
        log.error("续约失败,中断任务执行, taskId={}", taskId);
        Future<?> future = activeExecutions.get(taskId).future;
        if (future != null) {
            future.cancel(true);  // 中断线程
        }
    }
}

设计原则:

  • ✅ 续约失败立即中断任务,降低重复执行风险
  • ✅ 中断前检查任务是否已完成,避免误杀
  • ✅ 中断后任务由租约恢复机制接管

2. 任务完成与租约过期竞争

问题: Worker 写入成功状态时,租约刚好过期。

时序图:

Worker A                    Database                   Worker B
   |                           |                           |
   |--execute(task-1)--------->|                           |
   |                           |                           |
   |                           | lease_until=10:00:00      |
   |                           |                           |
   |        (10:00:00 到达,租约过期)                      |
   |                           |                           |
   |                           |<---recover(task-1)--------|
   |                           | status=QUEUED             |
   |                           |                           |
   |--complete(task-1)-------->|                           |
   |   (租约校验失败)          |                           |
   |<----update failed---------|                           |

处理策略:

-- 完成状态写入必须校验租约
UPDATE system_async_task
SET status = 'SUCCESS',
    lease_owner = NULL,
    lease_until = NULL,
    finished_at = NOW()
WHERE id = ?
  AND lease_owner = ?          -- 必须是租约持有者
  AND lease_until > NOW()      -- 租约仍有效
  AND status = 'RUNNING';      -- 状态仍是运行中

如果更新失败:

boolean updated = store.complete(taskId, workerId, SUCCESS, artifactId);
if (!updated) {
    // 1. 检查是否被取消
    if (store.isCancellationRequested(taskId)) {
        handleCancellation(task);
        return;
    }
    
    // 2. 租约已过期,放弃写入
    log.error("失去租约,任务将由其他 Worker 接管, taskId={}", taskId);
    // 不抛异常,由租约恢复机制处理
}

设计原则:

  • ✅ 状态写入失败不抛异常,避免误记录失败
  • ✅ 任务会被租约恢复机制重新排队
  • ✅ Handler 必须支持幂等,容忍重复执行

3. 多 Worker 同时恢复

问题: 多个 Worker 同时扫描到同一批过期任务。

处理策略:

public RecoveryResult recoverExpiredLeases(Instant now, int limit) {
    return transactionTemplate.execute(status -> {
        // 使用 FOR UPDATE 锁定待恢复任务
        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
        );
        
        // 在同一事务中更新状态
        for (AsyncTask task : expired) {
            requeueOrFail(task);
        }
        
        return result;
    });
}

设计原则:

  • ✅ 使用 FOR UPDATE SKIP LOCKED 避免锁等待
  • ✅ 每个 Worker 处理不重叠的任务集
  • ✅ 恢复批次限制避免长事务

4. Worker 宕机瞬间的任务状态

问题: Worker 正在执行任务时突然宕机(断电、OOM、Kill -9)。

场景分析:

宕机时机数据库状态恢复结果
领取任务后,执行前status=RUNNING, lease_until=有效租约过期后重新排队,Handler 未执行
执行中途status=RUNNING, lease_until=有效租约过期后重新排队,可能产生部分副作用
执行完毕,写入状态前status=RUNNING, lease_until=有效租约过期后重新排队,重复执行
写入状态后status=SUCCESS, lease_owner=NULL正常完成,无需恢复

风险点:

  • ⚠️ 执行中途宕机,业务逻辑已产生部分副作用
  • ⚠️ 执行完毕但未写入状态,业务逻辑会被重复执行

解决方案:Handler 幂等性

@Override
public TaskExecutionResult execute(TaskExecutionContext context, 
                                   ImportPayload payload) {
    // 方案1: 使用任务 ID 作为幂等键
    String idempotencyKey = "import_" + context.taskId();
    
    // 方案2: 使用业务唯一键
    String businessKey = payload.requestId();
    
    // 检查是否已处理
    if (importRecordExists(idempotencyKey)) {
        log.info("任务已处理,跳过执行, key={}", idempotencyKey);
        return TaskExecutionResult.success();
    }
    
    // 执行业务逻辑
    processData(payload);
    
    // 记录幂等标记(与业务数据在同一事务中)
    saveImportRecord(idempotencyKey, payload);
    
    return TaskExecutionResult.success();
}

幂等设计原则:

  • ✅ 业务写入必须带幂等键(任务 ID 或业务唯一键)
  • ✅ 幂等检查与业务操作在同一事务中
  • ✅ 分批处理时,每批都要记录批次号
  • ✅ 外部 API 调用使用调用凭证去重

性能优化

1. 批量续约优化

问题: 单实例有 10 个活跃任务,每次心跳需要 10 次数据库往返。

优化前:

for (String taskId : activeTaskIds) {
    store.renewLease(taskId, workerId, leaseUntil);
}
// 10 个任务 = 10 次数据库往返

优化后:

// 一次性续约所有任务
Set<String> renewed = store.renewLeases(activeTaskIds, workerId, leaseUntil);

SQL 实现:

-- 批量续约(MyBatis foreach)
UPDATE system_async_task
SET lease_until = ?,
    update_time = NOW()
WHERE id IN
<foreach collection="taskIds" item="id" open="(" separator="," close=")">
    #{id}
</foreach>
  AND lease_owner = ?
  AND status = 'RUNNING';

性能提升:

  • 10 次往返 → 1 次往返
  • 心跳延迟从 ~100ms 降至 ~10ms
  • 数据库 QPS 减少 90%

2. 索引优化

核心查询场景:

-- 场景1: 领取任务
SELECT * FROM system_async_task
WHERE status = 'QUEUED'
  AND task_type IN (?)
  AND (available_at IS NULL OR available_at <= NOW())
ORDER BY priority DESC, available_at ASC, create_time ASC
LIMIT 1;

-- 场景2: 租约恢复
SELECT * FROM system_async_task
WHERE status IN ('RUNNING', 'CANCEL_REQUESTED')
  AND lease_until < NOW()
ORDER BY lease_until ASC
LIMIT 100;

-- 场景3: Owner 查询
SELECT * FROM system_async_task
WHERE owner_id = ?
  AND status IN (?)
ORDER BY create_time DESC
LIMIT 20;

索引设计:

-- 领取任务索引
CREATE INDEX idx_claim ON system_async_task (
    status, task_type, available_at, priority DESC, create_time
);

-- 租约恢复索引
CREATE INDEX idx_lease_expire ON system_async_task (
    status, lease_until
);

-- Owner 查询索引
CREATE INDEX idx_owner_status ON system_async_task (
    owner_id, status, create_time DESC
);

3. 恢复批次限制

问题: 大量任务同时租约过期,恢复线程长时间持有事务锁。

优化策略:

@Scheduled(fixedDelay = 60000)
private void recover() {
    int batchSize = 100;  // 限制批次大小
    int maxBatches = 10;  // 限制批次数
    
    for (int i = 0; i < maxBatches; i++) {
        RecoveryResult result = store.recoverExpiredLeases(
            Instant.now(), 
            batchSize
        );
        
        if (result.totalCount() < batchSize) {
            break;  // 没有更多过期任务
        }
        
        // 批次间短暂休息,释放数据库连接
        Thread.sleep(100);
    }
}

设计原则:

  • ✅ 单批次处理 100 条,避免长事务
  • ✅ 最多处理 10 批次,避免恢复线程占用过长
  • ✅ 批次间休息 100ms,释放数据库资源

监控与告警

关键指标

// 租约续期成功率
metrics.recordLeaseRenewal(taskId, success);

// 租约过期恢复统计
metrics.recordRecovery(requeuedCount, failedCount, cancelledCount);

// 心跳失败计数
metrics.heartbeatFailure();

Grafana 看板

# 租约续期成功率
sum(rate(async_task_lease_renewal_success[5m])) /
sum(rate(async_task_lease_renewal_total[5m]))

# 租约过期任务数
async_task_lease_expired_total

# 平均恢复延迟
avg(async_task_recovery_delay_seconds)

告警规则

# 续约成功率低于 99%
- alert: LeaseRenewalLowSuccessRate
  expr: |
    sum(rate(async_task_lease_renewal_success[5m])) /
    sum(rate(async_task_lease_renewal_total[5m])) < 0.99
  for: 5m
  annotations:
    summary: "租约续期成功率过低"
    
# 大量租约过期
- alert: HighLeaseExpiration
  expr: rate(async_task_lease_expired_total[5m]) > 10
  for: 5m
  annotations:
    summary: "租约过期任务数异常"

总结

租约机制是分布式异步任务框架的核心,通过巧妙利用时间和数据库事务,实现了:

  • 无需心跳检测:租约自动过期,无需检测节点存活
  • 自动故障恢复:Worker 宕机后租约自然失效,其他实例接管
  • 强一致性保证:基于数据库事务,避免任务被多次执行
  • 可配置性:支持灵活调整租约时长、心跳频率、恢复间隔

但租约机制也带来了新的挑战:

  • ⚠️ 必须幂等:Handler 必须支持重复执行
  • ⚠️ 参数调优:租约时长、心跳间隔需要根据场景调优
  • ⚠️ 监控必要:需要监控续约成功率、恢复延迟等指标

在下一篇文章中,我们将深入探讨 AsyncTaskWorker 的调度策略与容错设计,敬请期待!


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

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