Amazon SageMaker多进程代码2CPU正常8CPU卡顿及修改后报错排查
问题排查与解决方案
一、多核实例卡顿问题(2核正常,8/16核卡在----------------1------------------)
卡顿发生在embedding_model.encode()调用后,核心原因是多进程环境下的资源竞争或全局状态冲突,结合SageMaker的Linux默认fork进程模式,具体分析:
模型继承的线程安全冲突
主进程加载的embedding_model会被fork出的子进程直接继承,但多数预训练模型内部包含线程池或全局状态,多核下多个子进程同时调用encode()会触发资源锁竞争,导致进程阻塞死锁。全局变量
df_text的访问竞争
子进程共享主进程的df_text字典,多核下大量进程同时执行work_id in df_text这类查询操作,会触发GIL(全局解释器锁)的频繁竞争,拖慢甚至阻塞进程执行。
解决方案:
- 子进程独立加载模型
避免继承主进程的模型实例,在process_data函数内部延迟加载模型,确保每个子进程拥有独立的模型资源:
def process_data(chunk): # 子进程内部单独加载模型,避免共享资源冲突 from sentence_transformers import SentenceTransformer # 替换为你的模型导入语句 embedding_model = SentenceTransformer('your-model-name') # 替换为你的模型初始化逻辑 results = [] for row in tqdm(chunk): # 原有业务逻辑保持不变 work_id = row[1] mentioning_work_id = row[3] if work_id in df_text and mentioning_work_id in df_text: # ... 剩余代码 return results
- 使用spawn模式创建进程池
spawn模式会重新启动Python解释器,完全隔离子进程与主进程的全局状态,彻底解决继承带来的资源冲突:
# 替换原进程池创建代码 import multiprocessing ctx = multiprocessing.get_context('spawn') pool = ctx.Pool(processes=num_cores)
- 优化
df_text的查询效率
将df_text转换为pandas DataFrame并设置索引,提升多进程下的查询速度,减少锁竞争:
# 主进程预处理df_text import pandas as pd df_text = pd.DataFrame(df_text).set_index('work_id') # 假设原df_text是字典结构 # 子进程内查询逻辑修改为 if work_id in df_text.index and mentioning_work_id in df_text.index: title1 = df_text.loc[work_id, 'title'] title2 = df_text.loc[mentioning_work_id, 'title']
二、修改代码后的TypeError问题
报错根源是函数参数不匹配:
- 原
process_data函数设计为接收一组数据行(chunk),内部通过for row in chunk遍历单个数据; - 修改后使用
pool.map(process_data, data, chunksize=2),此时data的每个元素会直接传给process_data,即函数接收的是单个数据行,但函数内部仍执行for row in tqdm(chunk),相当于遍历单个数据行的元素(比如行是列表时,会遍历出int类型元素),导致row[0]触发下标访问错误。
解决方案:
如果要使用pool.map的chunksize参数,需修改process_data为处理单个数据行的逻辑:
def process_data(row): # 移除外层遍历chunk的循环,直接处理单个row 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) return [row[0],row[1],row[2],row[3],row[4],similarity] else: return None # 调用时的逻辑调整 results = [] with tqdm(total=len(data)) as pbar: for result in pool.map(process_data, data, chunksize=100): # 按需调整chunksize大小 pbar.update() if result: # 过滤else分支的空结果 results.append(result)
若想保留原函数处理chunk的逻辑,继续使用手动拆分chunk的方式即可,无需依赖pool.map的chunksize参数。
内容的提问来源于stack exchange,提问作者Patthebug
相关产品推荐
相关产品推荐

