【问题标题】:Why does celery retry, but my job did not fail?为什么 celery 重试,但我的工作没有失败?
【发布时间】:2021-10-29 21:07:57
【问题描述】:

我有一个 celery 工作来运行 MySQL 数据库,但是,它总是得到Lock Wait Timeout。在深入研究数据库查询后,我意识到 celery 在 1800 秒后触发了另一项工作,并遇到了我的数据库问题。我不知道为什么——我的工作还没有失败!

@celery.task(bind=True, acks_late=True)
def etl_pipeline(dev=dev, test=test):

我可以告诉 MySQL 再次得到相同的查询,可能是 Celery 触发了相同的工作。为什么我在这里重试,默认重试是 180 秒(3 分钟)。

这是官方文档:

default_retry_delay = 180

重试任务之前的默认时间(以秒为单位)。默认 3 分钟。

但我的情况是 1800 秒。

另外,我的经纪人收到了一些其他警告,我不确定这是否相关:

AMQP 结果后端计划在 4.0 版中弃用并在 v5.0 版中移除。请使用 RPC 后端或持久化后端。

配置 RabbitMq

RABBITMQ_SERVER = 'amqp://{}:{}@{}'.format(
    os.getenv('RABBITMQ_USER'),
    os.getenv('RABBITMQ_PASS'),
    os.getenv('RABBITMQ_HOST')
)
broker_url = '{}/{}'.format(
    RABBITMQ_SERVER,
    os.getenv('RABBITMQ_VHOST'),
)
backend = 'amqp'

我该如何解决这个问题?谢谢!

芹菜:4.2.0

I am using job = chain(single_job), but i only have one single_job job() starting the job.

mysql> show processlist;
+-------+------+---------------+------------------+---------+------+-----------+
| Id    | User | Host          | db               | Command | Time | State     |
+-------+------+---------------+------------------+---------+------+-----------+
| 97189 | clp  | 172.11.17.202 | bain_ai_database | Query   |    0 | init      |
| 97488 | clp  | 172.11.11.252 | bain_ai_database | Query   | 1505 | executing |
| 97489 | clp  | 172.11.11.252 | bain_ai_database | Sleep   | 1851 |           |
| 97543 | clp  | 172.21.6.242  | bain_ai_database | Query   |   51 | updating  |
| 97544 | clp  | 172.21.6.242  | bain_ai_database | Sleep   |   51 |           |
+-------+------+---------------+------------------+---------+------+-----------+

【问题讨论】:

    标签: python socket.io rabbitmq celery session-timeout


    【解决方案1】:

    根据您执行 sql 查询的方式,我会尝试以下方法。 (1) 既然你有bind=True,那么任务应该是你函数的第一个参数。 celery 中的约定是调用第一个参数self。 (2) 您想尝试捕获正在发生的数据库级异常并忽略它。

    from celery.utils.log import get_task_logger
    
    log = get_task_logger(__name__)
    
    
    @celery.task(bind=True, acks_late=True)
    def etl_pipeline(self, dev=dev, test=test):
        try:
            # try querying the database here using sqlalchemy or mysqlconnect??
        except Exception as ex:
            # for now, log the exception and type so that you can drill down into what is happening
            log.info('[etl_pipeline]  exception of type %s.%s: %s', ex.__class__.__module__, ex.__class__.__name__, ex)
            raise       
    

    您将从日志记录中获得的调试应该帮助您确定您在客户端遇到的错误。

    【讨论】:

    • 如果我只使用@celery.task(bind=True),这是否会禁用retry?我只想禁用retry,因为我的sql查询不需要message queue重试。
    • 这就是我所做的,我升级celery==4.4.2,添加缓存global duplicated = False 并删除acks_late=True。我启动celery 并休眠1800 second 的sql 查询。看起来工人是一样的,没有触发新的工人。稍后我将在生产环境中对其进行测试。
    • 你说得对,我还需要添加数据库级别的异常处理来捕获数据库错误,它也可能是来自数据库的错误。非常感谢您的帮助。
    猜你喜欢
    • 1970-01-01
    • 2022-12-03
    • 2012-09-16
    • 1970-01-01
    • 1970-01-01
    • 2012-07-02
    • 2021-02-27
    • 2012-07-25
    • 1970-01-01
    相关资源
    最近更新 更多