【问题标题】:How to detect Celery task which doing similar job before run another task?如何在运行另一个任务之前检测执行类似工作的 Celery 任务?
【发布时间】:2018-12-24 14:14:39
【问题描述】:

我的 celery 任务是对一些数据库存储的实体进行耗时的计算。工作流程是这样的:从数据库中获取信息,将其编译为一些可序列化的对象,保存对象。其他任务正在加载的对象上进行其他计算(如渲染图像)。

但是序列化是耗时的,所以我希望每个实体有一个任务运行一段时间,它将序列化对象保存在内存中并处理客户端请求,通过消息队列(redis pubsub)传递。如果一段时间没有请求,则任务退出。之后,如果客户端需要完成某项工作,它会运行另一个任务,该任务加载对象、处理它并为其他工作等待一段时间。此任务应在启动时检查是否只有该特定实体上的一个工作人员以避免冲突。 那么检查是否有另一个任务正在为此实体运行的最佳策略是什么?

1) 第一个想法是将消息发送到与实体关联的某个通道,并等待响应。坏主意,目标任务可能忙于计算,等待超时响应只是在浪费时间。

2) 将 celery task-id 存储在 db 中更糟糕 - 任务可以被杀死,但记录会保留,因此我们需要确保目标任务是活动的。

3)第三个想法是检查工作人员是否正在运行任务,检查它的状态以获取实体 ID(哪个任务将在启动时提供)。似乎也可能发生一些冲突,即如果计划了多个任务但尚未运行。

目前我认为想法 1 是最好的修改如下:任务将在启动时将消息发送到实体通道及其启动时间,但随后立即开始工作,而不是等待响应。然后它检查消息队列,如果有人响应,他们会比较时间戳和任务与更大的时间戳退出。看起来够复杂,有没有更好的解决方案?

【问题讨论】:

    标签: python celery ipc


    【解决方案1】:

    最终的解决方案是在任务中启动主管线程,它会回复来自竞争任务的“发现”消息。

    所以工作流程就是这样。

    1. 任务启动,然后使用实体 ID 订阅 Redis PubSub 频道
    2. 任务向通道发送“发现”消息
    3. 任务稍等
    4. 如果找到退出,则在频道中的传入消息中搜索“回复”。
    5. 任务启动主管线程,该线程通过“回复”回复所有传入的“发现”消息

    这工作正常,除了几个任务同时启动,即在工作人员重新启动之后。为了避免这种需要使订阅过程原子化,使用 Redis 锁:

    class RedisChannel:
        def __init__(self, channel_id):
            self.channel_id = channel_id
            self.redis = StrictRedis()
            self.channel = self.redis.pubsub()
            with self.redis.lock(channel_id):
                self.channel.subscribe(channel_id)
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-05-31
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-04-09
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多