【问题标题】:Celery: access all previous results in a chain芹菜:访问链中所有以前的结果
【发布时间】:2015-04-09 19:52:09
【问题描述】:

所以基本上我有一个相当复杂的工作流程,看起来类似于:

>>> res = (add.si(2, 2) | add.s(4) | add.s(8))()
>>> res.get()
16

之后,我沿着结果链向上走并收集所有单独的结果是相当微不足道的:

>>> res.parent.get()
8

>>> res.parent.parent.get()
4

我的问题是,如果我的第三个任务依赖于知道第一个的结果,但在这个例子中只收到第二个的结果怎么办?

而且链很长,结果也不是那么小,所以仅仅通过输入作为结果会不必要地污染结果存储。哪个是 Redis,所以使用 RabbitMQ、ZeroMQ、...时的限制不适用。

【问题讨论】:

    标签: python redis celery


    【解决方案1】:

    也许你的设置太复杂了,但我喜欢使用group 结合noop 任务来完成类似的事情。我这样做是因为我想突出显示在我的管道中仍然同步的区域(通常这样它们可以被删除)。

    使用与您的示例类似的内容,我从一组如下所示的任务开始:

    tasks.py:

    from celery import Celery
    
    app = Celery('tasks', backend="redis", broker='redis://localhost')
    
    
    @app.task
    def add(x, y):
            return x + y
    
    
    @app.task
    def xsum(elements):
        return sum(elements)
    
    
    @app.task
    def noop(ignored):
        return ignored
    

    通过这些任务,我使用组创建了一个链来控制依赖于同步结果的结果:

    In [1]: from tasks import add,xsum,noop
    In [2]: from celery import group
    
    # First I run the task which I need the value of later, then I send that result to a group where the first task does nothing and the other tasks are my pipeline.
    In [3]: ~(add.si(2, 2) | group(noop.s(),  add.s(4) | add.s(8)))
    Out[3]: [4, 16]
    
    # At this point I have a list where the first element is the result of my original task and the second element has the result of my workflow.
    In [4]: ~(add.si(2, 2) | group(noop.s(),  add.s(4) | add.s(8)) | xsum.s())
    Out[4]: 20
    
    # From here, things can go back to a normal chain
    In [5]: ~(add.si(2, 2) | group(noop.s(),  add.s(4) | add.s(8)) | xsum.s() | add.s(1) | add.s(1))
    Out[5]: 22
    

    希望对你有用!

    【讨论】:

    • 这太棒了!开箱即用的想法和所有答案中最小的足迹。不幸的是,这是一个开放的 celery 错误,它可以防止嵌套组,所以我现在不能使用它,但我期待有一天能切换到这个!
    • 哦,有趣。您是否碰巧有该错误的链接?我有时会使用嵌套组,想了解更多信息。
    • @erik-e 有没有办法获得最后一个任务输出并在链上创建一个组?比如(return_range_task.s() | group(add.s(I, 4) for for I in range_task_output) | filter_things.s(),其中return_range_task 将返回数组,add 任务将该数组元素作为第一个参数并处理它们?
    【解决方案2】:

    我为每个链分配一个作业 ID,并通过将数据保存在数据库中来跟踪该作业。

    启动队列

    if __name__ == "__main__":
      # Generate unique id for the job
      job_id = uuid.uuid4().hex
      # This is the root parent
      parent_level = 1
      # Pack the data. The last value is your value to add
      parameters = job_id, parent_level, 2
      # Build the chain. I added an clean task that removes the data
      # created during the process (if you want it)
      add_chain = add.s(parameters, 2) | add.s(4) | add.s(8)| clean.s()
      add_chain.apply_async()
    

    现在的任务

    #Function for store the result. I used sqlalchemy (mysql) but you can
    # change it for whatever you want (distributed file system for example)
    @inject.params(entity_manager=EntityManager)
    def save_result(job_id, level, result, entity_manager):
      r = Result()
      r.job_id = job_id
      r.level = level
      r.result = result
      entity_manager.add(r)
      entity_manager.commit()
    
    #Restore a result from one parent
    @inject.params(entity_manager=EntityManager)
    def get_result(job_id, level, entity_manager):
      result = entity_manager.query(Result).filter_by(job_id=job_id, level=level).one()
      return result.result
    
    #Clear the data or do something with the final result
    @inject.params(entity_manager=EntityManager)
      def clear(job_id, entity_manager):
      entity_manager.query(Result).filter_by(job_id=job_id).delete()
    
    @app.task()
    def add(parameters, number):
      # Extract data from parameters list
      job_id, level, other_number = parameters
    
      #Load result from your second parent (level - 2)
      #For level 3 parent level - 3 and so on
      #second_parent_result = get_result(job_id, level - 2)
    
      # do your stuff, I guess you want to add numbers
      result = number + other_number
      save_result(job_id, level, result)
    
      #Return the result of the sum or anything you want, but you have to send something because the "add" function expects 3 values
      #Of course your should return the actual job and increment the parent level
      return job_id, level + 1, result
    
    @app.task()
    def clean(parameters):
      job_id, level, result = parameters
      #Do something with final result or not
      #Clear the data
      clear(job_id)
    

    我使用 entity_manager 来管理数据库操作。我的实体管理器使用 sql alchemy 和 mysql。我还使用了一个表格“结果”来存储部分结果。这部分应该针对您最好的存储系统进行更改(或者如果 mysql 适合您,请使用此部分)

    from sqlalchemy.orm import sessionmaker
    from sqlalchemy import create_engine
    import inject
    
    class EntityManager():
    
      session = None
    
      @inject.params(config=Configuration)
      def __init__(self, config):
        conf = config['persistence']
        uri = conf['driver'] + "://" + conf['username'] + ":@" + conf['host'] + "/" + conf['database']
    
        engine = create_engine(uri, echo=conf['debug'])
    
        Session = sessionmaker(bind=engine)
        self.session = Session()
    
      def query(self, entity_type):
        return self.session.query(entity_type)
    
      def add(self, entity):
        return self.session.add(entity)
    
      def flush(self):
        return self.session.flush()
    
      def commit(self):
        return self.session.commit()
    
    class Configuration:
      def __init__(self, params):
        f = open(os.environ.get('PYTHONPATH') + '/conf/config.yml')
        self.configMap = yaml.safe_load(f)
        f.close()
    
      def __getitem__(self, key: str):
        return self.configMap[key]
    
    class Result(Base):
      __tablename__ = 'result'
    
      id = Column(Integer, primary_key=True)
      job_id = Column(String(255))
      level = Column(Integer)
      result = Column(Integer)
    
      def __repr__(self):
        return "<Result (job='%s', level='%s', result='%s')>" % (self.job_id, str(self.level), str(self.result))
    

    我使用包注入来获取依赖注入器。注入包将重用该对象,因此您可以在每次需要时注入对数据库的访问,而无需担心连接。

    类配置是将数据库访问数据加载到配置文件中。您可以替换它并使用静态数据(硬编码的地图)进行测试。

    为任何其他适合你的东西更改依赖注入。这只是我的解决方案。我只是为了快速测试添加它。

    这里的关键是将部分结果保存在队列系统中的某个位置,并在任务中返回数据以访问这些结果(job_id 和父级)。您将发送这个额外(但很小)的数据,它是一个地址(job_id + parent 级别),指向真实数据(一些大的东西)。

    我在我的软件中使用的这个解决方案

    【讨论】:

    • 谢谢!老实说,我认为所有三个答案都很棒,值得赏金。我选择了你的,因为它通过将以前的结果排除在我的实际结果存储之外,从而产生最少的头痛。
    • 是否可以读取链中的前一个任务输出并创建一个类似(return_range_task.s() | group(add.s(I, 4) for for I in range_task_output) | filter_things.s())的组?
    【解决方案3】:

    一个简单的解决方法是将任务的结果存储在一个列表中并在您的任务中使用它们。

    from celery import Celery, chain
    from celery.signals import task_success
    
    results = []
    
    app = Celery('tasks', backend='amqp', broker='amqp://')
    
    
    @task_success.connect()
    def store_result(**kwargs):
        sender = kwargs.pop('sender')
        result = kwargs.pop('result')
        results.append((sender.name, result))
    
    
    @app.task
    def add(x, y):
        print("previous results", results)
        return x + y
    

    现在,在您的链中,可以从任何任务以任何顺序访问所有以前的结果。

    【讨论】:

    • (return_range_task.s() | group(add.s(I, 4) for for I in range_task_output) | filter_things.s())有什么办法吗?
    • 可以使用 celery 信号,然后在 task_success 信号中触发组。 @Nilesh
    • 它和stackoverflow.com/a/14995090/243031 一样,我可以在另一个任务中这样做,我正在寻找可以在同一任务链中访问的东西。
    猜你喜欢
    • 2019-07-22
    • 1970-01-01
    • 2014-06-22
    • 2016-06-03
    • 1970-01-01
    • 1970-01-01
    • 2016-07-14
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多