【问题标题】:Implementing Twisted style local multiple deferred callbacks in Celery在 Celery 中实现 Twisted 样式的本地多个延迟回调
【发布时间】:2013-12-09 16:26:49
【问题描述】:

我对使用 Celery 很陌生,想知道如何在 Celery 中实现 TWSITED 类型的多个延迟回调

MY TWISTE CODE 使用透视代理,如下所示。我有一个处理程序(服务器),它处理一些事件并返回结果。 Dispatcher(客户端)使用延迟回调打印返回的结果。

Handler.py(服务器)

from twisted.application import service, internet
from twisted.internet import reactor, task
from twisted.spread import pb
from Dispatcher import Event
from Dispatcher import CopyEvent

class ReceiverEvent(pb.RemoteCopy, Event):
    pass
pb.setUnjellyableForClass(CopyEvent, ReceiverEvent)


class Handler(pb.Root):

def remote_eventEnqueue(self, pond):
    d = task.deferLater(reactor,5,handle_event,sender=self)
    return d

def handle_event(sender):
    print "Do Something"
    return "did something"

if __name__ == '__main__':
    h=Handler()
    reactor.listenTCP(8739, pb.PBServerFactory(h))
    reactor.run()

现在是 Dispatcher.py(客户端)

from twisted.spread import pb, jelly
from twisted.python import log
from twisted.internet import reactor
from Event import Event

class CopyEvent(Event, pb.Copyable):
    pass

class Dispatcher:
    def __init__(self, event):
        self.event = event

    def dispatch_event(self, remote):
        d = remote.callRemote("eventEnqueue", self.event)   
        d.addCallback(self.printMessage)

    def printMessage(self, text):
        print text

def main():
    from Handler import CopyEvent
    event = CopyEvent()
    d = Dispatcher(event)
    factory = pb.PBClientFactory()
    reactor.connectTCP("localhost", 8739, factory)
    deferred = factory.getRootObject()
    deferred.addCallback(d.dispatch_event)
    reactor.run()

if __name__ == '__main__':
    main()

我尝试在 Celery 中实现这一点。

Handler.py(服务器)

from celery import Celery

app=Celery('tasks',backend='amqp',broker='amqp://guest@localhost//')

@app.task

def handle_event():
     print "Do Something"
     return "did something"

Dispatcher.py(客户端)

from Handler import handle_event
from datetime import datetime

def print_message(text):
    print text


t=handle_event.apply_async(countdown=10,link=print_message.s('Done'))  ##HOWTO?

我的确切问题是如何在 Celery 中的 print_message 等本地函数上实现延迟回调 TWISTED 样式。当 handle_Event 方法完成时,它会返回我想要另一个回调方法(print_message)的结果,它是本地的

在 Celery 中还有其他可能的设计工作流程吗?

谢谢

JR

【问题讨论】:

    标签: python twisted celery


    【解决方案1】:

    好吧,终于明白了。像 Twisted 风格那样直接在 Celery 客户端中添加回调是不太可能的。但 Celery 支持任务监控功能,使客户端能够监控不同类型的工作事件并在其上添加回调。

    一个简单的任务监视器(Task_Monitor.py)看起来像这样。 (详情可参见Celery实加工文档http://docs.celeryproject.org/en/latest/userguide/monitoring.html#real-time-processing

    Task_Monitor.py

    from celery import Celery
    
    def task_monitor(app):
        state = app.events.State()
    
        def announce_completed_tasks(event):
            state.event(event)
            task = state.tasks.get(event['uuid'])
    
            print('TASK SUCCEEDED: %s[%s] %s' % (task.name, task.uuid, task.info(), ))
    
        with app.connection() as connection:
            recv = app.events.Receiver(connection, handlers={'task-succeeded': announce_completed_tasks})
            recv.capture(limit=None, timeout=None, wakeup=True)
    
    if __name__ == '__main__':
        app = Celery(broker='amqp://guest@REMOTEHOST//')
        task_monitor(app)
    

    Task_Monitor.py 必须作为单独的进程(客户端)运行。除了需要使用 Celery 应用程序(服务器端)配置

    app.conf.CELERY_SEND_EVENTS = TRUE
    

    或在运行 celery 时使用 -E 选项

    以便它发送事件以便监控工作人员。

    【讨论】:

      【解决方案2】:

      我建议对Celery Canvas docs 使用链或类似机制之一。

      来自文档的示例:

      >>> from celery import chain
      >>> from proj.tasks import add, mul
      
      # (4 + 4) * 8 * 10
      >>> res = chain(add.s(4, 4), mul.s(8), mul.s(10))
      proj.tasks.add(4, 4) | proj.tasks.mul(8) | proj.tasks.mul(10)
      >>> res.apply_async()
      

      【讨论】:

      • 感谢您的评论。客户端和任务队列是远程拆分/分布的。当任务远程完成时,我想在客户端添加回调,就像 TWISTE 中支持的那样。在上面的示例中, mul 不是在添加成功完成时运行的远程任务(而不是客户端上的回调函数)吗?
      • 啊,您希望在回调中返回结果。好吧,您需要在调度程序中运行,例如 geventtwistedtulip。有了其中一个,您就可以生成一个等待异步返回 .get() 的 greenlet。
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2015-03-15
      • 1970-01-01
      • 1970-01-01
      • 2019-01-03
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多