You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何编写Kibana Timelion表达式获取Celery任务平均执行时间

实现Celery任务平均执行时间统计及周期变化(ELK 7.17)

完全可行,但不推荐用Timelion——它在7.17版本属于遗留功能,且本身不擅长处理需要关联"启动/结束"两类事件的计算。更高效的方案是通过Logstash预处理或Elasticsearch聚合实现,再用Kibana标准可视化展示周期变化。

核心思路

通过trace_id关联同一任务的「Task started」和「Task finished」日志,计算单任务执行时长,再按task_type聚合求平均值,最后按时间周期拆分展示变化。

方案一:Logstash预处理(推荐)

在日志写入Elasticsearch前,用Logstash的aggregate插件提前计算任务时长,后续统计会更高效。

示例Logstash配置片段:

filter {
  # 捕获任务启动事件,存储启动时间和任务类型
  if [message] == "Task started" {
    aggregate {
      task_id => "%{trace_id}"
      code => "map['start_time'] = event.get('@timestamp'); map['task_type'] = event.get('task_type')"
      map_action => "create"
    }
  }

  # 捕获任务结束事件,计算时长并补充到事件中
  if [message] == "Task finished" {
    aggregate {
      task_id => "%{trace_id}"
      code => "
        start_time = map['start_time'];
        finish_time = event.get('@timestamp');
        duration_in_seconds = finish_time.to_f - start_time.to_f;
        event.set('duration', duration_in_seconds);
        event.set('task_type', map['task_type']);
      "
      map_action => "update"
      end_of_task => true
      timeout => 3600  # 按任务最长执行时间调整超时阈值
    }
  }
}

处理后,每条「Task finished」日志会新增duration字段(单位:秒),直接用于后续统计。

方案二:Elasticsearch聚合计算(无需修改Logstash)

如果无法调整Logstash配置,可通过Elasticsearch的管道聚合直接计算。示例DSL查询:

GET /your-celery-logs/_search
{
  "size": 0,
  "aggs": {
    # 按trace_id分组,提取每个任务的启动/结束时间
    "group_trace": {
      "terms": {
        "field": "trace_id",
        "size": 10000  # 根据任务量调整,避免遗漏
      },
      "aggs": {
        "start_time": {"min": {"field": "@timestamp"}},
        "finish_time": {"max": {"field": "@timestamp"}},
        "task_type": {"terms": {"field": "task_type", "size": 10}},
        # 计算单任务时长(毫秒转秒)
        "task_duration": {
          "scripted_metric": {
            "init_script": "state.times = []",
            "map_script": "state.times.add(params._source['@timestamp'].toInstant().toEpochMilli())",
            "combine_script": "return state.times",
            "reduce_script": "
              def start = states.flatten().min();
              def finish = states.flatten().max();
              return (finish - start)/1000;
            "
          }
        }
      }
    },
    # 按task_type分组,求平均时长
    "group_task_type": {
      "terms": {"field": "task_type"},
      "aggs": {
        "avg_duration": {
          "avg_bucket": {"buckets_path": "group_trace>task_duration.value"}
        }
      }
    }
  }
}

可视化周期变化

当有了duration字段后,在Kibana中创建对应数据视图,然后:

  1. 选择「折线图」或「柱状图」可视化类型
  2. X轴选择@timestamp,设置时间间隔(如小时、天)
  3. Y轴选择「平均」聚合,字段为duration
  4. 添加「拆分系列」,按task_type分组
  5. 调整时间范围,即可查看指定周期内各任务类型的平均执行时间变化

注意事项

  • 确保所有任务都有对应的「启动/结束」日志,避免trace_id缺失导致统计偏差
  • Logstash的aggregate插件要设置合理的timeout,防止未完成任务占用内存
  • Elasticsearch聚合的size参数需覆盖统计周期内的所有trace_id,否则会遗漏数据

内容的提问来源于stack exchange,提问作者Dogrtt

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.21 06:37:18