Python 2.7中不可序列化类使用apply_async的可行方案咨询
Python 2.7 类内使用multiprocessing.apply_async的序列化问题解决
我来帮你搞定这个问题,先逐个解答你的疑问:
1. 不可序列化类的实例/类方法能不能用apply_async?
首先得明确:multiprocessing.Pool的所有异步方法(包括apply_async)都要求传递的可调用对象和参数能被序列化(默认用pickle)。
- 实例方法:因为实例方法绑定到具体实例,序列化它的时候需要把整个实例也序列化。如果你的类本身不可序列化(比如包含了不能被pickle的属性,像文件句柄、网络连接这类),那确实没法直接用实例方法。
- 类方法:类方法属于类本身,只要你的类是可pickle的(绝大多数默认类都满足),类方法是可以被序列化传递的。但如果类里有不可序列化的属性,而类方法又需要用到这些属性,那还是会踩坑。
2. 你的代码无输出无报错的原因及解决方案
你的代码在Python 2.7下没反应,核心问题就是实例方法self.execute_chunk需要序列化整个Try_multiprocess实例,而Python 2.7的pickle对实例方法的支持很有限,经常会出现静默失败(子进程崩溃但不把错误传回主进程),所以你看不到报错也看不到输出。
下面给你几个可行的解决方案,按需选择:
方案一:把execute_chunk改成静态方法
静态方法不绑定实例或类,不需要序列化整个实例,直接就能被Pool调用,适合不需要访问类/实例属性的场景:
import datetime from multiprocessing import Pool class Try_multiprocess(): def multyprocese_chunks(self, chunks): pool = Pool(processes=2) site_id = 564 site_st = 564 for chunk_ix, chunk in enumerate(chunks): pool.apply_async(self.execute_chunk, args=(chunk, chunk_ix, site_id, site_st,)) print "{} wait for join".format(datetime.datetime.now()) pool.close() pool.join() print "{} after for join".format(datetime.datetime.now()) @staticmethod def execute_chunk(chunk, chunk_ix, site_id, site_st): print "{} execute_chunk chunk : {} ".format(datetime.datetime.now(), chunk_ix)
方案二:改成类方法(如果需要访问类属性)
如果你的逻辑需要用到类级别的属性,用类方法更合适,只要类本身能被pickle就没问题:
import datetime from multiprocessing import Pool class Try_multiprocess(): # 示例类属性 class_site_id = 564 def multyprocese_chunks(self, chunks): pool = Pool(processes=2) site_st = 564 for chunk_ix, chunk in enumerate(chunks): # 类方法调用时要传递类本身作为第一个参数 pool.apply_async(Try_multiprocess.execute_chunk, args=(Try_multiprocess, chunk, chunk_ix, self.class_site_id, site_st,)) print "{} wait for join".format(datetime.datetime.now()) pool.close() pool.join() print "{} after for join".format(datetime.datetime.now()) @classmethod def execute_chunk(cls, chunk, chunk_ix, site_id, site_st): print "{} execute_chunk chunk : {} , class site id: {}".format(datetime.datetime.now(), chunk_ix, site_id)
方案三:把执行逻辑抽到类外的普通函数
如果你的方法必须依赖实例的属性,但实例又不可序列化,这个方案最稳妥——把需要用到的属性作为参数传递给类外的函数,完全避开实例序列化的问题:
import datetime from multiprocessing import Pool # 把执行逻辑抽到类外 def execute_chunk(chunk, chunk_ix, site_id, site_st): print "{} execute_chunk chunk : {} ".format(datetime.datetime.now(), chunk_ix) class Try_multiprocess(): def multyprocese_chunks(self, chunks): pool = Pool(processes=2) site_id = 564 site_st = 564 for chunk_ix, chunk in enumerate(chunks): pool.apply_async(execute_chunk, args=(chunk, chunk_ix, site_id, site_st,)) print "{} wait for join".format(datetime.datetime.now()) pool.close() pool.join() print "{} after for join".format(datetime.datetime.now())
小技巧:捕获子进程的错误
Python 2.7中这种静默失败很烦人,你可以给apply_async加个error_callback参数,这样就能看到具体的错误信息了:
def handle_error(e): print "子进程出错啦:{}".format(str(e)) # 在调用apply_async时加上 pool.apply_async(self.execute_chunk, args=(chunk, chunk_ix, site_id, site_st,), error_callback=handle_error)
内容的提问来源于stack exchange,提问作者Omer Anisfeld
相关产品推荐
相关产品推荐

