多进程修复字典数据类型遇cannot pickle '_thread.lock' object错误原因咨询
问题描述
尝试用multiprocessing批量修复字典值的数据类型时,持续报错:TypeError: cannot pickle '_thread.lock' object。
相关代码如下:
def update_orders(self, created_after: str, marketplace=Marketplaces.US): orders = self._get_orders(created_after, marketplace) order_items = [self._get_order_items(order["AmazonOrderId"]) for order in orders] pool = mp.Pool(3) results = pool.map(self._fix_dtype, [orders, order_items]) pool.close() pool.join() pprint(results) def _fix_dtype(self, value: Any) -> Any: return fix_dtype_cython(value, {"ObjectID": id(self), "Childs": {"AmazonApiManagerID": id(self.api)}}) @cython.boundscheck(False) @cython.wraparound(False) def fix_dtype_cython(value, debug_information): cdef str lower_val cdef int n cdef list list_copy cdef dict dict_copy if isinstance(value, (int, float, bool, type(None), date, datetime)): return value elif isinstance(value, str): if value.isnumeric(): if len(value) >= 10: return value return int(value) try: return float(value) except ValueError: lower_val = value.lower() if lower_val == "false": return False elif lower_val == "true": return True elif lower_val == "" or lower_val == "null" or lower_val == "none": return None else: return value elif isinstance(value, list): list_copy = value n = len(list_copy) if n == 0: return None return [fix_dtype_cython(list_copy[i], debug_information) for i in range(n)] elif isinstance(value, dict): dict_copy = value if len(dict_copy) == 0: return None return {key: fix_dtype_cython(val, debug_information) for key, val in dict_copy.items()} else: logging.critical(f"Failed to fix value '{value}' Therefore, returned None. {debug_information}") return None
最初怀疑是Cython导致的问题,移除Cython后报错依旧。后来发现注释掉__init__方法中的self.client(pymongo.MongoClient实例)后,程序可正常运行,但不清楚具体原因。
__init__代码如下:
def __init__(self, client: pymongo.MongoClient) -> None: self.client = client self.api = AmazonApiManager()
问题原因
Python的multiprocessing.Pool创建子进程时,会通过pickle序列化将主进程的数据传递给子进程。而pymongo.MongoClient内部包含_thread.lock这类无法被pickle序列化的对象(锁、网络连接等资源均不可序列化)。
调用pool.map(self._fix_dtype, ...)时,类实例self会被传递给子进程,而self持有self.client这个MongoClient实例,导致pickle尝试序列化它时失败,抛出cannot pickle '_thread.lock' object错误。
解决方案
有几种可行的处理方式:
- 避免在类实例中持有MongoClient:如果子进程不需要用到MongoClient,不要将其存为类实例属性;如果需要,在子进程内部重新创建MongoClient连接,不要从主进程传递。
- 改用独立函数替代类方法:将
_fix_dtype修改为不依赖类实例的独立函数,仅传入需要处理的数据和必要的调试信息,避免传递整个类实例。 - 使用
starmap传递独立参数:把需要处理的数据和调试信息打包成元组,通过starmap传递给独立函数,不涉及类实例的传递。
以下是修改示例,将_fix_dtype改为独立函数:
def fix_dtype_standalone(value: Any, debug_info: dict) -> Any: return fix_dtype_cython(value, debug_info) # 在update_orders中调用 def update_orders(self, created_after: str, marketplace=Marketplaces.US): orders = self._get_orders(created_after, marketplace) order_items = [self._get_order_items(order["AmazonOrderId"]) for order in orders] # 准备任务,仅传递数据和调试信息,不传递self tasks = [ (orders, {"ObjectID": id(self), "Childs": {"AmazonApiManagerID": id(self.api)}}), (order_items, {"ObjectID": id(self), "Childs": {"AmazonApiManagerID": id(self.api)}}) ] pool = mp.Pool(3) results = pool.starmap(fix_dtype_standalone, tasks) pool.close() pool.join() pprint(results)
这样就不会把包含MongoClient的类实例传递给子进程,避免了pickle序列化不可序列化对象的问题。
内容的提问来源于stack exchange,提问作者Salih Arda Mermer
相关产品推荐
相关产品推荐

