- iOS/Objective-C 元类和类别
- objective-c - -1001 错误,当 NSURLSession 通过 httpproxy 和/etc/hosts
- java - 使用网络类获取 url 地址
- ios - 推送通知中不播放声音
我正在使用 multiprocessing.queues.JoinableQueue 如下:
一个非常长时间运行的线程(多天)轮询 XML 的 API。执行此操作的线程只是将 XML 解析为对象并将它们插入队列。
处理每个对象比解析 XML 花费的时间要多得多,并且绝不依赖于从 API 读取的线程。因此,这种多处理的实现非常简单。
创建和清理进程的代码在这里:
def queueAdd(self, item):
try:
self.queue.put(item)
except AssertionError:
#queue has been closed, remake it (let the other GC)
logger.warn('Queue closed early.')
self.queue = BufferQueue(ctx=multiprocessing.get_context())
self.queue.put(item)
except BrokenPipeError:
#workaround for pipe issue
logger.warn('Broken pipe, Forcing creation of new queue.')
# all reading procesess should suicide and new ones spawned.
self.queue = BufferQueue(ctx=multiprocessing.get_context())
# address = 'localhost'
# if address in multiprocessing.managers.BaseProxy._address_to_local:
# del BaseProxy._address_to_local[address][0].connection
self.queue.put(item)
except Exception as e:
#general thread exception.
logger.error('Buffer queue exception %s' % e)
#TODO: continue trying/trap exceptions?
raise
# check for finished consumers and clean them up before we check to see
# if we need to add additional consumers.
for csmr in self.running:
if not csmr.is_alive():
debug('Child dead, releasing.')
self.running.remove(csmr)
#see if we should start a consumer...
# TODO: add min/max processes (default and override)
if not self.running:
debug('Spawning consumer.')
new_consumer = self.consumer(
queue=self.queue,
results_queue=self.results_queue,
response_error=self.response_error)
new_consumer.start()
self.running.append(new_consumer)
消费者进程控制循环也非常简单:
def run(self):
'''Consumes the queue in the framework, passing off each item to the
ItemHandler method.
'''
while True:
try:
item = self.queue.get(timeout=3)
#the base class just logs this stuff
rval = self.singleItemHandler(item)
self.queue.task_done()
if rval and self.results_queue:
self.results_queue.put(rval)
except queue.Empty:
logging.debug('Queue timed out after 3 seconds.')
break
except EOFError:
logging.info(
'%s has finished consuming queue.' % (__class__.__name__))
break
except Exception as e:
#general thread exception.
self.logger.error('Consumer exception %s' % e)
#TODO: continue trying/trap exceptions?
raise
一段时间后(大约一个小时的成功处理),我收到一条日志消息,表明消费者进程因超时而终止 DEBUG:root:Queue timed out after 3 seconds.
,但队列仍处于打开状态,并且显然仍在由原始线程写入。该线程似乎并不认为消费者进程已终止(请参阅 queueAdd 方法)并且不会尝试启动一个新进程。队列似乎并不为空,只是从中读取似乎已超时。
我不明白为什么经理认为 child 还活着。
编辑
由于对 BrokenPipeError 记录方式的代码更改以及删除断开的连接清理,我修改了原始问题。我认为这是一个单独的问题。
最佳答案
问题是由 multiprocessing.Queue 的微妙现实引起的。任何调用 queue.put 的进程都将运行一个写入命名管道的后台线程。
在我的特殊情况下,虽然没有大量数据被发布到结果队列(由于某种原因无法处理的项目),但它仍然足以“填满”管道并导致消费者无法退出,即使它没有运行。这反过来会导致写入队列缓慢填满。
解决方案是我为 API 调用的下一次迭代修改了我的非阻塞调用,以读取目前为止所有可用的结果,但最后一次(阻塞)调用除外,以确保获取所有结果。
def finish(self, block=True, **kwargs):
'''
Notifies the buffer that we are done filling it.
This command binds to any processes still running and lets them
finish and then copies and flushes the managed results list.
'''
# close the queue and wait until it is consumed
if block:
self.queue.close()
self.queue.join_thread()
# make sure the consumers are done consuming the queue
for csmr in self.running:
#get everything on the results queue right now.
try:
while csmr.is_alive():
self.results_list.append(
self.results_queue.get(timeout=0.5))
self.results_queue.task_done()
except queue.Empty:
if csmr.is_alive():
logger.warn('Result queue empty but consumer alive.')
logger.warn('joining %s.' % csmr.name)
csmr.join()
del self.running[:]
if self.callback:
return self.callback(self.results_list)
else:
#read results immediately available.
try:
while True:
self.results_list.append(self.results_queue.get_nowait())
self.results_queue.task_done()
except queue.Empty:
#got everything on the queue so far
pass
return self.results_list
关于python - multiprocessing.Queue 似乎消失了?操作系统(管道破坏)与 Python?,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/38637282/
我正在处理一组标记为 160 个组的 173k 点。我想通过合并最接近的(到 9 或 10 个组)来减少组/集群的数量。我搜索过 sklearn 或类似的库,但没有成功。 我猜它只是通过 knn 聚类
我有一个扁平数字列表,这些数字逻辑上以 3 为一组,其中每个三元组是 (number, __ignored, flag[0 or 1]),例如: [7,56,1, 8,0,0, 2,0,0, 6,1,
我正在使用 pipenv 来管理我的包。我想编写一个 python 脚本来调用另一个使用不同虚拟环境(VE)的 python 脚本。 如何运行使用 VE1 的 python 脚本 1 并调用另一个 p
假设我有一个文件 script.py 位于 path = "foo/bar/script.py"。我正在寻找一种在 Python 中通过函数 execute_script() 从我的主要 Python
这听起来像是谜语或笑话,但实际上我还没有找到这个问题的答案。 问题到底是什么? 我想运行 2 个脚本。在第一个脚本中,我调用另一个脚本,但我希望它们继续并行,而不是在两个单独的线程中。主要是我不希望第
我有一个带有 python 2.5.5 的软件。我想发送一个命令,该命令将在 python 2.7.5 中启动一个脚本,然后继续执行该脚本。 我试过用 #!python2.7.5 和http://re
我在 python 命令行(使用 python 2.7)中,并尝试运行 Python 脚本。我的操作系统是 Windows 7。我已将我的目录设置为包含我所有脚本的文件夹,使用: os.chdir("
剧透:部分解决(见最后)。 以下是使用 Python 嵌入的代码示例: #include int main(int argc, char** argv) { Py_SetPythonHome
假设我有以下列表,对应于及时的股票价格: prices = [1, 3, 7, 10, 9, 8, 5, 3, 6, 8, 12, 9, 6, 10, 13, 8, 4, 11] 我想确定以下总体上最
所以我试图在选择某个单选按钮时更改此框架的背景。 我的框架位于一个类中,并且单选按钮的功能位于该类之外。 (这样我就可以在所有其他框架上调用它们。) 问题是每当我选择单选按钮时都会出现以下错误: co
我正在尝试将字符串与 python 中的正则表达式进行比较,如下所示, #!/usr/bin/env python3 import re str1 = "Expecting property name
考虑以下原型(prototype) Boost.Python 模块,该模块从单独的 C++ 头文件中引入类“D”。 /* file: a/b.cpp */ BOOST_PYTHON_MODULE(c)
如何编写一个程序来“识别函数调用的行号?” python 检查模块提供了定位行号的选项,但是, def di(): return inspect.currentframe().f_back.f_l
我已经使用 macports 安装了 Python 2.7,并且由于我的 $PATH 变量,这就是我输入 $ python 时得到的变量。然而,virtualenv 默认使用 Python 2.6,除
我只想问如何加快 python 上的 re.search 速度。 我有一个很长的字符串行,长度为 176861(即带有一些符号的字母数字字符),我使用此函数测试了该行以进行研究: def getExe
list1= [u'%app%%General%%Council%', u'%people%', u'%people%%Regional%%Council%%Mandate%', u'%ppp%%Ge
这个问题在这里已经有了答案: Is it Pythonic to use list comprehensions for just side effects? (7 个答案) 关闭 4 个月前。 告
我想用 Python 将两个列表组合成一个列表,方法如下: a = [1,1,1,2,2,2,3,3,3,3] b= ["Sun", "is", "bright", "June","and" ,"Ju
我正在运行带有最新 Boost 发行版 (1.55.0) 的 Mac OS X 10.8.4 (Darwin 12.4.0)。我正在按照说明 here构建包含在我的发行版中的教程 Boost-Pyth
学习 Python,我正在尝试制作一个没有任何第 3 方库的网络抓取工具,这样过程对我来说并没有简化,而且我知道我在做什么。我浏览了一些在线资源,但所有这些都让我对某些事情感到困惑。 html 看起来
我是一名优秀的程序员,十分优秀!