没有足够的信息可以确定,但问题很可能是
Slave.do_work
正在引发未处理的异常。(您的代码中有很多行可以在各种不同的条件下做到这一点。)
当您这样做时,子进程将直接退出。
在POSIX系统上,完整的细节有点复杂,但在简单的情况下(这里有),退出的子进程将作为
<defunct>
处理直到收获(因为父对象
wait
s在上面,或者退出)。由于您的父代码在队列结束之前不会等待子代码,所以这正是发生的情况。
所以,有一个简单的管道胶带修复:
def do_work(self):
self.log(str(os.getpid()))
while True:
try:
# the rest of your code
except Exception as e:
self.log("something appropriate {}".format(e))
# you may also want to post a reply back to the parent
你可能还想打破巨大的
try
分成不同的阶段,这样你就可以区分所有可能出现问题的不同阶段(尤其是如果其中一些阶段意味着你需要回复,而另一些阶段则意味着你不需要回复)。
然而,看起来你试图做的是复制
multiprocessing.Pool
,但有几个地方错过了酒吧。这就提出了一个问题:为什么不直接使用
Pool
首先?然后,您可以通过使用
map
家庭方法。例如,你的整个
Master.run
可以减少到:
self.init()
pool = multiprocessing.Pool(Master.SLAVE_COUNT, initializer=slave_setup)
pool.map(slave_job, tables)
pool.join()
这将为您处理异常,并允许您在以后需要时返回值/异常,并且允许您使用内置的
logging
库,而不是试图构建自己的库,等等。而且只需要几十行小的代码更改就可以
Slave
,然后你就完了。
如果您想从作业中提交新作业,最简单的方法可能是使用
Future
-基于API(它可以扭转局面,使未来的结果成为焦点,使池/执行器成为提供它们的哑对象,而不是使池成为焦点,并使结果成为它返回的哑对象),但有多种方法
水塘
也例如,现在,你没有从每一份工作中返回任何东西,所以,你可以只返回一个列表
tables
执行。下面是一个简单的例子,展示了如何做到这一点:
import multiprocessing
def foo(x):
print(x, x**2)
return list(range(x))
if __name__ == '__main__':
pool = multiprocessing.Pool(2)
jobs = [5]
while jobs:
jobs, oldjobs = [], jobs
for job in oldjobs:
jobs.extend(pool.apply(foo, [job]))
pool.close()
pool.join()
显然,你可以通过将整个循环替换为,例如,一个列表理解,来浓缩这一点
itertools.chain
,你可以通过向每个作业传递“一个提交者”对象并添加到该对象中,而不是返回一个新作业列表,等等,来让它看起来更干净。但我想让它尽可能明确,以显示它的内容有多少。
无论如何,如果你认为显式队列更容易理解和管理,那就去做吧
multiprocessing.worker
和/或
concurrent.futures.ProcessPoolExecutor
看看你自己需要做什么。这并没有那么难,但有足够多的事情你可能会出错(就我个人而言,当我自己尝试做这样的事情时,我总是忘记至少一个边缘情况),那就是看代码才能把它做好。
或者,这似乎是你不能使用的唯一原因
concurrent.futures.ProcessPoolExecutor
这里需要初始化一些每个进程的状态(
boto.s3.key.Key
,
MySqlWrap
等等),这可能是非常好的缓存原因。(如果这涉及到web服务查询、数据库连接等,你当然不想每次任务都这样做一次!)但有几种不同的方法可以解决这个问题。
但你可以细分
ProcessPoolExecutor
并覆盖未记录的函数
_adjust_process_count
(参见
the source
因为它有多简单)来传递你的设置函数,这就是你所要做的。
或者你可以混搭。包裹
将来
从…起
concurrent.futures
围绕
AsyncResult
从…起
multiprocessing
.