【问题标题】:How to capture the celery warning(s) from logs?如何从日志中捕获芹菜警告?
【发布时间】:2022-06-12 23:47:43
【问题描述】:

我将celeryfastapi 结合使用,结果通过命令记录在一个名为celery.log 的文件中

celery worker --app=app.celery_worker.celery --loglevel=info --logfile=app/logs/celery.log

代码触发时,日志写入celery.log文件如下:

[2022-05-20 11:38:35,148: INFO/MainProcess] Received task: kwept_calculation[bd80737a-92cd-4fea-8a68-c010d5ab3ed3]  
[2022-05-20 11:38:46,249: WARNING/ForkPoolWorker-7] System mit Neuinstallation(en) nicht betreibbar (z.B. wegen nicht deckbarer Wärmenachfrage)
[2022-05-20 11:38:53,401: WARNING/ForkPoolWorker-7] System ohne Neuinstallation(en) nicht betreibbar (z.B. wegen nicht deckbarer Wärmenachfrage)
[2022-05-20 11:38:53,402: ERROR/ForkPoolWorker-7] Task kwept_calculation[bd80737a-92cd-4fea-8a68-c010d5ab3ed3] raised unexpected: UnboundLocalError("local variable 'kpis_df' referenced before assignment")
Traceback (most recent call last):
  File "/usr/local/lib/python3.7/site-packages/celery/app/trace.py", line 412, in trace_task
    R = retval = fun(*args, **kwargs)
  File "/usr/local/lib/python3.7/site-packages/celery/app/trace.py", line 704, in __protected_call__
    return self.run(*args, **kwargs)
  File "/kwept/app/celery_worker.py", line 77, in calculation
    created_at_ts
  File "/kwept/app/services/planning_calc.py", line 571, in plan_calc
    kpis_df.reset_index(drop=False, inplace=True)
UnboundLocalError: local variable 'kpis_df' referenced before assignment

要获取有关任务的信息,我会这样做

from celery.result import AsyncResult
task_result = AsyncResult(task_id) # in the above case bd80737a-92cd-4fea-8a68-c010d5ab3ed3
task_info = task_result.info

当我这样做时,会捕获来自上述task_info 的错误

赋值前引用的局部变量“kpis_df”

有没有办法捕获警告消息?在上面的例子中,警告是:

System mit Neuinstallation(en) nicht betreibbar (z.B. wegen nicht Deckbarer Wärmenachfrage)

System ohne Neuinstallation(en) nicht betreibbar (z.B. wegen nicht Deckbarer Wärmenachfrage)

【问题讨论】:

  • 警告正在写入日志文件,您还想如何捕获它们?

标签: python redis celery celery-task


【解决方案1】:

如果您希望将警告消息作为 Celery 任务本身的自定义字段的一部分,您可以尝试这样的操作并通过 update_state() method 将其添加到元字段中的自定义条目中:

@app.task(bind=True)
def task(self):
try:
    raise ValueError('Some error')
except Exception as ex:
    self.update_state(
        state=states.FAILURE,
        meta={
            'exc_type': type(ex).__name__,
            'exc_message': traceback.format_exc().split('\n')
            'custom': '...'
        })
    raise Ignore()

在您的任务定义中,您必须更新代码以不仅记录警告消息,而且将其放入元字段中。

当您需要检查任务时,您可以在此处找到您的自定义字段:

>>> print(task.info)
{'custom': '...'}

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-04-09
    • 2015-03-29
    • 1970-01-01
    • 1970-01-01
    • 2013-02-25
    • 2022-01-21
    相关资源
    最近更新 更多