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):pProcess(targetworker,args(fworker-{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(processes4)aspool:resultspool.map(cpu_heavy,[1000000,2000000,3000000,4000000])print(results)进程池自动把任务分发到多个 workerwith保证结束后正确回收。apply_async 异步提交withPool()aspool:futures[pool.apply_async(cpu_heavy,(x,))forxintasks]forfinfutures:print(f.get())apply_async立即返回不阻塞主流程调用get()时再取结果可带超时见“七、退出与超时”。三、进程间通信QueuefrommultiprocessingimportProcess,Queuedefproducer(q):foriinrange(10):q.put(fitem-{i})q.put(None)# 结束标记defconsumer(q):whileTrue:itemq.get()ifitemisNone:breakprint(fconsumed{item})if__name____main__:qQueue()pProcess(targetproducer,args(q,))cProcess(targetconsumer,args(q,))p.start();c.start()p.join();c.join()生产-消费模型。注意有几个消费者就要放几个结束标记或改用q.task_done()q.join()等待消费完成。PipefrommultiprocessingimportProcess,Pipedefworker(conn):conn.send(hello)msgconn.recv()print(fgot{msg})if__name____main__:parent_conn,child_connPipe()pProcess(targetworker,args(child_conn,))p.start()print(parent_conn.recv())parent_conn.send(world)p.join()点对点通信。Pipe()默认双向duplexTrue两端都可收发只需要单向时传Pipe(duplexFalse)。共享内存变量Value / ArrayfrommultiprocessingimportProcess,Valuedefincrement(counter):for_inrange(1000):withcounter.get_lock():counter.value1if__name____main__:counterValue(i,0)# 默认带锁procs[Process(targetincrement,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_recordscollect_raw()# 大批量原始报文# 拆成 chunkchunks[raw_records[i:i1000]foriinrange(0,len(raw_records),1000)]withPool()aspool:resultspool.map(parse_modbus_batch,chunks)parsed[itemforchunkinresultsforiteminchunk]把大批量报文分块后并行解析最后按原顺序拼回。场景 2加解密并行fromcryptography.fernetimportFernetdefencrypt_batch(args):key,data_listargs cipherFernet(key)return[cipher.encrypt(d)fordindata_list]if__name____main__:keyFernet.generate_key()data_chunks[...]# 待加密数据分块withPool()aspool:encrypted_batchespool.map(encrypt_batch,[(key,c)forcindata_chunks])重点密钥必须显式传进 worker。如果把key Fernet.generate_key()放在模块顶层spawn 模式下每个子进程会重新生成一把随机密钥各块密文无法用同一把钥匙解密。场景 3数据聚合并行Map-Reducedefaggregate_chunk(records):stats{}forrinrecords:dr[device_id]sstats.setdefault(d,{sum:0.0,count:0})s[sum]r[value]s[count]1returnstatswithPool()aspool:partialspool.map(aggregate_chunk,chunks)# 合并先汇总求和与计数最后再统一求平均final{}forpartialinpartials:fordevice,sinpartial.items():curfinal.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,targettaskifcmdconnect:returnconnect_device(target)elifcmdpoll:returnpoll_device(target)elifcmdupgrade:returnupgrade_device(target)withPool()aspool:resultspool.map(run_task,tasks)互不依赖的批量操作直接并行若任务间需要按顺序或条件推进用 asyncio 编排更合适。场景 5numpy 数值计算并行importnumpyasnpdefcompute_chunk(arr):returnnp.fft.fft(arr).real chunksnp.array_split(big_array,4)withPool()aspool:resultspool.map(compute_chunk,list(chunks))直接把 numpy 数组作为 chunk 传入即可不要先.tolist()再传——白多一次序列化开销。注意numpy 底层会释放 GIL纯 numpy 的大计算可以先用多线程压测数据量小、计算快时多进程的序列化和进程启动开销反而更慢。五、ProcessPoolExecutor现代推荐基本用法fromconcurrent.futuresimportProcessPoolExecutorwithProcessPoolExecutor(max_workers4)asexecutor:resultsexecutor.map(cpu_heavy,tasks)forrinresults:print(r)executor.map按任务顺序返回结果。配合 as_completedfromconcurrent.futuresimportProcessPoolExecutor,as_completedwithProcessPoolExecutor()asexecutor:futures{executor.submit(cpu_heavy,x):xforxintasks}forfutureinas_completed(futures):xfutures[future]try:resultfuture.result()print(f{x}-{result})exceptExceptionase:print(f{x}failed:{e})谁先完成先处理谁天然支持单任务级错误隔离接口也更现代。六、错误处理整批失败try:withPool()aspool:resultspool.map(risky_func,tasks)exceptExceptionase:print(ftask failed:{e})pool.map遇到第一个异常就会抛出适合“要么全部成功、要么整体重试”的批次任务。单任务级隔离fromconcurrent.futuresimportProcessPoolExecutor,as_completedwithProcessPoolExecutor()asexecutor:futures{executor.submit(risky_func,x):xforxintasks}forfutureinas_completed(futures):try:resultfuture.result()exceptExceptionase:print(f{futures[future]}failed:{e})单个任务失败不影响其他任务继续执行。注意异常对象需要能被 pickle才能从子进程传回主进程。七、退出与超时超时importmultiprocessingasmpwithPool()aspool:resultpool.apply_async(cpu_heavy,(x,))try:valueresult.get(timeout10)exceptmp.TimeoutError:print(timeout)重要get(timeout...)只是放弃等待任务仍会在 worker 里继续跑。如果超时后必须强制终止需要pool.terminate()或给 Pool 设置maxtasksperchild定期回收 worker。优雅退出importsignalimportsysfrommultiprocessingimportPool poolNone# 全局引用供信号处理器访问defhandler(sig,frame):print(收到中断正在终止进程池...)ifpoolisnotNone:pool.terminate()pool.join()sys.exit(0)signal.signal(signal.SIGINT,handler)if__name____main__:withPool()asp:poolp resultsp.map(cpu_heavy,tasks)也可以在with Pool()内捕获KeyboardInterrupt但要注意 worker 可能仍在执行必须显式terminate()才能真正停掉。八、内存共享 shared_memoryPython 3.8frommultiprocessingimportshared_memoryimportnumpyasnp# 主进程创建共享内存arrnp.array([1,2,3,4,5],dtypenp.int64)shmshared_memory.SharedMemory(createTrue,sizearr.nbytes)shared_arrnp.ndarray(arr.shape,dtypearr.dtype,buffershm.buf)shared_arr[:]arr[:]# 子进程按名字挂载同一块内存defworker(name,shape,dtype):existing_shmshared_memory.SharedMemory(namename)arrnp.ndarray(shape,dtypedtype,bufferexisting_shm.buf)print(arr.sum())existing_shm.close()# 子进程只 close不 unlinkpProcess(targetworker,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抛异常整体重试单任务失败submitas_completed隔离长期运行记录失败任务并补偿重跑实践 5监控观察进程数、队列积压、任务耗时监控 worker 异常与内存占用边缘环境内存小定期检查僵尸进程未 join 的子进程十、几个常见的坑坑 1数据可序列化跨进程传参和返回值都要能被 pickleWindows / spawn 模式下 worker 必须是模块顶层函数lambda、闭包不行锁、socket、文件句柄等对象不能直接传递应对传简单数据类型复杂对象在子进程内部重新构造。坑 2fork 与资源继承Linux 默认 fork子进程会继承父进程全部资源数据库连接、锁、日志句柄等容易出问题macOSPython 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 构建混合并行架构