我正在尝试使用 Celery 作为 Twisted 应用程序的控制 channel 。我的 Twisted 应用程序是一个抽象层,它为各种本地运行的进程(通过 ProcessProtocol)提供标准接口(interface)。我想使用 Celery 来远程控制它——AMQP 似乎是从中央位置控制许多 Twisted 应用程序的理想方法,我想利用 Celery 基于任务的功能,例如任务重试、子任务等
这并没有像我计划的那样工作,我希望有人能帮助我指明正确的方向以实现这一目标。
我在运行脚本时试图实现的行为是:
“稍微修改过的 celery ”是 celeryd有一个小的修改,允许任务通过 self.app.twisted 访问 Twisted react 器,并通过 self.app.process 访问生成的进程。为了简单起见,我使用了 Celery 的“单独”进程池实现,它不会为任务 worker 创建新进程。
当我尝试使用 Celery 任务来初始化 ProcessProtocol(即启动外部进程)时,我的问题就出现了。进程正确启动,但 ProcessProtocol 的 childDataReceived 永远不会被调用。我认为这与未正确继承/设置文件描述符有关。
下面是一些示例代码,基于 ProcessProtocol 文档中的“wc”示例。它包括两个 Celery 任务——一个用于启动 wc 进程,另一个用于计算某些文本中的单词(使用先前启动的 wc 进程)。
这个示例相当人为设计,但如果我能让它正常工作,它将作为实现我的 ProcessProtocols 的良好起点,这些 ProcessProtocols 是长期运行的进程,将响应写入标准输入的命令。
我首先通过运行 Celery 守护进程来测试它:
python2.6 mycelery.py -l info -P solo
然后,在另一个窗口中,运行发送两个任务的脚本:
python2.6 命令测试.py
command_test.py 的预期行为是执行两个命令 - 一个启动 wc 进程,另一个向 CountWordsTask 发送一些文本。实际发生的是:
任何人都可以阐明这一点,或者就如何最好地使用 Celery 作为 Twisted ProcessProtocols 的控制 channel 提供一些建议吗?
为 Celery 编写一个 Twisted-backed ProcessPool 实现会更好吗?我通过 reactor.callLater 调用 WorkerCommand.execute_from_commandline 的方法是否是确保一切都发生在 Twisted 线程内的正确方法?
我已经阅读了有关 AMPoule 的资料,我认为它可以提供其中的一些功能,但如果可能的话我想坚持使用 Celery,因为我在我的应用程序的其他部分使用它。
任何帮助或协助将不胜感激!
from functools import partial
from celery.app import App
from celery.bin.celeryd import WorkerCommand
from twisted.internet import reactor
class MyCeleryApp(App):
def __init__(self, twisted, *args, **kwargs):
self.twisted = twisted
super(MyCeleryApp, self).__init__(*args, **kwargs)
def main():
get_my_app = partial(MyCeleryApp, reactor)
worker = WorkerCommand(get_app=get_my_app)
reactor.callLater(1, worker.execute_from_commandline)
reactor.run()
if __name__ == '__main__':
main()
from twisted.internet import protocol
from twisted.internet.defer import Deferred
class WCProcessProtocol(protocol.ProcessProtocol):
def __init__(self, text):
self.text = text
self._waiting = {} # Dict to contain deferreds, keyed by command name
def connectionMade(self):
if 'startup' in self._waiting:
self._waiting['startup'].callback('process started')
def outReceived(self, data):
fieldLength = len(data) / 3
lines = int(data[:fieldLength])
words = int(data[fieldLength:fieldLength*2])
chars = int(data[fieldLength*2:])
self.transport.loseConnection()
self.receiveCounts(lines, words, chars)
if 'countWords' in self._waiting:
self._waiting['countWords'].callback(words)
def processExited(self, status):
print 'exiting'
def receiveCounts(self, lines, words, chars):
print >> sys.stderr, 'Received counts from wc.'
print >> sys.stderr, 'Lines:', lines
print >> sys.stderr, 'Words:', words
print >> sys.stderr, 'Characters:', chars
def countWords(self, text):
self._waiting['countWords'] = Deferred()
self.transport.write(text)
return self._waiting['countWords']
from celery.task import Task
from protocol import WCProcessProtocol
from twisted.internet.defer import Deferred
from twisted.internet import reactor
class StartProcTask(Task):
def run(self):
self.app.proc = WCProcessProtocol('testing')
self.app.proc._waiting['startup'] = Deferred()
self.app.twisted.spawnProcess(self.app.proc,
'wc',
['wc'],
usePTY=True)
return self.app.proc._waiting['startup']
class CountWordsTask(Task):
def run(self):
return self.app.proc.countWords('test test')
最佳答案
Celery 可能会在等待来自网络的新消息时阻塞。由于您在一个单线程进程中与 Twisted react 器一起运行它,因此它会阻止 react 器运行。这将禁用大部分 Twisted,这需要 react 堆实际运行(您调用了 reactor.run,但由于 Celery 阻止了它,它实际上没有运行)。
reactor.callLater 只是延迟了 Celery 的启动。一旦 Celery 启动,它仍然会阻塞 react 器。
您需要避免的问题是阻塞 react 堆。
一种解决方案是在一个线程中运行 Celery,在另一个线程中运行 react 器。使用 reactor.callFromThread 从 Celery 线程向 Twisted 发送消息(“在 react 器线程中调用函数”)。如果您需要从 Twisted 线程将消息发送回 Celery,请使用 Celery 等效项。
另一种解决方案是将 Celery 协议(protocol)(AMQP? - 请参阅 txAMQP)作为原生 Twisted 库来实现,并使用它来无阻塞地处理 Celery 消息。
关于python - 使用 Celery 作为 Twisted 应用程序的控制 channel ,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/8137277/
我正在学习如何使用Nokogiri,根据这段代码我遇到了一些问题:require'rubygems'require'mechanize'post_agent=WWW::Mechanize.newpost_page=post_agent.get('http://www.vbulletin.org/forum/showthread.php?t=230708')puts"\nabsolutepathwithtbodygivesnil"putspost_page.parser.xpath('/html/body/div/div/div/div/div/table/tbody/tr/td/div
我有一个Ruby程序,它使用rubyzip压缩XML文件的目录树。gem。我的问题是文件开始变得很重,我想提高压缩级别,因为压缩时间不是问题。我在rubyzipdocumentation中找不到一种为创建的ZIP文件指定压缩级别的方法。有人知道如何更改此设置吗?是否有另一个允许指定压缩级别的Ruby库? 最佳答案 这是我通过查看rubyzip内部创建的代码。level=Zlib::BEST_COMPRESSIONZip::ZipOutputStream.open(zip_file)do|zip|Dir.glob("**/*")d
类classAprivatedeffooputs:fooendpublicdefbarputs:barendprivatedefzimputs:zimendprotecteddefdibputs:dibendendA的实例a=A.new测试a.foorescueputs:faila.barrescueputs:faila.zimrescueputs:faila.dibrescueputs:faila.gazrescueputs:fail测试输出failbarfailfailfail.发送测试[:foo,:bar,:zim,:dib,:gaz].each{|m|a.send(m)resc
很好奇,就使用rubyonrails自动化单元测试而言,你们正在做什么?您是否创建了一个脚本来在cron中运行rake作业并将结果邮寄给您?git中的预提交Hook?只是手动调用?我完全理解测试,但想知道在错误发生之前捕获错误的最佳实践是什么。让我们理所当然地认为测试本身是完美无缺的,并且可以正常工作。下一步是什么以确保他们在正确的时间将可能有害的结果传达给您? 最佳答案 不确定您到底想听什么,但是有几个级别的自动代码库控制:在处理某项功能时,您可以使用类似autotest的内容获得关于哪些有效,哪些无效的即时反馈。要确保您的提
假设我做了一个模块如下:m=Module.newdoclassCendend三个问题:除了对m的引用之外,还有什么方法可以访问C和m中的其他内容?我可以在创建匿名模块后为其命名吗(就像我输入“module...”一样)?如何在使用完匿名模块后将其删除,使其定义的常量不再存在? 最佳答案 三个答案:是的,使用ObjectSpace.此代码使c引用你的类(class)C不引用m:c=nilObjectSpace.each_object{|obj|c=objif(Class===objandobj.name=~/::C$/)}当然这取决于
我正在尝试使用ruby和Savon来使用网络服务。测试服务为http://www.webservicex.net/WS/WSDetails.aspx?WSID=9&CATID=2require'rubygems'require'savon'client=Savon::Client.new"http://www.webservicex.net/stockquote.asmx?WSDL"client.get_quotedo|soap|soap.body={:symbol=>"AAPL"}end返回SOAP异常。检查soap信封,在我看来soap请求没有正确的命名空间。任何人都可以建议我
我需要在客户计算机上运行Ruby应用程序。通常需要几天才能完成(复制大备份文件)。问题是如果启用sleep,它会中断应用程序。否则,计算机将持续运行数周,直到我下次访问为止。有什么方法可以防止执行期间休眠并让Windows在执行后休眠吗?欢迎任何疯狂的想法;-) 最佳答案 Here建议使用SetThreadExecutionStateWinAPI函数,使应用程序能够通知系统它正在使用中,从而防止系统在应用程序运行时进入休眠状态或关闭显示。像这样的东西:require'Win32API'ES_AWAYMODE_REQUIRED=0x0
关闭。这个问题是opinion-based.它目前不接受答案。想要改进这个问题?更新问题,以便editingthispost可以用事实和引用来回答它.关闭4年前。Improvethisquestion我想在固定时间创建一系列低音和高音调的哔哔声。例如:在150毫秒时发出高音调的蜂鸣声在151毫秒时发出低音调的蜂鸣声200毫秒时发出低音调的蜂鸣声250毫秒的高音调蜂鸣声有没有办法在Ruby或Python中做到这一点?我真的不在乎输出编码是什么(.wav、.mp3、.ogg等等),但我确实想创建一个输出文件。
我在我的项目目录中完成了compasscreate.和compassinitrails。几个问题:我已将我的.sass文件放在public/stylesheets中。这是放置它们的正确位置吗?当我运行compasswatch时,它不会自动编译这些.sass文件。我必须手动指定文件:compasswatchpublic/stylesheets/myfile.sass等。如何让它自动运行?文件ie.css、print.css和screen.css已放在stylesheets/compiled。如何在编译后不让它们重新出现的情况下删除它们?我自己编译的.sass文件编译成compiled/t
对于具有离线功能的智能手机应用程序,我正在为Xml文件创建单向文本同步。我希望我的服务器将增量/差异(例如GNU差异补丁)发送到目标设备。这是计划:Time=0Server:hasversion_1ofXmlfile(~800kiB)Client:hasversion_1ofXmlfile(~800kiB)Time=1Server:hasversion_1andversion_2ofXmlfile(each~800kiB)computesdeltaoftheseversions(=patch)(~10kiB)sendspatchtoClient(~10kiBtransferred)Cl