使用multiprocessing的Python代码Windows正常,Ubuntu运行异常求助
问题原因分析
在Ubuntu等Unix-like系统中,multiprocessing默认使用fork机制创建子进程,而Windows采用spawn机制。你在全局作用域提前加载了Whisper模型和Processor,fork后的子进程会直接继承这些内存中的对象,但PyTorch模型依赖的底层资源(如CUDA上下文、内部线程池)无法在fork后的进程间安全共享,导致模型执行generate方法时出现死锁或停滞。
而Windows的spawn机制会重新执行整个脚本,子进程会重新初始化模型,因此不会出现这个问题。
解决方案
针对这个问题,需要对代码做以下几处修改:
- 将模型初始化移到子进程内部:每个子进程独立加载模型,避免继承父进程的模型对象导致资源冲突。
- 优化队列读取逻辑:去掉
queue.empty()的轮询,直接使用queue.get()的阻塞特性,更高效且避免空轮询的问题。 - 添加进程退出信号:当生产者进程完成任务后,向队列中放入与子进程数量相等的结束标记,让子进程能正常退出,避免无限循环。
修改后的代码
import multiprocessing as mc import os import time import librosa from transformers import WhisperProcessor, WhisperForConditionalGeneration def speech_to_text(path_to_audio): # 每个子进程独立加载模型和processor processor = WhisperProcessor.from_pretrained("openai/whisper-tiny") model = WhisperForConditionalGeneration.from_pretrained("openai/whisper-tiny") model.config.forced_decoder_ids = None sound, sr = librosa.load(path_to_audio, sr=16000) input_features = processor(sound, sampling_rate=sr, return_tensors="pt").input_features predicted_ids = model.generate(input_features) transcription = processor.batch_decode(predicted_ids, skip_special_tokens=True) return transcription class Test: def __init__(self): self.queue = mc.Queue() def put_items(self, num_process): for i in range(100): self.queue.put((i, 'speech.wav')) time.sleep(2) # 放入结束标记,每个子进程一个 for _ in range(num_process): self.queue.put(None) def process(self): while True: # 阻塞等待队列元素,无需轮询 item = self.queue.get() if item is None: # 收到结束标记,退出循环 break i, audio_path = item print("Take", i, os.getpid()) result = speech_to_text(audio_path)[0] print(i, result, os.getpid()) def run(self, num_process): put = mc.Process(target=self.put_items, args=(num_process,)) put.start() processes = [] for _ in range(num_process): p = mc.Process(target=self.process) processes.append(p) p.start() put.join() # 等待所有子进程完成 for p in processes: p.join() if __name__ == '__main__': t = Test() t.run(num_process=2)
额外优化建议
如果模型加载耗时较长,可以让每个子进程只加载一次模型,而不是每次调用speech_to_text都加载。可以修改为在子进程启动时加载模型,然后循环处理队列中的任务:
def process_worker(queue): # 子进程启动时加载一次模型 processor = WhisperProcessor.from_pretrained("openai/whisper-tiny") model = WhisperForConditionalGeneration.from_pretrained("openai/whisper-tiny") model.config.forced_decoder_ids = None while True: item = queue.get() if item is None: break i, audio_path = item print("Take", i, os.getpid()) sound, sr = librosa.load(audio_path, sr=16000) input_features = processor(sound, sampling_rate=sr, return_tensors="pt").input_features predicted_ids = model.generate(input_features) transcription = processor.batch_decode(predicted_ids, skip_special_tokens=True)[0] print(i, transcription, os.getpid()) # 修改Test类的process方法为调用这个worker class Test: # ... 其他代码不变 ... def process(self): process_worker(self.queue)
这样可以避免重复加载模型,提升处理效率。
内容的提问来源于stack exchange,提问作者Gawain
相关产品推荐
相关产品推荐

