【发布时间】: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
【问题讨论】: