异步任务框架的监控、运维与故障排查实战

引言

一个企业级框架的价值不仅在于其功能完备性,更在于其可观测性可维护性。当系统在生产环境遇到问题时,能否快速定位根因、评估影响范围、制定恢复方案,是衡量框架成熟度的重要标准。

本文将深入探讨 tudicloud-async-task 框架的监控体系、运维工具和故障排查方法,帮助你构建一个"看得见、管得住、修得快"的异步任务系统。

监控体系设计

1. 四个黄金信号

根据 Google SRE 方法论,我们重点监控四个黄金信号:

1.1 延迟(Latency)

定义: 从任务提交到开始执行的等待时长 + 任务执行时长

监控指标:

// 队列等待时长
metrics.recordQueueWait(Duration.between(task.createdAt(), claimedAt));

// 任务执行时长
metrics.recordExecution(taskType, Duration.ofNanos(endNanos - startNanos));

Prometheus 查询:

# P50 队列等待时长
histogram_quantile(0.5, 
  rate(async_task_queue_wait_seconds_bucket[5m])
)

# P99 执行时长(按任务类型)
histogram_quantile(0.99, 
  rate(async_task_execution_seconds_bucket{task_type="IMPORT.MEDICAL"}[5m])
)

# 平均端到端延迟
avg(async_task_queue_wait_seconds + async_task_execution_seconds)

告警规则:

- alert: HighTaskQueueWait
  expr: |
    histogram_quantile(0.99, 
      rate(async_task_queue_wait_seconds_bucket[5m])
    ) > 60
  for: 5m
  annotations:
    summary: "任务队列等待时长 P99 超过 60 秒"
    description: "当前 P99 等待时长: {{ $value }}s"

1.2 流量(Traffic)

定义: 任务提交速率和完成速率

监控指标:

// 任务提交计数
metrics.recordSubmission(taskType);

// 任务完成计数
metrics.recordCompletion(taskType, status);

Prometheus 查询:

# 任务提交速率(按类型)
sum(rate(async_task_submitted_total[5m])) by (task_type)

# 任务完成速率(按状态)
sum(rate(async_task_completed_total[5m])) by (status)

# 任务积压数(提交速率 - 完成速率)
sum(rate(async_task_submitted_total[5m])) - 
sum(rate(async_task_completed_total[5m]))

告警规则:

- alert: TaskBacklogIncreasing
  expr: |
    (
      sum(rate(async_task_submitted_total[5m])) - 
      sum(rate(async_task_completed_total[5m]))
    ) > 10
  for: 10m
  annotations:
    summary: "任务积压持续增长"
    description: "积压速率: {{ $value }} 任务/秒"

1.3 错误(Errors)

定义: 任务失败率和错误类型分布

监控指标:

// 任务失败计数(按错误码)
metrics.recordFailure(taskType, errorCode);

// 重试计数
metrics.recordRetry(taskType);

// 死信队列计数
metrics.recordDeadLetter(taskType);

Prometheus 查询:

# 任务失败率
sum(rate(async_task_completed_total{status="FAILED"}[5m])) /
sum(rate(async_task_completed_total[5m]))

# 失败任务数(按错误码)
sum(rate(async_task_failures_total[5m])) by (error_code)

# 死信队列大小
async_task_dead_letter_queue_size

告警规则:

- alert: HighTaskFailureRate
  expr: |
    sum(rate(async_task_completed_total{status="FAILED"}[5m])) /
    sum(rate(async_task_completed_total[5m])) > 0.1
  for: 5m
  annotations:
    summary: "任务失败率超过 10%"
    
- alert: DeadLetterQueueGrowing
  expr: rate(async_task_dead_letter_queue_size[10m]) > 5
  for: 5m
  annotations:
    summary: "死信队列持续增长"

1.4 饱和度(Saturation)

定义: 系统资源使用情况

监控指标:

// 活跃任务数
metrics.activeIncrement();
metrics.activeDecrement();

// 并发槽位使用率
metrics.recordConcurrencyUtilization(
    activeCount, localMaxConcurrency
);

Prometheus 查询:

# 活跃任务数
async_task_active_count

# 并发槽位使用率
async_task_active_count / async_task_local_max_concurrency

# 数据库连接池使用率
hikaricp_connections_active / hikaricp_connections_max

告警规则:

- alert: HighConcurrencySaturation
  expr: |
    async_task_active_count / 
    async_task_local_max_concurrency > 0.9
  for: 5m
  annotations:
    summary: "并发槽位使用率超过 90%"
    
- alert: DatabaseConnectionPoolExhausted
  expr: |
    hikaricp_connections_active / 
    hikaricp_connections_max > 0.9
  for: 2m
  annotations:
    summary: "数据库连接池接近耗尽"

2. 租约健康度监控

租约机制是框架的核心,必须重点监控:

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

// 租约续期失败计数
metrics.recordLeaseRenewalFailure();

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

Prometheus 查询:

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

# 租约过期任务恢复速率
sum(rate(async_task_lease_expired_recovered_total[5m]))

# 平均租约持有时长
avg(async_task_lease_duration_seconds)

Grafana 面板配置:

{
  "title": "租约健康度",
  "panels": [
    {
      "title": "租约续期成功率",
      "targets": [
        {
          "expr": "sum(rate(async_task_lease_renewal_success_total[5m])) / sum(rate(async_task_lease_renewal_total[5m]))"
        }
      ],
      "alert": {
        "conditions": [
          {
            "evaluator": {
              "params": [0.99],
              "type": "lt"
            }
          }
        ]
      }
    }
  ]
}

3. 业务指标监控

除了框架级指标,还需要监控业务维度的指标:

@Component
public class BusinessMetrics {
    private final MeterRegistry registry;
    
    // 按任务类型统计成功率
    public void recordBusinessSuccess(String taskType, boolean success) {
        registry.counter(
            "business.task.result",
            "task_type", taskType,
            "result", success ? "success" : "failed"
        ).increment();
    }
    
    // 按用户统计任务数
    public void recordUserTask(String userId, String taskType) {
        registry.counter(
            "business.user.task",
            "user_id", userId,
            "task_type", taskType
        ).increment();
    }
    
    // 记录业务处理速率(如导入速率)
    public void recordProcessingRate(String taskType, long itemsProcessed) {
        registry.counter(
            "business.items.processed",
            "task_type", taskType
        ).increment(itemsProcessed);
    }
}

业务看板示例:

# 医保导入任务成功率
sum(rate(business_task_result{task_type="MEDICAL_INSURANCE.IMPORT", result="success"}[5m])) /
sum(rate(business_task_result{task_type="MEDICAL_INSURANCE.IMPORT"}[5m]))

# Top 10 活跃用户
topk(10, 
  sum(rate(business_user_task[1h])) by (user_id)
)

# 每小时导入患者数
sum(rate(business_items_processed{task_type="MEDICAL_INSURANCE.IMPORT"}[1h])) * 3600

健康检查设计

1. 控制面健康检查

@Component
public class AsyncTaskHealthIndicator extends AbstractHealthIndicator {
    
    private final AsyncTaskStore store;
    
    @Override
    protected void doHealthCheck(Health.Builder builder) throws Exception {
        try {
            // 1. 检查数据库连接
            store.verifySchema();
            
            // 2. 检查表结构
            long totalTasks = store.countAllTasks();
            
            // 3. 检查队列深度
            long queuedTasks = store.countQueuedTasks();
            
            // 4. 检查活跃任务数
            long runningTasks = store.countRunningTasks();
            
            builder.up()
                .withDetail("total_tasks", totalTasks)
                .withDetail("queued_tasks", queuedTasks)
                .withDetail("running_tasks", runningTasks);
                
            // 告警阈值检查
            if (queuedTasks > 1000) {
                builder.status("DEGRADED")
                    .withDetail("warning", "队列积压超过 1000 个任务");
            }
            
        } catch (Exception ex) {
            builder.down()
                .withException(ex)
                .withDetail("error", "数据库连接失败或表结构异常");
        }
    }
}

2. Worker 健康检查

@Component
public class AsyncTaskWorkerHealthIndicator extends AbstractHealthIndicator {
    
    private final AsyncTaskWorker worker;
    
    @Override
    protected void doHealthCheck(Health.Builder builder) {
        if (!worker.isRunning()) {
            builder.down()
                .withDetail("status", "Worker 未运行");
            return;
        }
        
        int activeCount = worker.getActiveTaskCount();
        int maxConcurrency = worker.getLocalMaxConcurrency();
        double utilization = (double) activeCount / maxConcurrency;
        
        builder.up()
            .withDetail("worker_id", worker.getWorkerId())
            .withDetail("active_tasks", activeCount)
            .withDetail("max_concurrency", maxConcurrency)
            .withDetail("utilization", String.format("%.2f%%", utilization * 100))
            .withDetail("registered_task_types", worker.getRegisteredTaskTypes());
        
        // 高负载告警
        if (utilization > 0.9) {
            builder.status("DEGRADED")
                .withDetail("warning", "Worker 负载超过 90%");
        }
    }
}

3. Kubernetes 健康探针

apiVersion: v1
kind: Pod
metadata:
  name: async-task-worker
spec:
  containers:
  - name: app
    image: medical-system:1.0.0
    
    # 存活探针(Worker 是否存活)
    livenessProbe:
      httpGet:
        path: /actuator/health/liveness
        port: 8080
      initialDelaySeconds: 30
      periodSeconds: 10
      timeoutSeconds: 5
      failureThreshold: 3
    
    # 就绪探针(Worker 是否就绪)
    readinessProbe:
      httpGet:
        path: /actuator/health/readiness
        port: 8080
      initialDelaySeconds: 10
      periodSeconds: 5
      timeoutSeconds: 3
      failureThreshold: 2
    
    # 启动探针(Worker 启动检查)
    startupProbe:
      httpGet:
        path: /actuator/health/startup
        port: 8080
      initialDelaySeconds: 0
      periodSeconds: 5
      timeoutSeconds: 3
      failureThreshold: 30

日志规范

1. 日志级别定义

级别使用场景示例
ERROR任务执行失败、租约续期失败、数据库异常租约续期失败,中断任务执行
WARN租约过期恢复、慢任务告警、配置异常任务执行超过慢任务阈值
INFO任务提交、开始执行、完成、取消异步任务执行完成
DEBUG领取任务、心跳续约、进度上报心跳续约成功

2. 日志格式规范

// ✅ 好的日志格式(结构化、可检索)
log.info("异步任务执行完成, taskId={}, taskType={}, attempt={}, status={}, duration={}ms",
    task.id(), task.taskType(), task.attemptCount(), task.status(), durationMs);

// ❌ 不好的日志格式(难以解析)
log.info("任务 " + task.id() + " 执行完成,状态:" + task.status());

3. 敏感信息脱敏

// ✅ 正确:不记录敏感信息
log.error("异步任务执行异常, taskId={}, taskType={}, errorCode={}, traceId={}",
    task.id(), task.taskType(), failure.errorCode(), failure.traceId());

// ❌ 错误:记录了完整异常消息(可能包含敏感数据)
log.error("异步任务执行异常, taskId={}, error={}", 
    task.id(), throwable.getMessage());

4. 日志聚合与检索

Elasticsearch 查询示例:

// 查询某个任务的完整生命周期
GET /logs-*/_search
{
  "query": {
    "bool": {
      "must": [
        {"match": {"message": "taskId=task-12345"}},
        {"range": {"@timestamp": {"gte": "now-1h"}}}
      ]
    }
  },
  "sort": [{"@timestamp": "asc"}]
}

// 统计失败任务的错误码分布
GET /logs-*/_search
{
  "query": {
    "match": {"level": "ERROR"}
  },
  "aggs": {
    "error_codes": {
      "terms": {"field": "errorCode.keyword"}
    }
  }
}

故障排查手册

1. 任务一直排队不执行

现象: 任务状态一直是 QUEUED​,从未变为 RUNNING

排查步骤:

# 1. 检查 Worker 是否启动
curl http://localhost:8080/actuator/health/asyncTaskWorker

# 2. 检查 Worker 是否注册了该任务类型
curl http://localhost:8080/actuator/metrics/async.task.registered.types

# 3. 检查全局并发是否已满
mysql> SELECT COUNT(*) FROM system_async_task WHERE status = 'RUNNING';

# 4. 检查类型并发是否已满
mysql> SELECT COUNT(*) FROM system_async_task 
       WHERE status = 'RUNNING' AND task_type = 'IMPORT.MEDICAL';

# 5. 检查任务是否被暂停
mysql> SELECT * FROM system_async_task WHERE id = 'task-12345';
-- 检查 status 字段是否为 'PAUSED'

# 6. 检查任务的 available_at 是否未到达
mysql> SELECT id, available_at, NOW() FROM system_async_task 
       WHERE id = 'task-12345';

# 7. 检查任务的 deadline_at 是否已过期
mysql> SELECT id, deadline_at, NOW() FROM system_async_task 
       WHERE id = 'task-12345';

# 8. 检查动态策略是否禁用
mysql> SELECT * FROM system_async_task_policy 
       WHERE policy_key = 'TYPE:IMPORT.MEDICAL';

常见原因与解决方案:

原因解决方案
Worker 未启动检查配置 worker.enabled=true,重启应用
未注册 Handler确认 Handler 类上有 @Component 注解
全局并发已满调高 default-global-concurrency 或等待任务完成
类型并发已满调高类型策略的 maxConcurrency
任务被暂停调用 asyncTaskService.resume(taskId)
available_at 未到达等待到达时间或修改任务
deadline_at 已过期任务已失效,需要重新提交
策略被禁用更新策略 enabled=true

2. 任务反复进入 RECOVERING

现象: 任务状态在 RUNNING​ 和 RECOVERING 之间反复切换。

排查步骤:

# 1. 检查租约配置
cat application.yml | grep -A 10 "worker:"

# 2. 检查租约续期成功率
curl http://localhost:8080/actuator/metrics/async.task.lease.renewal.success.rate

# 3. 检查数据库连接池
curl http://localhost:8080/actuator/metrics/hikaricp.connections.active

# 4. 检查系统时钟同步
ntpstat
timedatectl status

# 5. 查看 Worker 日志
tail -f /var/log/app.log | grep "租约续期失败"

# 6. 检查任务执行时长
mysql> SELECT id, create_time, update_time, 
       TIMESTAMPDIFF(SECOND, create_time, update_time) AS duration_sec
       FROM system_async_task 
       WHERE id = 'task-12345';

常见原因与解决方案:

原因解决方案
lease-timeout 过短调大至 60s 或更高
heartbeat-interval 过长调小至 lease-timeout / 4
数据库连接池耗尽增大连接池大小
数据库慢查询优化 SQL,增加索引
系统时钟不同步配置 NTP 同步
Worker 负载过高增加 Worker 实例或降低并发

3. 任务被重复执行

现象: 同一个任务的业务逻辑被执行了多次,产生重复数据。

排查步骤:

# 1. 检查 Handler 是否实现幂等
grep -r "idempotencyKey" src/main/java/

# 2. 检查任务的执行记录
mysql> SELECT * FROM system_async_task_attempt 
       WHERE task_id = 'task-12345' 
       ORDER BY started_at;

# 3. 检查租约续期失败日志
tail -f /var/log/app.log | grep "租约续期失败"

# 4. 检查 Worker 重启记录
kubectl logs -n medical-system async-task-worker --previous

# 5. 检查业务数据是否重复
mysql> SELECT batch_no, COUNT(*) FROM medical_insurance_import_detail 
       GROUP BY batch_no HAVING COUNT(*) > 1;

常见原因与解决方案:

原因解决方案
Handler 未实现幂等添加幂等检查逻辑
租约续期频繁失败参考"任务反复进入 RECOVERING"的解决方案
Worker 频繁重启排查 Worker 宕机原因
数据库事务未正确提交检查事务配置和异常处理

4. 任务执行超时

现象: 任务被标记为 FAILED​,错误码为 TIMEOUT

排查步骤:

# 1. 检查任务配置的超时时间
mysql> SELECT id, deadline_at, create_time FROM system_async_task 
       WHERE id = 'task-12345';

# 2. 检查 Worker 配置的最大执行时间
cat application.yml | grep "max-execution-time"

# 3. 检查任务实际执行时长
mysql> SELECT 
         a.task_id, 
         a.started_at, 
         a.finished_at,
         TIMESTAMPDIFF(SECOND, a.started_at, a.finished_at) AS duration_sec
       FROM system_async_task_attempt a
       WHERE a.task_id = 'task-12345'
       ORDER BY a.started_at DESC
       LIMIT 1;

# 4. 检查是否有性能瓶颈
# - 数据库慢查询
mysql> SHOW PROCESSLIST;

# - 外部 API 调用延迟
tail -f /var/log/app.log | grep "API调用"

# - CPU/内存使用率
top -p <worker_pid>

常见原因与解决方案:

原因解决方案
deadline_at 设置过短提交任务时设置更长的截止时间
max-execution-time 过短调大配置值(如 1h)
数据库慢查询优化 SQL,增加索引
外部 API 超时增加超时配置,增加重试
数据量过大优化批次大小,减少单批处理量

5. 死信队列堆积

现象: 死信队列任务数持续增长。

排查步骤:

# 1. 查询死信队列大小
curl http://localhost:8080/actuator/metrics/async.task.dead.letter.queue.size

# 2. 查询死信任务列表
mysql> SELECT * FROM system_async_task_dead_letter 
       ORDER BY moved_at DESC 
       LIMIT 20;

# 3. 统计死信任务的错误码分布
mysql> SELECT error_code, COUNT(*) AS count
       FROM system_async_task_dead_letter 
       GROUP BY error_code 
       ORDER BY count DESC;

# 4. 查看失败任务的详细信息
mysql> SELECT * FROM system_async_task_failure 
       WHERE task_id IN (
         SELECT original_task_id FROM system_async_task_dead_letter
       )
       LIMIT 10;

处理策略:

// 1. 分析死信任务,修复根因后批量重放
List<String> deadLetterIds = deadLetterQueueService.listDeadLetters(
    "IMPORT.MEDICAL",
    Instant.now().minus(Duration.ofDays(1)),
    Instant.now(),
    100
);

// 2. 批量重放
deadLetterQueueService.replayBatch(deadLetterIds);

// 3. 无法修复的任务,归档并清理
deadLetterQueueService.archiveAndClean(deadLetterIds);

运维工具

1. 管理端 API

/**
 * 任务管理 API(仅管理员访问)
 */
@RestController
@RequestMapping("/api/admin/async-task")
@PreAuthorize("hasRole('ADMIN')")
public class AsyncTaskAdminController {
    
    /**
     * 暂停任务
     */
    @PostMapping("/{taskId}/pause")
    public BaseResult<Void> pauseTask(@PathVariable String taskId) {
        asyncTaskService.pause(taskId, AsyncTaskQueryScope.global());
        return BaseResult.success();
    }
    
    /**
     * 恢复任务
     */
    @PostMapping("/{taskId}/resume")
    public BaseResult<Void> resumeTask(@PathVariable String taskId) {
        asyncTaskService.resume(taskId, AsyncTaskQueryScope.global());
        return BaseResult.success();
    }
    
    /**
     * 调整全局并发
     */
    @PostMapping("/policy/global/concurrency")
    public BaseResult<Void> adjustGlobalConcurrency(
        @RequestParam int maxConcurrency
    ) {
        policyService.save(new AsyncTaskPolicy(
            AsyncTaskPolicy.GLOBAL_POLICY_KEY,
            null,
            maxConcurrency,
            null,  // 保持原重试配置
            null,
            true
        ));
        return BaseResult.success();
    }
    
    /**
     * 批量取消任务
     */
    @PostMapping("/batch-cancel")
    public BaseResult<Integer> batchCancel(
        @RequestBody List<String> taskIds
    ) {
        int count = 0;
        for (String taskId : taskIds) {
            try {
                asyncTaskService.requestCancellation(
                    taskId, 
                    AsyncTaskQueryScope.global()
                );
                count++;
            } catch (Exception ex) {
                log.error("取消任务失败, taskId={}", taskId, ex);
            }
        }
        return BaseResult.success(count);
    }
    
    /**
     * 清理历史任务
     */
    @PostMapping("/cleanup")
    public BaseResult<AsyncTaskCleanupResult> cleanup(
        @RequestParam int retentionDays,
        @RequestParam(defaultValue = "500") int batchSize
    ) {
        Instant finishedBefore = Instant.now()
            .minus(Duration.ofDays(retentionDays));
        
        AsyncTaskCleanupResult result = store.purgeTerminalTasks(
            finishedBefore, 
            batchSize
        );
        
        return BaseResult.success(result);
    }
}

2. 诊断脚本

2.1 任务状态诊断

#!/bin/bash
# diagnose-task.sh

TASK_ID=$1

if [ -z "$TASK_ID" ]; then
  echo "用法: $0 <task_id>"
  exit 1
fi

echo "========== 任务状态诊断 =========="
echo "任务ID: $TASK_ID"
echo ""

# 1. 查询任务基本信息
echo "=== 任务基本信息 ==="
mysql -e "SELECT * FROM system_async_task WHERE id = '$TASK_ID' \G"

# 2. 查询执行记录
echo "=== 执行记录 ==="
mysql -e "SELECT * FROM system_async_task_attempt WHERE task_id = '$TASK_ID' ORDER BY started_at DESC \G"

# 3. 查询失败明细
echo "=== 失败明细 ==="
mysql -e "SELECT * FROM system_async_task_failure WHERE task_id = '$TASK_ID' ORDER BY failed_at DESC LIMIT 5 \G"

# 4. 检查租约状态
echo "=== 租约状态 ==="
mysql -e "
  SELECT 
    id, 
    status, 
    lease_owner, 
    lease_until,
    CASE 
      WHEN lease_until IS NULL THEN 'NO_LEASE'
      WHEN lease_until < NOW() THEN 'EXPIRED'
      ELSE 'ACTIVE'
    END AS lease_status,
    TIMESTAMPDIFF(SECOND, NOW(), lease_until) AS remaining_sec
  FROM system_async_task 
  WHERE id = '$TASK_ID' \G
"

# 5. 检查 Worker 健康度
echo "=== Worker 健康度 ==="
curl -s http://localhost:8080/actuator/health/asyncTaskWorker | jq .

2.2 性能分析脚本

#!/bin/bash
# analyze-performance.sh

TASK_TYPE=$1
HOURS=${2:-24}

echo "========== 任务性能分析 =========="
echo "任务类型: $TASK_TYPE"
echo "分析时长: 最近 $HOURS 小时"
echo ""

# 1. 任务数量统计
echo "=== 任务数量统计 ==="
mysql -e "
  SELECT 
    status,
    COUNT(*) AS count
  FROM system_async_task
  WHERE task_type = '$TASK_TYPE'
    AND create_time > NOW() - INTERVAL $HOURS HOUR
  GROUP BY status
  ORDER BY count DESC
"

# 2. 执行时长分析
echo "=== 执行时长分析 ==="
mysql -e "
  SELECT 
    AVG(TIMESTAMPDIFF(SECOND, started_at, finished_at)) AS avg_duration_sec,
    MIN(TIMESTAMPDIFF(SECOND, started_at, finished_at)) AS min_duration_sec,
    MAX(TIMESTAMPDIFF(SECOND, started_at, finished_at)) AS max_duration_sec,
    COUNT(*) AS sample_count
  FROM system_async_task_attempt a
  JOIN system_async_task t ON a.task_id = t.id
  WHERE t.task_type = '$TASK_TYPE'
    AND a.started_at > NOW() - INTERVAL $HOURS HOUR
    AND a.finished_at IS NOT NULL
"

# 3. 队列等待时长分析
echo "=== 队列等待时长分析 ==="
mysql -e "
  SELECT 
    AVG(TIMESTAMPDIFF(SECOND, t.create_time, a.started_at)) AS avg_wait_sec,
    MIN(TIMESTAMPDIFF(SECOND, t.create_time, a.started_at)) AS min_wait_sec,
    MAX(TIMESTAMPDIFF(SECOND, t.create_time, a.started_at)) AS max_wait_sec
  FROM system_async_task t
  JOIN system_async_task_attempt a ON t.id = a.task_id
  WHERE t.task_type = '$TASK_TYPE'
    AND t.create_time > NOW() - INTERVAL $HOURS HOUR
"

# 4. 失败原因分析
echo "=== 失败原因分析 ==="
mysql -e "
  SELECT 
    error_code,
    COUNT(*) AS count,
    COUNT(*) * 100.0 / SUM(COUNT(*)) OVER () AS percentage
  FROM system_async_task_failure f
  JOIN system_async_task t ON f.task_id = t.id
  WHERE t.task_type = '$TASK_TYPE'
    AND f.failed_at > NOW() - INTERVAL $HOURS HOUR
  GROUP BY error_code
  ORDER BY count DESC
  LIMIT 10
"

3. Grafana 看板模板

完整的 Grafana Dashboard JSON 配置:

{
  "dashboard": {
    "title": "异步任务监控",
    "panels": [
      {
        "title": "任务提交速率",
        "targets": [{
          "expr": "sum(rate(async_task_submitted_total[5m])) by (task_type)"
        }]
      },
      {
        "title": "任务完成速率",
        "targets": [{
          "expr": "sum(rate(async_task_completed_total[5m])) by (status)"
        }]
      },
      {
        "title": "活跃任务数",
        "targets": [{
          "expr": "async_task_active_count"
        }]
      },
      {
        "title": "队列等待时长 P99",
        "targets": [{
          "expr": "histogram_quantile(0.99, rate(async_task_queue_wait_seconds_bucket[5m]))"
        }]
      },
      {
        "title": "执行时长 P99(按类型)",
        "targets": [{
          "expr": "histogram_quantile(0.99, rate(async_task_execution_seconds_bucket[5m])) by (task_type)"
        }]
      },
      {
        "title": "任务失败率",
        "targets": [{
          "expr": "sum(rate(async_task_completed_total{status='FAILED'}[5m])) / sum(rate(async_task_completed_total[5m]))"
        }]
      },
      {
        "title": "租约续期成功率",
        "targets": [{
          "expr": "sum(rate(async_task_lease_renewal_success_total[5m])) / sum(rate(async_task_lease_renewal_total[5m]))"
        }]
      },
      {
        "title": "死信队列大小",
        "targets": [{
          "expr": "async_task_dead_letter_queue_size"
        }]
      }
    ]
  }
}

总结

一个成熟的异步任务框架不仅要功能完善,更要具备强大的可观测性和可维护性。本文介绍的监控、运维和故障排查体系包括:

监控体系

  • 四个黄金信号:延迟、流量、错误、饱和度
  • 租约健康度:续期成功率、过期恢复、平均持有时长
  • 业务指标:成功率、处理速率、用户分布
  • 健康检查:控制面、Worker、Kubernetes 探针

运维工具

  • 管理端 API:暂停/恢复、批量操作、策略调整、清理历史
  • 诊断脚本:任务状态诊断、性能分析
  • 可视化看板:Grafana Dashboard

故障排查

  • 常见问题:排队不执行、反复恢复、重复执行、超时、死信堆积
  • 排查流程:现象 → 排查步骤 → 常见原因 → 解决方案
  • 日志规范:级别定义、格式规范、敏感信息脱敏

通过这套完整的监控和运维体系,你可以构建一个"看得见、管得住、修得快"的异步任务系统,为业务提供稳定可靠的服务。


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

完整系列总结:

通过这五篇文章,我们完整地介绍了 tudicloud-async-task 异步任务框架:

  1. 整体架构:设计理念、核心组件、关键特性
  2. 租约机制:分布式协调、边界问题、性能优化
  3. Worker 实现:调度策略、执行流程、容错设计
  4. 业务接入:需求分析、代码实现、上线部署
  5. 监控运维:指标体系、故障排查、运维工具

希望这个系列能帮助你深入理解企业级异步任务框架的设计与实现!

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