资讯动态

Python并发编程实战:多线程与多进程应用指南

发布时间:2026/9/2 5:14:31 来源:尧图企业网站定制
在实际 Python 项目中当任务量增大或需要同时处理多个 I/O 操作时单线程顺序执行的模式很快就会成为性能瓶颈。无论是开发一个需要同时处理多个用户请求的 Web 服务还是编写一个需要批量下载文件或处理数据的脚本并发编程都是绕不开的核心技能。对于 Python 开发者而言理解并正确使用多线程、多进程以及它们之间的同步与通信机制是从编写简单脚本到构建健壮应用的关键一步。本文旨在提供一个从零基础到实战的 Python 并发编程全景指南。我们将从并发的基本概念入手逐步深入到threading和multiprocessing模块的具体使用详细探讨线程同步、进程通信等核心问题并分析ThreadLocal这类高级工具的应用场景。最终你将能够根据不同的任务类型I/O 密集型 vs CPU 密集型选择合适的并发模型并规避常见的陷阱如全局解释器锁GIL的影响、死锁和数据竞争。1. 理解 Python 并发编程的核心概念与 GIL在动手写代码之前必须先理清几个关键概念否则很容易陷入“为什么我用了多线程速度反而更慢了”的困惑。1.1 并发 vs 并行目标与实现并发和并行是两个经常被混淆的术语。并发指的是系统具有处理多个任务的能力这些任务在宏观上看起来是同时进行的但在微观上可能是交替执行的。例如一个单核 CPU 通过时间片轮转快速切换执行多个线程给用户一种“同时”的错觉。并行则指的是系统真正在同一时刻同时执行多个任务这通常需要多核 CPU 的支持。Python 的并发编程主要围绕两个模块threading多线程和multiprocessing多进程。选择哪一个很大程度上取决于任务性质和对全局解释器锁GIL的考量。1.2 全局解释器锁GILPython 多线程的“守门人”GIL 是 Python 解释器特指 CPython中的一个互斥锁它规定任何时候都只有一个线程可以执行 Python 字节码。这意味着即使在多核 CPU 上一个 Python 进程中的多个线程也无法实现真正的并行计算。这引出了一个至关重要的结论对于 CPU 密集型任务如科学计算、图像处理使用多线程通常无法提升性能甚至可能因为线程切换的开销而变慢。此时应使用multiprocessing创建多个进程每个进程有独立的 Python 解释器和内存空间从而绕过 GIL 实现真正的并行。相反对于I/O 密集型任务如网络请求、文件读写、数据库查询线程大部分时间在等待 I/O 操作完成处于阻塞状态此时 GIL 会被释放其他线程可以运行。因此多线程能有效提升 I/O 密集型任务的吞吐量让程序在等待一个 I/O 时去处理另一个 I/O。任务类型特点推荐并发模型原因I/O 密集型大量时间在等待外部响应CPU 空闲。多线程 (threading)线程在 I/O 阻塞时释放 GIL其他线程可运行提升整体效率。CPU 密集型大量时间在进行计算CPU 持续工作。多进程 (multiprocessing)多进程可绕过 GIL利用多核实现真正并行计算。混合型既有计算又有 I/O。多进程 进程内多线程或concurrent.futures根据瓶颈灵活选择或使用高级抽象模块。2. 环境准备与项目结构在开始编码前确保你的 Python 环境已就绪。本文示例基于 Python 3.8因为concurrent.futures等模块在此版本后更为成熟稳定。2.1 检查 Python 环境打开终端或命令行运行以下命令确认版本python --version # 或 python3 --version输出应类似Python 3.8.10。如果版本过低建议从 Python 官网 下载安装最新稳定版。安装时务必勾选“Add Python to PATH”。2.2 创建项目目录与文件建议为本次学习创建一个独立的项目目录结构清晰便于管理。mkdir python_concurrency_tutorial cd python_concurrency_tutorial在该目录下我们将创建多个 Python 文件来演示不同概念io_bound_threading.py: I/O 密集型任务的多线程示例。cpu_bound_multiprocessing.py: CPU 密集型任务的多进程示例。thread_synchronization.py: 线程同步机制示例。process_communication.py: 进程通信示例。thread_local_demo.py: ThreadLocal 使用示例。你可以使用任何文本编辑器或 IDE如 VS Code、PyCharm来编写代码。如果使用 VS Code确保安装了 Python 扩展并正确配置了 Python 解释器路径。3. 多线程 (threading) 实战处理 I/O 密集型任务让我们从一个实际的 I/O 密集型场景开始模拟批量下载多个网页。我们将对比单线程和多线程的执行时间直观感受其差异。3.1 单线程版本顺序下载首先我们创建一个模拟网络下载的函数它只是休眠一段时间来模拟网络延迟。# io_bound_threading.py import time import threading def download_site(url, delay0.5): 模拟下载一个网页休眠 delay 秒模拟网络延迟 thread_name threading.current_thread().name print(f[{thread_name}] 开始下载: {url}) time.sleep(delay) # 模拟 I/O 等待 print(f[{thread_name}] 完成下载: {url}) return f{url} 的内容 def download_all_sites_single_thread(sites): 单线程顺序下载所有站点 start_time time.time() results [] for site in sites: results.append(download_site(site)) duration time.time() - start_time print(f\n单线程下载 {len(sites)} 个站点耗时: {duration:.2f} 秒) return results if __name__ __main__: # 模拟 5 个需要下载的 URL sites [fhttps://site{i}.com for i in range(5)] download_all_sites_single_thread(sites)运行这个脚本你会看到每个站点依次开始和完成下载总耗时大约是站点数 * delay即 5 * 0.5 ≈ 2.5 秒。3.2 多线程版本并发下载现在我们使用threading.Thread来创建多个线程同时执行下载任务。# 接上 io_bound_threading.py 文件 def download_all_sites_multi_thread(sites): 使用多线程并发下载所有站点 start_time time.time() threads [] results [] # 注意这个列表不是线程安全的 # 定义一个线程要执行的函数 def download_and_store(url): result download_site(url) results.append(result) # 这里存在数据竞争风险 # 创建并启动线程 for site in sites: thread threading.Thread(targetdownload_and_store, args(site,)) threads.append(thread) thread.start() # 等待所有线程完成 for thread in threads: thread.join() duration time.time() - start_time print(f\n多线程下载 {len(sites)} 个站点耗时: {duration:.2f} 秒) return results if __name__ __main__: sites [fhttps://site{i}.com for i in range(5)] print( 单线程执行 ) download_all_sites_single_thread(sites) print(\n 多线程执行 ) download_all_sites_multi_thread(sites)运行修改后的脚本你会看到多个线程几乎同时开始下载总耗时接近单个任务的延迟时间约 0.5 秒而不是 2.5 秒。这清晰地展示了多线程在处理 I/O 等待时的优势。注意上面的results.append(result)操作在线程间共享一个可变对象列表这在多线程环境下是不安全的可能导致数据丢失或损坏。这引出了下一个核心主题——线程同步。4. 线程同步保护共享资源当多个线程需要读写同一个变量、文件或数据结构时就会发生数据竞争。为了避免结果不可预测必须使用同步原语来协调线程间的访问。4.1 使用互斥锁 (threading.Lock)互斥锁是最基本的同步工具它保证同一时刻只有一个线程能进入被保护的代码块临界区。# thread_synchronization.py import threading import time # 一个共享的计数器 counter 0 # 创建一个锁对象 counter_lock threading.Lock() def unsafe_increment(): 不安全的递增操作会导致数据竞争 global counter for _ in range(100000): counter 1 # 这个操作不是原子的 def safe_increment_with_lock(): 使用锁保护的安全递增操作 global counter for _ in range(100000): with counter_lock: # 获取锁退出 with 块时自动释放 counter 1 def test_counter(use_lockFalse): 测试函数启动多个线程修改计数器 global counter counter 0 # 重置计数器 threads [] func safe_increment_with_lock if use_lock else unsafe_increment # 创建 10 个线程 for _ in range(10): t threading.Thread(targetfunc) threads.append(t) t.start() # 等待所有线程结束 for t in threads: t.join() expected 10 * 100000 print(f使用锁: {use_lock}, 最终计数: {counter}, 期望值: {expected}, 正确: {counter expected}) if __name__ __main__: print(测试不安全递增可能出错:) for i in range(3): # 多跑几次看错误结果 test_counter(use_lockFalse) print(\n测试使用锁的安全递增:) for i in range(3): test_counter(use_lockTrue)运行此脚本你会看到在不使用锁的情况下最终计数值几乎总是小于期望值 1,000,000这是因为counter 1这个操作读取、计算、写入不是原子的可能被其他线程打断。而使用锁后结果始终正确。4.2 使用队列 (queue.Queue) 进行线程间通信除了保护共享变量线程间更安全、更常用的通信方式是使用queue.Queue。它是一个线程安全的 FIFO先进先出队列完美契合“生产者-消费者”模型。# 接上 thread_synchronization.py 文件 import queue import random def producer(q, item_count): 生产者线程生成数据放入队列 for i in range(item_count): item f产品-{i} time.sleep(random.uniform(0.01, 0.1)) # 模拟生产时间 q.put(item) print(f[生产者] 生产了: {item}) # 放入结束信号 q.put(None) def consumer(q, consumer_id): 消费者线程从队列取出并处理数据 while True: item q.get() # 阻塞直到有数据可取 if item is None: q.put(None) # 将结束信号放回让其他消费者也能结束 print(f[消费者{consumer_id}] 收到结束信号退出。) break time.sleep(random.uniform(0.05, 0.15)) # 模拟处理时间 print(f[消费者{consumer_id}] 处理了: {item}) q.task_done() # 通知队列该任务已完成 def producer_consumer_demo(): 生产者-消费者模型演示 q queue.Queue(maxsize5) # 设置队列最大容量为5 # 创建生产者 prod_thread threading.Thread(targetproducer, args(q, 10)) # 创建两个消费者 cons_threads [threading.Thread(targetconsumer, args(q, i)) for i in range(2)] prod_thread.start() for t in cons_threads: t.start() prod_thread.join() for t in cons_threads: t.join() print(生产-消费任务完成。) if __name__ __main__: print(\n--- 测试队列通信 ---) producer_consumer_demo()Queue的put()和get()方法是线程安全的并且支持阻塞操作当队列空时get()会等待当队列满时put()会等待。task_done()和join()方法可以用于等待队列中所有任务处理完毕这在批量任务处理中非常有用。5. 多进程 (multiprocessing) 实战加速 CPU 密集型任务对于计算密集型任务我们需要使用multiprocessing模块来创建多个进程利用多核 CPU。5.1 基础多进程计算素数个数我们以一个计算某个范围内素数个数的任务为例。# cpu_bound_multiprocessing.py import math import time from multiprocessing import Process, Queue def is_prime(n): 判断一个数是否为素数 if n 2: return False for i in range(2, int(math.sqrt(n)) 1): if n % i 0: return False return True def count_primes_in_range(start, end, result_queue): 计算从 start 到 end不含范围内的素数个数并将结果放入队列 count 0 for num in range(start, end): if is_prime(num): count 1 result_queue.put(count) def single_process_count(limit): 单进程计算 start_time time.time() count 0 for num in range(limit): if is_prime(num): count 1 duration time.time() - start_time print(f单进程计算 0-{limit} 的素数个数: {count}, 耗时: {duration:.2f} 秒) return count, duration def multi_process_count(limit, num_processes4): 多进程计算将任务平均分给多个进程 start_time time.time() chunk_size limit // num_processes processes [] result_queue Queue() # 创建并启动进程 for i in range(num_processes): start i * chunk_size # 最后一个进程处理剩余部分 end limit if i num_processes - 1 else (i 1) * chunk_size p Process(targetcount_primes_in_range, args(start, end, result_queue)) processes.append(p) p.start() # 等待所有进程结束 for p in processes: p.join() # 从队列中收集结果 total_primes 0 while not result_queue.empty(): total_primes result_queue.get() duration time.time() - start_time print(f{num_processes}进程计算 0-{limit} 的素数个数: {total_primes}, 耗时: {duration:.2f} 秒) return total_primes, duration if __name__ __main__: limit 200000 print(f计算 0 到 {limit} 之间的素数个数) single_count, single_time single_process_count(limit) multi_count, multi_time multi_process_count(limit, num_processes4) if single_count multi_count: print(f\n性能提升: {single_time/multi_time:.2f} 倍) else: print(\n错误单进程与多进程计算结果不一致)运行此脚本注意计算量较大可根据机器性能调整limit值你会看到多进程版本的计算时间显著少于单进程版本提升倍数接近你的 CPU 核心数。这证明了多进程在 CPU 密集型任务上的有效性。5.2 使用进程池 (multiprocessing.Pool)手动管理进程和队列比较繁琐。multiprocessing.Pool提供了一个更高级的接口它管理一个工作进程池并提供了像map这样的便捷方法。# 接上 cpu_bound_multiprocessing.py 文件 from multiprocessing import Pool def count_primes_chunk(args): 供进程池使用的函数参数打包为一个元组 start, end args count 0 for num in range(start, end): if is_prime(num): count 1 return count def multi_process_with_pool(limit, num_processes4): 使用进程池进行计算 start_time time.time() chunk_size limit // num_processes # 准备参数列表[(0, chunk_size), (chunk_size, 2*chunk_size), ...] chunks [(i * chunk_size, limit if i num_processes - 1 else (i 1) * chunk_size) for i in range(num_processes)] with Pool(processesnum_processes) as pool: # 使用 map 方法将任务分发给进程池 results pool.map(count_primes_chunk, chunks) total_primes sum(results) duration time.time() - start_time print(f进程池({num_processes}进程)计算 0-{limit} 的素数个数: {total_primes}, 耗时: {duration:.2f} 秒) return total_primes, duration if __name__ __main__: limit 200000 print(f\n--- 使用进程池计算 ---) pool_count, pool_time multi_process_with_pool(limit, num_processes4)Pool.map方法将可迭代对象如我们的参数列表中的每个元素应用到函数上并自动将工作分配给池中的进程最后收集所有结果。代码比手动管理进程和队列简洁得多。6. 进程间通信 (IPC)进程拥有独立的内存空间不能像线程那样直接共享变量。multiprocessing模块提供了多种进程间通信机制如Queue,Pipe,Value,Array以及Manager。6.1 使用multiprocessing.Queuemultiprocessing.Queue是一个跨进程的队列其接口与queue.Queue类似但底层实现了进程间的数据传递。# process_communication.py from multiprocessing import Process, Queue import time def producer_process(queue, items): 生产者进程 for item in items: print(f[生产者进程] 发送: {item}) queue.put(item) time.sleep(0.1) queue.put(None) # 发送结束信号 def consumer_process(queue, consumer_id): 消费者进程 while True: item queue.get() if item is None: print(f[消费者进程{consumer_id}] 收到结束信号) break print(f[消费者进程{consumer_id}] 收到: {item}) time.sleep(0.2) if __name__ __main__: mp_queue Queue() items_to_send [数据A, 数据B, 数据C, 数据D] prod Process(targetproducer_process, args(mp_queue, items_to_send)) cons Process(targetconsumer_process, args(mp_queue, 1)) prod.start() cons.start() prod.join() cons.join() print(进程间队列通信示例结束。)6.2 使用multiprocessing.Manager共享状态Manager可以创建一个服务进程该进程管理共享对象如列表、字典其他进程通过代理来访问这些对象。注意这种方式比进程内共享内存慢但可以共享更复杂的数据结构。# 接上 process_communication.py 文件 from multiprocessing import Process, Manager def worker_with_manager(shared_list, process_id): 工作进程向共享列表添加数据 import os pid os.getpid() shared_list.append(f来自进程 {process_id} (PID: {pid}) 的数据) print(f进程 {process_id} 添加了数据。当前列表: {shared_list}) if __name__ __main__: print(\n--- 使用 Manager 共享列表 ---) with Manager() as manager: shared_list manager.list() # 创建一个由 Manager 管理的共享列表 processes [] for i in range(3): p Process(targetworker_with_manager, args(shared_list, i)) processes.append(p) p.start() for p in processes: p.join() print(f最终共享列表内容: {shared_list})重要提示Manager对象支持的类型如list,dict,Namespace在修改时是同步的但性能开销较大。对于高性能需求应优先考虑Value或Array进行简单数据共享或使用Queue进行消息传递。7. ThreadLocal线程的私有存储空间有时你需要一些数据对每个线程都是“全局”的但又不希望在不同线程间共享。典型的应用场景是数据库连接、Web 请求上下文等。threading.local()可以创建线程本地存储对象每个线程对它属性的修改其他线程都看不到。# thread_local_demo.py import threading import time import random # 创建一个 ThreadLocal 实例 local_data threading.local() def get_connection(): 模拟获取一个‘线程专属’的数据库连接这里用字符串代替 # 每个线程第一次访问 local_data.connection 时它都不存在 if not hasattr(local_data, connection): thread_name threading.current_thread().name # 模拟创建连接的开销 time.sleep(random.uniform(0.05, 0.1)) local_data.connection f{thread_name} 的数据库连接 print(f[{thread_name}] 创建了新的连接: {local_data.connection}) return local_data.connection def database_operation(operation_id): 模拟数据库操作需要用到连接 conn get_connection() # 获取本线程的连接 thread_name threading.current_thread().name time.sleep(0.05) # 模拟操作耗时 print(f[{thread_name}] 使用连接 {conn} 执行操作 {operation_id}) def worker(worker_id): 工作线程函数执行多次数据库操作 for i in range(3): database_operation(i) if __name__ __main__: threads [] for i in range(3): t threading.Thread(targetworker, args(i,), namefWorker-{i}) threads.append(t) t.start() for t in threads: t.join() print(所有线程操作完成。)运行此脚本你会看到每个线程Worker-0, Worker-1, Worker-2都创建了自己唯一的连接并在后续操作中重复使用它而不会与其他线程冲突。如果不使用ThreadLocal就需要在函数间传递连接参数或者使用复杂的锁机制来管理一个共享连接池代码会变得冗长且容易出错。8. 常见问题排查与最佳实践并发编程容易出错下面是一些典型问题及其解决方案。8.1 常见问题排查表问题现象可能原因检查与解决思路多线程程序速度没有提升甚至变慢1. 任务是 CPU 密集型受 GIL 限制。2. 线程创建/切换开销大于收益。3. 存在严重的锁竞争锁粒度太粗。1. 使用multiprocessing替代threading。2. 使用线程池 (ThreadPoolExecutor) 复用线程。3. 细化锁的粒度或使用无锁数据结构如queue.Queue。程序偶尔出现数据错误或丢失数据竞争。多个线程同时读写共享变量未加锁。1. 使用threading.Lock或RLock保护临界区。2. 使用线程安全的数据结构如queue.Queue。3. 将共享状态改为通过队列传递消息。程序卡死不再响应死锁两个或多个线程互相等待对方释放锁。1. 确保锁的获取顺序在所有线程中一致。2. 使用with lock:上下文管理器避免忘记释放锁。3. 使用threading.Timer或设置锁的超时参数。多进程程序启动报错或行为异常1. 在 Windows 上未将代码放在if __name__ __main__:下。2. 传递了不可序列化pickle的对象给子进程。1.必须将进程启动代码放在if __name__ __main__:块内。2. 确保传递给Process或Pool的参数可以被pickle模块序列化。进程池 (Pool) 中的任务没有执行可能是在交互式环境如 Jupyter中使用不当或者主进程提前退出。1. 确保使用with Pool() as pool:或手动调用pool.close()和pool.join()。2. 在脚本中运行而非某些受限的交互环境。8.2 最佳实践清单明确任务类型I/O 密集型用多线程CPU 密集型用多进程。混合型任务考虑组合或使用concurrent.futures。优先使用高级抽象在大多数场景下使用concurrent.futures.ThreadPoolExecutor或ProcessPoolExecutor比手动管理线程/进程更安全、更简洁。避免共享状态线程间尽量通过队列 (queue.Queue) 传递消息而非共享变量。进程间使用multiprocessing.Queue,Pipe或Manager。锁的使用要谨慎只在必要时加锁并尽量缩短持有锁的时间。使用with lock:上下文管理器确保锁总能被释放。避免在持有一个锁时再去获取另一个锁容易导致死锁。设置守护线程/进程如果某些线程或进程是后台服务不需要等待其结束可以将其daemon属性设为True。这样主程序退出时它们会自动被终止。处理异常子线程或子进程中的异常默认不会传递到主线程/主进程。确保在任务函数内部做好异常捕获和日志记录。对于concurrent.futures可以使用future.result()来获取异常。控制并发度不要无限制地创建线程或进程。使用池Pool,Executor来限制最大数量防止耗尽系统资源。进程编程的特殊要求在 Windows 上创建多进程的代码必须放在if __name__ __main__:块中否则会引发递归创建进程的错误。8.3 下一步学习与扩展方向掌握了threading和multiprocessing的基础后你可以进一步探索以下方向来构建更健壮、高效的并发应用concurrent.futures模块这是 Python 标准库中更高层次的异步执行接口。ThreadPoolExecutor和ProcessPoolExecutor提供了更简单的线程池和进程池管理并且支持Future对象便于获取任务结果和异常。异步编程 (asyncio)对于高并发 I/O 操作如大量网络连接asyncio提供的协程模型比多线程更轻量级、效率更高。它使用async/await语法是编写现代 Python 网络服务的重要工具。第三方库joblib: 特别适用于科学计算场景的轻量级流水线并行。celery: 分布式任务队列用于处理后台作业支持多种消息中间件。dask: 用于并行计算的灵活库可以处理大于内存的数据集。性能分析与调试学习使用cProfile、line_profiler等工具分析程序性能瓶颈。使用threading模块的enumerate()函数或faulthandler模块来调试死锁。并发编程是 Python 进阶的必经之路理解其核心概念并谨慎实践能让你开发的程序在处理复杂任务时游刃有余。从今天起在遇到 I/O 等待或大量计算时不妨先思考一下这个任务是否可以用并发来优化

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价