Python 的 GIL 决定了普通多线程无法并行执行纯 Python 的 CPU 密集计算。工业边缘里协议解析、加解密、数据聚合这类场景需要真正的多核并行,multiprocessing(多进程)是主要方案。本文用可运行的示例讲清楚怎么用、怎么通信、怎么避坑。
一、什么时候用 multiprocessing
适用场景
- CPU 密集计算:协议解析、加解密、数值计算等
- 相互独立的任务:彼此不依赖,可以同时跑
- 需要利用多核 CPU
- 长耗时任务,且能接受进程级隔离与通信开销
与 asyncio 对比
- asyncio:适合 IO 密集(网络请求、文件读写、等待类操作),单线程事件循环,并发而不并行
- multiprocessing:适合 CPU 密集,真正多核并行
- 两者可以配合:asyncio 负责并发 IO 与调度,把重计算丢给进程池,通过队列 / future 回收结果
与多线程对比
- 多线程:受 GIL 限制,纯 Python 的 CPU 计算无法并行;但 IO 密集场景线程依然好用,且线程间共享内存方便
- 多进程:每个进程有独立解释器,真并行;代价是进程间不共享内存,通信与序列化有额外开销
- 补充:numpy 等底层 C 库在计算时会释放 GIL,纯 numpy 的大计算有时用多线程也够,不必无脑上多进程
二、基础用法
最简单的进程
frommultiprocessingimportProcessdefworker(name):print(f"{name}processing")if__name__=="__main__":procs=[]foriinrange(4):p=Process(target=worker,args=(f"worker-{i}",))procs.append(p)p.start()forpinprocs:p.join()4 个进程并行执行。if __name__ == "__main__"保护必不可少,见“常见坑 5”。
Pool 进程池
frommultiprocessingimportPooldefcpu_heavy(x):# CPU 密集计算returnsum(i*iforiinrange(x))if__name__=="__main__":withPool(processes=4)aspool:results=pool.map(cpu_heavy,[1000000,2000000,3000000,4000000])print(results)进程池自动把任务分发到多个 worker,with保证结束后正确回收。
apply_async 异步提交
withPool()aspool:futures=[pool.apply_async(cpu_heavy,(x,))forxintasks]forfinfutures:print(f.get())apply_async立即返回,不阻塞主流程;调用get()时再取结果(可带超时,见“七、退出与超时”)。
三、进程间通信
Queue
frommultiprocessingimportProcess,Queuedefproducer(q):foriinrange(10):q.put(f"item-{i}")q.put(None)# 结束标记defconsumer(q):whileTrue:item=q.get()ifitemisNone:breakprint(f"consumed{item}")if__name__=="__main__":q=Queue()p=Process(target=producer,args=(q,))c=Process(target=consumer,args=(q,))p.start();c.start()p.join();c.join()生产-消费模型。注意:有几个消费者就要放几个结束标记(或改用q.task_done()+q.join()等待消费完成)。
Pipe
frommultiprocessingimportProcess,Pipedefworker(conn):conn.send("hello")msg=conn.recv()print(f"got{msg}")if__name__=="__main__":parent_conn,child_conn=Pipe()p=Process(target=worker,args=(child_conn,))p.start()print(parent_conn.recv())parent_conn.send("world")p.join()点对点通信。Pipe()默认双向(duplex=True),两端都可收发;只需要单向时传Pipe(duplex=False)。
共享内存变量(Value / Array)
frommultiprocessingimportProcess,Valuedefincrement(counter):for_inrange(1000):withcounter.get_lock():counter.value+=1if__name__=="__main__":counter=Value("i",0)# 默认带锁procs=[Process(target=increment,args=(counter,))for_inrange(4)]forpinprocs:p.start()forpinprocs:p.join()print(counter.value)适合传小体积的共享变量(计数、标志位、配置项);加锁保证并发安全。
四、工业边缘实战场景
场景 1:协议解析并行
defparse_modbus_batch(records):return[parse_modbus(r)forrinrecords]# parse_modbus 为实际解析函数if__name__=="__main__":raw_records=collect_raw()# 大批量原始报文# 拆成 chunkchunks=[raw_records[i:i+1000]foriinrange(0,len(raw_records),1000)]withPool()aspool:results=pool.map(parse_modbus_batch,chunks)parsed=[itemforchunkinresultsforiteminchunk]把大批量报文分块后并行解析,最后按原顺序拼回。
场景 2:加解密并行
fromcryptography.fernetimportFernetdefencrypt_batch(args):key,data_list=args cipher=Fernet(key)return[cipher.encrypt(d)fordindata_list]if__name__=="__main__":key=Fernet.generate_key()data_chunks=[...]# 待加密数据分块withPool()aspool:encrypted_batches=pool.map(encrypt_batch,[(key,c)forcindata_chunks])重点:密钥必须显式传进 worker。如果把key = Fernet.generate_key()放在模块顶层,spawn 模式下每个子进程会重新生成一把随机密钥,各块密文无法用同一把钥匙解密。
场景 3:数据聚合并行(Map-Reduce)
defaggregate_chunk(records):stats={}forrinrecords:d=r["device_id"]s=stats.setdefault(d,{"sum":0.0,"count":0})s["sum"]+=r["value"]s["count"]+=1returnstatswithPool()aspool:partials=pool.map(aggregate_chunk,chunks)# 合并:先汇总求和与计数,最后再统一求平均final={}forpartialinpartials:fordevice,sinpartial.items():cur=final.setdefault(device,{"sum":0.0,"count":0})cur["sum"]+=s["sum"]cur["count"]+=s["count"]final_avg={d:s["sum"]/s["count"]ford,sinfinal.items()}不要对分块平均值直接再取平均:各块样本数不同时会算错。正确做法是 Map 阶段算“和 + 计数”,Reduce 阶段汇总后统一除以总计数。
场景 4:独立任务并行
tasks=[("connect","device_1"),("poll","device_2"),("upgrade","device_3"),]defrun_task(task):cmd,target=taskifcmd=="connect":returnconnect_device(target)elifcmd=="poll":returnpoll_device(target)elifcmd=="upgrade":returnupgrade_device(target)withPool()aspool:results=pool.map(run_task,tasks)互不依赖的批量操作直接并行;若任务间需要按顺序或条件推进,用 asyncio 编排更合适。
场景 5:numpy 数值计算并行
importnumpyasnpdefcompute_chunk(arr):returnnp.fft.fft(arr).real chunks=np.array_split(big_array,4)withPool()aspool:results=pool.map(compute_chunk,list(chunks))直接把 numpy 数组作为 chunk 传入即可,不要先.tolist()再传——白多一次序列化开销。注意:numpy 底层会释放 GIL,纯 numpy 的大计算可以先用多线程压测;数据量小、计算快时,多进程的序列化和进程启动开销反而更慢。
五、ProcessPoolExecutor(现代推荐)
基本用法
fromconcurrent.futuresimportProcessPoolExecutorwithProcessPoolExecutor(max_workers=4)asexecutor:results=executor.map(cpu_heavy,tasks)forrinresults:print(r)executor.map按任务顺序返回结果。
配合 as_completed
fromconcurrent.futuresimportProcessPoolExecutor,as_completedwithProcessPoolExecutor()asexecutor:futures={executor.submit(cpu_heavy,x):xforxintasks}forfutureinas_completed(futures):x=futures[future]try:result=future.result()print(f"{x}->{result}")exceptExceptionase:print(f"{x}failed:{e}")谁先完成先处理谁,天然支持单任务级错误隔离,接口也更现代。
六、错误处理
整批失败
try:withPool()aspool:results=pool.map(risky_func,tasks)exceptExceptionase:print(f"task failed:{e}")pool.map遇到第一个异常就会抛出,适合“要么全部成功、要么整体重试”的批次任务。
单任务级隔离
fromconcurrent.futuresimportProcessPoolExecutor,as_completedwithProcessPoolExecutor()asexecutor:futures={executor.submit(risky_func,x):xforxintasks}forfutureinas_completed(futures):try:result=future.result()exceptExceptionase:print(f"{futures[future]}failed:{e}")单个任务失败不影响其他任务继续执行。注意:异常对象需要能被 pickle,才能从子进程传回主进程。
七、退出与超时
超时
importmultiprocessingasmpwithPool()aspool:result=pool.apply_async(cpu_heavy,(x,))try:value=result.get(timeout=10)exceptmp.TimeoutError:print("timeout")重要:get(timeout=...)只是放弃等待,任务仍会在 worker 里继续跑。如果超时后必须强制终止,需要pool.terminate(),或给 Pool 设置maxtasksperchild定期回收 worker。
优雅退出
importsignalimportsysfrommultiprocessingimportPool pool=None# 全局引用,供信号处理器访问defhandler(sig,frame):print("收到中断,正在终止进程池...")ifpoolisnotNone:pool.terminate()pool.join()sys.exit(0)signal.signal(signal.SIGINT,handler)if__name__=="__main__":withPool()asp:pool=p results=p.map(cpu_heavy,tasks)也可以在with Pool()内捕获KeyboardInterrupt,但要注意 worker 可能仍在执行,必须显式terminate()才能真正停掉。
八、内存共享 shared_memory(Python 3.8+)
frommultiprocessingimportshared_memoryimportnumpyasnp# 主进程:创建共享内存arr=np.array([1,2,3,4,5],dtype=np.int64)shm=shared_memory.SharedMemory(create=True,size=arr.nbytes)shared_arr=np.ndarray(arr.shape,dtype=arr.dtype,buffer=shm.buf)shared_arr[:]=arr[:]# 子进程:按名字挂载同一块内存defworker(name,shape,dtype):existing_shm=shared_memory.SharedMemory(name=name)arr=np.ndarray(shape,dtype=dtype,buffer=existing_shm.buf)print(arr.sum())existing_shm.close()# 子进程只 close,不 unlinkp=Process(target=worker,args=(shm.name,arr.shape,arr.dtype))p.start();p.join()shm.close()shm.unlink()# 只能由创建方 unlink,否则资源泄漏适合在进程间零拷贝共享大数组。共享内存不会随进程退出自动释放,务必在创建方close()+unlink(),子进程只close()。
九、几个工程实践
实践 1:合理设置进程数
- 参考 CPU 核数(
os.cpu_count()),再留出系统与 IO 的余量 - 结合业务与内存预算:每多一个进程就多一份内存
- 边缘设备核少、资源紧,别盲目开满
实践 2:任务拆分粒度
- 拆得太小:进程调度与序列化开销占比大,反而变慢
- 拆得太大:并行度不足,负载不均衡
- 以“单块任务耗时远大于通信开销”为平衡点
实践 3:数据传输方式
- 数据小:直接传参(自动 pickle)
- 数据大:用 shared_memory 或磁盘映射,避免反复拷贝
- 长任务、频繁传大对象:考虑按批处理 + 队列削峰
实践 4:错误处理
- 整批失败:
pool.map抛异常,整体重试 - 单任务失败:
submit+as_completed隔离 - 长期运行:记录失败任务并补偿重跑
实践 5:监控
- 观察进程数、队列积压、任务耗时
- 监控 worker 异常与内存占用(边缘环境内存小)
- 定期检查僵尸进程(未 join 的子进程)
十、几个常见的坑
坑 1:数据可序列化
- 跨进程传参和返回值都要能被 pickle
- Windows / spawn 模式下 worker 必须是模块顶层函数,lambda、闭包不行
- 锁、socket、文件句柄等对象不能直接传递
应对:传简单数据类型;复杂对象在子进程内部重新构造。
坑 2:fork 与资源继承
- Linux 默认 fork:子进程会继承父进程全部资源,数据库连接、锁、日志句柄等容易出问题
- macOS(Python 3.8+)和 Windows 默认 spawn:子进程重新导入模块、独立初始化
应对:生产环境显式mp.set_start_method("spawn"),连接类资源在子进程内独立创建。
坑 3:内存占用
- 多进程内存按进程数累加,边缘设备内存有限
应对:合理控制进程数 + 监控内存,必要时用 shared_memory 减少复制。
坑 4:僵尸进程
- 子进程退出后不
join(),会积累僵尸进程
应对:用with上下文或显式join()回收。
坑 5:缺 ifname== “main”
- Windows / spawn 下缺少保护会导致无限递归创建进程
应对:所有可执行入口一律加if __name__ == "__main__":。
十一、在运行时(如EdgeOS)中的角色
- CPU 密集场景(协议解析、加解密、数据聚合)用进程池承担
- 与 asyncio 配合:asyncio 管 IO 密集与调度,重计算丢给进程池,结果通过队列 / future 回收
- 进程数、任务粒度按设备资源压测与监测
- 长任务要有超时、取消与优雅退出机制
了解 Zenova EdgeOS 完整方案 →
下一步建议
- 盘点运行时里的 CPU 密集点(协议解析 / 加解密 / 聚合)
- 按场景选型:短任务用 Pool,长任务用 Executor + as_completed,大数组用 shared_memory
- 合理拆分任务,压测进程数与耗时
- 上线监测:进程数、内存、任务耗时、异常
- 持续优化:结合 asyncio 构建混合并行架构