Amazon Sagemaker中pool.map()多CPU实例运行异常求助
多CPU SageMaker实例上多进程嵌入计算停滞问题
代码在2核Amazon SageMaker实例运行正常,但在更高CPU配置的实例上会陷入停滞,无进度输出。终止内核后触发KeyboardInterrupt错误,同时生成大量cores核心转储文件。
原代码
import sentence_transformers import multiprocessing from tqdm import tqdm from multiprocessing import Pool import numpy as np embedding_model = sentence_transformers.SentenceTransformer('sentence-transformers/all-mpnet-base-v2') data = [[100227, 7382501.0, 'view', 30065006, False, ''], [100227, 7382501.0, 'view', 57072062, True, ''], [100227, 7382501.0, 'view', 66405922, True, ''], [100227, 7382501.0, 'view', 5221475, False, ''], [100227, 7382501.0, 'view', 63283995, True, '']] df_text = dict() df_text[7382501] = {'title': 'The Geography of the Internet Industry, Venture Capital, Dot-coms, and Local Knowledge - MATTHEW A. ZOOK', 'abstract': '23', 'highlight': '12'} df_text[30065006] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[57072062] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[66405922] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[5221475] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[63283995] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} # Define the function to be executed in parallel def process_data(chunk): results = [] for row in chunk: print(row[0]) work_id = row[1] mentioning_work_id = row[3] print(work_id) if work_id in df_text and mentioning_work_id in df_text: title1 = df_text[work_id]['title'] title2 = df_text[mentioning_work_id]['title'] embeddings_title1 = embedding_model.encode(title1,convert_to_numpy=True) embeddings_title2 = embedding_model.encode(title2,convert_to_numpy=True) similarity = np.matmul(embeddings_title1, embeddings_title2.T) results.append([row[0],row[1],row[2],row[3],row[4],similarity]) else: continue return results # Define the number of CPU cores to use num_cores = multiprocessing.cpu_count() # Split the data into chunks chunk_size = len(data) // num_cores chunks = [data[i:i+chunk_size] for i in range(0, len(data), chunk_size)] # Create a pool of worker processest pool = multiprocessing.Pool(processes=num_cores) results = [] with tqdm(total=len(data)) as pbar: for i, result_chunk in enumerate(pool.map(process_data, chunks)): # Update the progress bar pbar.update() # Add the results to the list results += result_chunk # Concatenate the results final_result = results
错误栈信息
--------------------------------------------------------------------------- KeyboardInterrupt Traceback (most recent call last) <ipython-input-18-19449c86abd3> in <module> 1 results = [] 2 with tqdm(total=len(chunks)) as pbar: ----> 3 for i, result_chunk in enumerate(pool.map(process_data, chunks)): 4 # Update the progress bar 5 pbar.update() /opt/conda/lib/python3.7/multiprocessing/pool.py in map(self, func, iterable, chunksize) 266 in a list that is returned. 267 ''' --> 268 return self._map_async(func, iterable, mapstar, chunksize).get() 269 270 def starmap(self, func, iterable, chunksize=None): /opt/conda/lib/python3.7/multiprocessing/pool.py in get(self, timeout) 649 650 def get(self, timeout=None): --> 651 self.wait(timeout) 652 if not self.ready(): 653 raise TimeoutError /opt/conda/lib/python3.7/multiprocessing/pool.py in wait(self, timeout) 646 647 def wait(self, timeout=None): --> 648 self._event.wait(timeout) 649 650 def get(self, timeout=None): /opt/conda/lib/python3.7/threading.py in wait(self, timeout) 550 signaled = self._flag 551 if not signaled: --> 552 signaled = self._cond.wait(timeout) 553 return signaled 554 /opt/conda/lib/python3.7/threading.py in wait(self, timeout) 294 try: # restore state no matter what (e.g., KeyboardInterrupt) 295 if timeout is None: --> 296 waiter.acquire() 297 gotit = True 298 else: KeyboardInterrupt:
问题原因与解决方案
1. 跨进程模型共享冲突
sentence-transformers的模型加载包含PyTorch权重和相关资源,这些资源无法安全地在多进程间共享。主进程加载模型后,子进程继承的模型实例会因资源竞争或状态不一致导致卡死。
修复方案:在子进程内部初始化模型,而不是主进程。修改process_data函数:
def process_data(chunk): # 子进程内单独加载模型 embedding_model = sentence_transformers.SentenceTransformer('sentence-transformers/all-mpnet-base-v2') results = [] for row in chunk: work_id = row[1] mentioning_work_id = row[3] if work_id in df_text and mentioning_work_id in df_text: title1 = df_text[work_id]['title'] title2 = df_text[mentioning_work_id]['title'] embeddings_title1 = embedding_model.encode(title1, convert_to_numpy=True) embeddings_title2 = embedding_model.encode(title2, convert_to_numpy=True) similarity = np.matmul(embeddings_title1, embeddings_title2.T) results.append([row[0], row[1], row[2], row[3], row[4], similarity]) return results
2. 进程数量过载
直接使用multiprocessing.cpu_count()会启用所有vCPU,但SageMaker实例的vCPU可能包含超线程,而模型计算是CPU密集型,过多进程会导致上下文切换开销剧增,甚至资源耗尽。
修复方案:限制进程数量为物理核心数,或减半:
# 用物理核心数,或根据实例类型调整,比如减半 num_cores = max(1, multiprocessing.cpu_count() // 2)
3. 子进程输出阻塞
子进程中的print语句会占用stdout资源,多进程同时输出可能导致IO阻塞。
修复方案:移除process_data中的print语句,或使用日志模块替代。
4. 核心转储文件处理
cores文件是进程崩溃时的内存快照,解决上述问题后会自动停止生成。若需临时禁用,可在SageMaker终端执行:
ulimit -c 0
修改后的完整代码
import sentence_transformers import multiprocessing from tqdm import tqdm from multiprocessing import Pool import numpy as np data = [[100227, 7382501.0, 'view', 30065006, False, ''], [100227, 7382501.0, 'view', 57072062, True, ''], [100227, 7382501.0, 'view', 66405922, True, ''], [100227, 7382501.0, 'view', 5221475, False, ''], [100227, 7382501.0, 'view', 63283995, True, '']] df_text = dict() df_text[7382501] = {'title': 'The Geography of the Internet Industry, Venture Capital, Dot-coms, and Local Knowledge - MATTHEW A. ZOOK', 'abstract': '23', 'highlight': '12'} df_text[30065006] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[57072062] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[66405922] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[5221475] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} df_text[63283995] = {'title': 'Determination of the Effect of Lipophilicity on the in vitro Permeability and Tissue Reservoir Characteristics of Topically Applied Solutes in Human Skin Layers', 'abstract': '12', 'highlight': '12'} def process_data(chunk): embedding_model = sentence_transformers.SentenceTransformer('sentence-transformers/all-mpnet-base-v2') results = [] for row in chunk: work_id = row[1] mentioning_work_id = row[3] if work_id in df_text and mentioning_work_id in df_text: title1 = df_text[work_id]['title'] title2 = df_text[mentioning_work_id]['title'] embeddings_title1 = embedding_model.encode(title1, convert_to_numpy=True) embeddings_title2 = embedding_model.encode(title2, convert_to_numpy=True) similarity = np.matmul(embeddings_title1, embeddings_title2.T) results.append([row[0], row[1], row[2], row[3], row[4], similarity]) return results # 限制进程数量为CPU核心数的一半 num_cores = max(1, multiprocessing.cpu_count() // 2) chunk_size = len(data) // num_cores # 处理剩余数据,避免遗漏 if len(data) % num_cores != 0: chunks = [data[i:i+chunk_size] for i in range(0, len(data)-len(data)%num_cores, chunk_size)] + [data[-len(data)%num_cores:]] else: chunks = [data[i:i+chunk_size] for i in range(0, len(data), chunk_size)] pool = multiprocessing.Pool(processes=num_cores) results = [] with tqdm(total=len(data)) as pbar: for result_chunk in pool.map(process_data, chunks): pbar.update(len(result_chunk)) results += result_chunk final_result = results
内容的提问来源于stack exchange,提问作者Patthebug
相关产品推荐
相关产品推荐

