使用multiprocessing.pool.starmap()触发ExceptionWithTraceback错误排查
我正在开发从API获取数据并生成输出文件的代码,为提升处理效率采用multiprocessing库实现任务并行化。目前代码能返回部分结果,但会随机(返回2-4条结果后)抛出如下错误:
multiprocessing.pool.MaybeEncodingError: Error sending result: '<multiprocessing.pool.ExceptionWithTraceback object at 0x000002CA8BA44BD0>'. Reason: 'PicklingError("Can't pickle <function HTTPResponse.getheaders at 0x000002CA8BA06980>: it's not the same object as http.client.HTTPResponse.getheaders")'
相关代码
function1(主函数)
def function1 (origin: str, departureDate: datetime, returnDate: datetime): if __name__ == '__main__': destinations = ["PAR", "LON", "AMS", "ROM", "BER", "BRU", "MUC"] res = [] with Pool() as pool: pool.starmap(function2, zip(list(repeat(origin,7)), args, list(repeat(departureDate, 7)), list(repeat(returnDate, 7)))) x = [r[0] for r in res] with open("res.json", "w") as outfile: json.dump(x, outfile, indent=4, sort_keys=True) function1("MAD", datetime.datetime(2023,3,1), datetime.datetime(2023,3,8))
function2(任务处理函数)
def function2 (origin: str, destination: str, departureDate: datetime, returnDate: datetime): vuelo = v.function3(origin, destination, departureDate, returnDate) habitacion = h.function4(destino, departureDate, returnDate) if(len(habitacion)!=0): habitacion = habitacion[0][0] result = { "destination": destination, "totalPrice": (float(room["offers"][0]["price"]["total"]) + float(flight["price"]["total"])), "hotelPrice": float(room["offers"][0]["price"]["total"])/2, "departurePrice": float(flight["price"]["departure"]), "returnPrice": float(flight["price"]["return"]) print(result) #this is to test the output return result
function3(航班API调用)
def function3 (origin: str, destination: str, departureDate: datetime, returnDate: datetime): data = {#request parameters} try: response = provider.post(data) return response.data except ResponseError as error: raise error
function4(酒店列表处理)
def function4 (destination: str, departureDate: datetime, returnDate: datetime): #if __name__ == '__main__': hotels = GetHotels(destination) res = [] for hotel in hotels: res.append(function5 (hotel,departureDate,returnDate)) rooms = list(filter(None, res[0])) return rooms
function5(酒店价格API调用)
def function5 (hotels, departureDate: datetime, returnDate: datetime): roomList = [] try: hotel_offers = provider.get(#params) if(hotel_offers.data != []): roomList.append(hotel_offers.data) except Exception: 0 if roomList == []: roomList.append("") return roomList
完整错误回溯
Traceback (most recent call last): File "c:\Users\user\Documents\folder\folder\user\function2.py", line 49, in <module> function1("MAD", datetime.datetime(2023,3,1), datetime.datetime(2023,3,8)) File "c:\Users\user\Documents\folder\folder\user\function2.py", line 42, in function1 pool.starmap(function2, zip(list(repeat(origin,7)), args, list(repeat(departureDate, 7)), list(repeat(returnDate, 7)))) File "C:\Python311\Lib\multiprocessing\pool.py", line 375, in starmap return self._map_async(func, iterable, starmapstar, chunksize).get() ^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^^ File "C:\Python311\Lib\multiprocessing\pool.py", line 774, in get raise self._value multiprocessing.pool.MaybeEncodingError: Error sending result: '<multiprocessing.pool.ExceptionWithTraceback object at 0x00000246DDDA4C50>'. Reason: 'PicklingError("Can't pickle <function HTTPResponse.getheaders at 0x00000246DDD66980>: it's not the same object as http.client.HTTPResponse.getheaders")'
我已尝试多种并行化实现方式,甚至重写了处理网络数据获取的function3、function4和function5,但问题仍未解决,无法定位错误根源,特此寻求帮助。
问题根源与解决方案
错误核心原因
Python多进程间传递数据依赖pickle序列化,但HTTPResponse.getheaders这类对象无法被正常序列化——要么是API返回的响应对象(或其嵌套结构)包含了未序列化的HTTPResponse实例,要么是第三方库的响应对象内部持有了不可pickle的方法引用。同时代码还存在变量名错误、结果处理逻辑漏洞等问题。
具体修复步骤
确保返回纯可序列化数据
在function3和function5中,避免返回第三方库的响应对象,确保只传递Python原生类型(字典、列表、字符串等)。如果response.data仍包含不可pickle的元素,手动转换为纯字典:import json def function3(origin: str, destination: str, departureDate: datetime, returnDate: datetime): data = {#request parameters} try: response = provider.post(data) # 强制转换为纯Python字典,消除嵌套的不可序列化对象 return json.loads(json.dumps(response.data)) except ResponseError as error: raise error修正变量名错误
function1中未定义的args改为destinationsfunction2中room改为habitacion,flight改为vuelo,与实际定义的变量名匹配
修复
function4的结果处理逻辑
原代码只取第一个酒店的结果,改为收集所有有效酒店数据:def function4(destination: str, departureDate: datetime, returnDate: datetime): hotels = GetHotels(destination) res = [] for hotel in hotels: res.append(function5(hotel, departureDate, returnDate)) # 过滤所有空结果并合并 rooms = [] for item in res: rooms.extend(list(filter(None, item))) return rooms调整多进程代码结构
把if __name__ == '__main__':移到最外层,避免子进程重复执行初始化代码,同时接收starmap的返回结果:def function1(origin: str, departureDate: datetime, returnDate: datetime): destinations = ["PAR", "LON", "AMS", "ROM", "BER", "BRU", "MUC"] with Pool() as pool: # 修正参数传递,接收返回结果 results = pool.starmap(function2, zip(repeat(origin, len(destinations)), destinations, repeat(departureDate, len(destinations)), repeat(returnDate, len(destinations)))) # 直接写入结果 with open("res.json", "w") as outfile: json.dump(results, outfile, indent=4, sort_keys=True) if __name__ == '__main__': function1("MAD", datetime.datetime(2023,3,1), datetime.datetime(2023,3,8))优化异常处理
function5中无效的异常处理改为明确的错误记录:def function5(hotels, departureDate: datetime, returnDate: datetime): roomList = [] try: hotel_offers = provider.get(#params) if hotel_offers.data: roomList.extend(hotel_offers.data) # 用extend避免嵌套列表 except Exception as e: print(f"获取酒店[{hotels}]数据失败: {str(e)}") return roomList
内容的提问来源于stack exchange,提问作者kogota984

