pybind11封装C++函数结合多进程时pool.join()挂起问题排查
问题:Python多进程调用pybind11封装的C++函数在大数据量下挂起
我尝试在Python中并行运行C++实现的编辑距离计算函数,通过pybind11封装该函数后,在Python中定义包装函数,使用multiprocessing.pool.map_async()进行调用。
本地笔记本上处理约100万条字符串时运行正常,但在远程服务器(配备500GB内存,足以容纳2500万条待处理字符串)上执行时,总会在pool.join()处挂起。最终会剩下1个进程,占用0%CPU和约20%内存,12小时以上后进程被杀死,并出现如下资源泄漏警告:
/usr/lib/python3.8/multiprocessing/resource_tracker.py:216: UserWarning: resource_tracker: There appear to be 6 leaked semaphore objects to clean up at shutdown
warnings.warn('resource_tracker: There appear to be %d 'w
笔记本和服务器均运行Ubuntu 20.04,Python版本为3.8.10,仅待处理字符串列表规模不同。
相关代码
C++代码
#include <pybind11/pybind11.h> #include <pybind11/stl.h> #include <vector> #include <string> #include <iostream> #include <string.h> #include <sstream> using namespace std; bool LevenshteinDistanceCutoff(const char* str1, const char* str2, const uint_fast16_t l1, const uint_fast16_t l2, const uint_fast16_t cutoff, uint_fast16_t** d) { // 代码省略 } vector<vector<uint_fast32_t>> CheckPairwiseLevensteinDistOfList(vector<string> &str_list, vector<uint_fast16_t> &strLens, const uint_fast16_t cutoff) { uint_fast32_t size = str_list.size(); cout << "length of str_list = " << size << endl; uint_fast16_t maxLen = strLens[size-1]; uint_fast16_t penultimateLen = strLens[size-2]; cout << "maxLen " << maxLen << endl; cout << "penultimateLen " << penultimateLen << endl; vector<vector<uint_fast32_t>> indPotMatches(size); cout << "current size indPotMatches " << indPotMatches.size() << endl; cout << "max size indPotMatches " << indPotMatches.max_size() << endl; uint_fast16_t** d = new uint_fast16_t*[penultimateLen+1]; for (uint_fast16_t i=0; i<penultimateLen+1; ++i) { d[i] = new uint_fast16_t[maxLen+1]; } for (uint_fast32_t j=0; j<size-1; ++j) { for (uint_fast32_t k=j+1; k<size; ++k) { if (LevenshteinDistanceCutoff(str_list[j].c_str(), str_list[k].c_str(), strLens[j], strLens[k], cutoff, d)) { indPotMatches[j].push_back(k); } } } for (uint_fast16_t i=0; i<penultimateLen+1; ++i) { delete[] d[i]; } delete[] d; return indPotMatches; } PYBIND11_MODULE(CHelpers, m) { m.def("CheckPairwiseLevensteinDistOfList", &CheckPairwiseLevensteinDistOfList, pybind11::return_value_policy::take_ownership); }
Python代码
import multiprocessing as mp import sys sys.path.append("C++编译生成的so文件所在目录") import CHelpers sys.path.remove("C++编译生成的so文件所在目录") NBR_PROC = 12 # C++函数的Python包装器 def CheckPairwiseLEvensteinDistOFListWrapper(arg): locStartValRange = arg[0] locEndValRange = arg[1] locDisplayNameLst = arg[2] locDisplayNameLensLst = arg[3] locCutoff = arg[4] tmpInd = CHelpers.CheckPairwiseLevensteinDistOfList(locDisplayNameLst, locDisplayNameLensLst, locCutoff) return locStartValRange, locEndValRange, locCutoff, tmpInd # 全局变量等其他代码 if __name__ == '__main__': # 准备传入多进程池的参数列表 lstArgsForMap = [(startValRangeLst[j], endValRangeLst[j], displayNameLst[startValRangeLst[j]:endValRangeLst[j]].copy(), displayNameLensLst[startValRangeLst[j]:endValRangeLst[j]].copy(), cutoffLst[j]) for j in range(len(startValRangeLst))] mp.set_start_method('spawn') pool = mp.Pool(processes = NBR_PROC, maxtasksperchild=1) res = pool.map_async( CheckPairwiseLEvensteinDistOFListWrapper, lstArgsForMap) pool.close() pool.join() res = res.get() # 结果后处理代码
已尝试的排查措施
- 将进程启动方式从fork改为spawn,避免fork引发的内存问题
- 添加
maxtasksperchild=1参数,确保每个任务使用全新进程 - 为pybind11绑定函数设置
return_value_policy::take_ownership,让Python负责管理返回的vector内存
我是pybind11和multiprocessing的新手,不清楚后续排查方向,最坏打算会改用C++线程重写,希望得到帮助。
内容的提问来源于stack exchange,提问作者tarquaeron
相关产品推荐
相关产品推荐

