【问题标题】:Celery time statistics per-task-name每个任务名称的芹菜时间统计
【发布时间】:2016-05-30 22:22:35
【问题描述】:

我有一些相当繁忙的芹菜队列,但不确定哪些任务是有问题的。有没有办法汇总结果以找出哪些任务需要很长时间?我在 2-4 台服务器上有 10-20 名工作人员。

使用 redis 作为代理和结果后端。我注意到 Flower 上的队列很忙,但不知道如何汇总每个任务的时间统计信息。

【问题讨论】:

    标签: python-2.7 celery django-celery flower


    【解决方案1】:

    方法一:

    如果您在启动 celery worker 时启用了日志记录,它们会记录每个任务所花费的时间。

    $ celery worker -l info -A your_app --logfile celery.log
    

    这将生成这样的日志

    [2016-06-04 13:21:30,749: INFO/MainProcess] Task sig.add[a8b648eb-9674-44f0-90bd-71cfebe22f2f] succeeded in 0.00979363399983s: 3
    [2016-06-04 13:21:30,973: INFO/MainProcess] Received task: sig.add[7fd422e6-8f48-4dd2-90de-e213afbedc38]
    [2016-06-04 13:21:30,982: WARNING/Worker-2] called by small_task. LOL {'signal': <Signal: Signal>, 'result': 3, 'sender': <@task: sig.add of tasks:0x7fdf33146c50>}
    

    您可以过滤具有succeeded in 的行。使用[: 作为分隔符拆分这些行,打印任务名称和每行花费的时间,然后对所有行进行排序。

    $ grep ' succeeded in ' celery.log  | awk -F'[ :\[]' '{print $9, $13}' | sort 
    awk: warning: escape sequence `\[' treated as plain `['
    sig.add 0.00775764500031s
    sig.add 0.00802627899975s
    sig.foo 12.00813863099938s
    sig.foo 15.00871706100043s
    sig.foo 12.00979363399983s
    

    如您所见,add 非常快,foo 很慢。

    方法二:

    Celery 有task_prerun_handler,task_postrun_handler 信号,它们在任务之前/之后运行。您可以连接跟踪时间的函数,然后在某处记录时间。

    from time import time
    from celery.signals import task_prerun, task_postrun
    
    
    tasks = {}
    task_avg_time = {}
    Average = namedtuple('Average', 'cum_avg count')
    
    
    @task_prerun.connect
    def task_prerun_handler(signal, sender, task_id, task, args, kwargs):
        tasks[task_id] = time()
    
    
    @task_postrun.connect
    def task_postrun_handler(signal, sender, task_id, task, args, kwargs, retval, state):
        try:
            cost = time() - tasks.pop(task_id)
        except KeyError:
            cost = None
    
        if not cost:
            return
    
        try:
            cum_avg, count = task_avg_time[task.name]
            new_count = count + 1
            new_avg = ((cum_avg * count) + cost) / new_count
            task_avg_time[task.name] = Average(new_avg, new_count)
        except KeyError:
            task_avg_time[task.name] = Average(cost, 1)
    
        # write to redis: task_avg_time
    

    参考:https://stackoverflow.com/a/31731622/2698552

    【讨论】:

      猜你喜欢
      • 2018-04-28
      • 2018-01-28
      • 2020-05-02
      • 1970-01-01
      • 2022-06-10
      • 2012-03-24
      • 2021-01-02
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多