Python并发编程实战:从多线程、多进程到线程同步与进程通信
2026/9/2 9:42:34 网站建设 项目流程

这次我们来看一套完整的 Python 并发编程教程。对于任何想要提升程序性能、处理 I/O 密集型或计算密集型任务的开发者来说,并发编程都是绕不开的核心技能。这套教程从零基础概念讲起,一直深入到多线程、多进程、线程同步、进程通信以及 ThreadLocal 等实战应用,目标是让你不仅能理解原理,更能写出高效、稳定、无 Bug 的并发代码。

本文将带你快速梳理 Python 并发编程的核心知识体系,并通过实战代码演示关键技术的应用。无论你是想优化爬虫速度、加速数据处理,还是构建高并发的网络服务,这里的内容都能提供直接的帮助。我们会重点关注 GIL(全局解释器锁)对多线程的影响、如何选择多线程与多进程、各种同步原语的使用场景,以及进程间通信的几种高效方式。

1. 核心能力速览

在深入细节之前,我们先通过一个表格快速了解 Python 并发编程的主要技术模块及其核心要点:

能力项说明与关键点
技术范畴多线程 (threading)、多进程 (multiprocessing)、异步 I/O (asyncio)
核心目标提升程序执行效率,充分利用多核 CPU,处理高并发 I/O 任务
适用场景I/O 密集型任务(网络请求、文件读写)、计算密集型任务(需绕开 GIL)
硬件/环境门槛无特殊要求,现代操作系统(Windows/macOS/Linux)和 Python 环境即可
性能影响关键GIL(全局解释器锁):导致 CPython 中多线程无法并行执行 CPU 密集型代码
同步机制锁 (Lock)、递归锁 (RLock)、信号量 (Semaphore)、条件变量 (Condition)、事件 (Event)、屏障 (Barrier)
进程通信队列 (Queue)、管道 (Pipe)、共享内存 (Value,Array)、管理器 (Manager)
线程局部数据threading.localThreadLocal,为每个线程存储独立状态
启动与监控线程/进程的启动 (start())、等待 (join())、守护模式 (daemon)
常见风险竞态条件、死锁、线程/进程间数据污染、资源泄露

2. 适用场景与使用边界

Python 并发编程并非银弹,选对技术方案才能事半功倍。

适合使用多线程 (threading) 的场景:

  • I/O 密集型操作:这是多线程的主场。例如,编写网络爬虫同时请求多个网页,开发 Web 服务器处理多个客户端连接,或者批量读写本地文件。线程在等待 I/O 响应时会释放 GIL,其他线程可以继续执行,从而显著提升整体吞吐量。
  • 需要维护响应性的 GUI 应用:在图形界面程序中,将耗时操作(如文件处理、网络通信)放入子线程,可以防止主界面“卡死”。
  • 处理大量短期、独立的任务:任务本身计算量小,但数量多,且涉及等待。

适合使用多进程 (multiprocessing) 的场景:

  • CPU 密集型计算:如图像处理、科学计算、数据加密解密、模型推理等。多进程可以绕过 GIL 限制,真正利用多核 CPU 进行并行计算。
  • 需要更高隔离性的任务:每个进程拥有独立的内存空间,一个进程崩溃不会直接影响其他进程,稳定性更好。
  • 利用多台机器进行分布式计算multiprocessing模块的某些组件(如Manager)为分布式扩展提供了基础。

需要谨慎或避免的情况:

  • 过度并发:盲目创建大量线程或进程会导致系统调度开销剧增,性能反而下降。通常需要使用线程池 (ThreadPoolExecutor) 或进程池 (ProcessPoolExecutor) 来限制并发数。
  • 共享状态过于复杂:如果线程或进程间需要频繁通信和同步复杂的数据结构,设计和调试难度会指数级上升,此时应考虑是否能用更简单的架构(如消息队列)来解耦。
  • 对实时性要求极高:Python 的 GIL 和垃圾回收机制会带来不确定的延迟,不适合硬实时系统。

3. 环境准备与前置条件

开始实战之前,确保你的环境已就绪。

1. 基础环境:

  • 操作系统:Windows 10/11, macOS, 或 Linux 发行版(如 Ubuntu)均可。本文示例代码在主流系统上通用。
  • Python 版本:推荐使用Python 3.7 及以上asyncioconcurrent.futures模块在较新版本中功能更完善、性能更好。使用以下命令检查版本:
python --version # 或 python3 --version

2. 开发工具:

  • 代码编辑器/IDE:Visual Studio Code (VSCode)、PyCharm、Jupyter Notebook 等任选其一。确保已配置好 Python 扩展和调试环境。
  • 包管理:通常使用pip。确保能正常安装第三方包(非必须,但后续示例可能用到)。

3. 理解核心模块:Python 标准库已包含我们所需的大部分工具,无需额外安装:

  • threading: 用于创建和管理线程。
  • multiprocessing: 用于创建和管理进程。
  • concurrent.futures: 提供了高级的线程池和进程池接口。
  • queue(Queue): 线程安全的队列,用于线程间通信。multiprocessing模块也有自己的Queue
  • asyncio: 用于编写单线程并发代码(协程),本文重点在线程和进程,但会简要提及它与前两者的区别。

4. 从基础到实战:多线程编程

4.1 创建与启动线程

最基础的方式是继承threading.Thread类。

import threading import time class MyThread(threading.Thread): def __init__(self, name): super().__init__() self.name = name def run(self): # 线程启动后执行的任务 print(f"线程 {self.name} 开始运行") time.sleep(2) # 模拟耗时操作 print(f"线程 {self.name} 运行结束") if __name__ == "__main__": print("主线程开始") t1 = MyThread("Thread-1") t2 = MyThread("Thread-2") t1.start() # 启动线程,异步执行 t2.start() t1.join() # 主线程等待 t1 结束 t2.join() # 主线程等待 t2 结束 print("主线程结束")

关键点

  • start(): 启动线程,调用后会执行类中的run()方法。
  • join(): 阻塞当前线程(通常是主线程),直到被调用join()的线程执行完毕。如果不调用join(),主线程可能提前结束,导致子线程被强制终止。
  • daemon属性:如果将线程设置为守护线程 (t1.daemon = True),则主线程结束时,无论守护线程是否完成,都会随主线程一起退出。

4.2 使用线程池简化管理

手动管理大量线程很繁琐,concurrent.futures模块的ThreadPoolExecutor是更优选择。

from concurrent.futures import ThreadPoolExecutor, as_completed import time def task(n): """模拟一个耗时任务""" print(f"开始执行任务 {n}") time.sleep(n) # 睡眠 n 秒模拟任务耗时 return f"任务 {n} 完成" if __name__ == "__main__": # 创建一个最大线程数为 3 的线程池 with ThreadPoolExecutor(max_workers=3) as executor: # 提交任务到线程池,返回 Future 对象 future_to_task = {executor.submit(task, i): i for i in [5, 2, 3, 1, 4]} # 使用 as_completed 获取已完成的任务结果 for future in as_completed(future_to_task): task_id = future_to_task[future] try: result = future.result() # 获取任务返回值,如果发生异常会在此抛出 print(f"获取到结果: {result}") except Exception as exc: print(f"任务 {task_id} 产生了异常: {exc}")

优势

  • 资源复用:避免频繁创建销毁线程的开销。
  • 流量控制:通过max_workers限制并发数,防止系统过载。
  • 结果获取灵活submit返回Future对象,支持result()阻塞获取,或as_completed()wait()等非阻塞方式。

5. 深入多进程编程

当遇到 CPU 密集型任务时,多进程是突破 GIL 限制的关键。

5.1 创建进程与进程池

创建进程的方式与线程类似,但使用的是multiprocessing.Process

import multiprocessing import os import time def cpu_bound_task(number): """一个模拟的CPU密集型计算任务""" print(f"进程 {os.getpid()} 开始计算数字 {number}") result = sum(i * i for i in range(number)) print(f"进程 {os.getpid()} 计算完成,结果: {result}") return result if __name__ == '__main__': # 多进程编程必须有的保护 print(f"主进程 ID: {os.getpid()}") start_time = time.time() processes = [] numbers = [1000000 + x for x in range(4)] # 四个计算任务 # 创建进程 for num in numbers: p = multiprocessing.Process(target=cpu_bound_task, args=(num,)) processes.append(p) p.start() # 等待所有进程结束 for p in processes: p.join() end_time = time.time() print(f"多进程总耗时: {end_time - start_time:.2f} 秒")

重要提示:多进程代码必须放在if __name__ == '__main__':保护块下,尤其是在 Windows 系统上,这是为了避免子进程无限递归创建。

5.2 使用进程池处理批量任务

与线程池类似,进程池 (ProcessPoolExecutor) 是管理进程的最佳实践。

from concurrent.futures import ProcessPoolExecutor import math PRIMES = [ 112272535095293, 112582705942171, 112272535095293, 115280095190773, 115797848077099, 1099726899285419 ] def is_prime(n): """判断一个数是否为质数(CPU密集型)""" if n < 2: return False if n == 2: return True if n % 2 == 0: return False sqrt_n = int(math.floor(math.sqrt(n))) for i in range(3, sqrt_n + 1, 2): if n % i == 0: return False return True if __name__ == '__main__': with ProcessPoolExecutor() as executor: # 将任务映射到进程池 results = executor.map(is_prime, PRIMES) for number, result in zip(PRIMES, results): print(f"{number} 是质数吗? {result}")

executor.map的优势:它保持了原始列表的顺序,返回结果的顺序与输入参数的顺序一致,非常直观。

6. 线程同步与数据安全

当多个线程访问和修改同一共享资源时,如果不加控制,就会导致竞态条件,产生不可预知的结果。同步机制就是用来解决这个问题的。

6.1 使用锁 (Lock) 保护临界区

锁是最基本的同步原语,确保同一时刻只有一个线程能进入“临界区”代码。

import threading class BankAccount: def __init__(self, initial_balance=0): self.balance = initial_balance self.lock = threading.Lock() # 创建一把锁 def deposit(self, amount): # 存款操作需要加锁保护 with self.lock: # 使用 with 语句自动获取和释放锁 new_balance = self.balance + amount # 模拟一点延迟,增加竞态条件发生概率 threading.Event().wait(0.0001) self.balance = new_balance print(f"存入 {amount}, 当前余额: {self.balance}") def withdraw(self, amount): # 取款操作同样需要加锁 with self.lock: if self.balance >= amount: new_balance = self.balance - amount threading.Event().wait(0.0001) self.balance = new_balance print(f"取出 {amount}, 当前余额: {self.balance}") return True else: print(f"取款 {amount} 失败,余额不足") return False def perform_transactions(account): for _ in range(100): account.deposit(5) account.withdraw(5) if __name__ == "__main__": account = BankAccount(100) print(f"初始余额: {account.balance}") threads = [] for i in range(10): # 创建10个线程并发操作账户 t = threading.Thread(target=perform_transactions, args=(account,)) threads.append(t) t.start() for t in threads: t.join() print(f"最终余额 (应为 100): {account.balance}")

运行结果:如果去掉with self.lock,最终余额很可能不是 100。加上锁后,无论运行多少次,最终余额都稳定为 100。

6.2 其他同步工具速览

  • RLock(可重入锁):允许同一个线程多次获取同一把锁,防止单线程内自己锁死自己。
  • Semaphore(信号量):控制同时访问特定资源的线程数量。例如,限制数据库连接池的最大并发数。
    semaphore = threading.Semaphore(5) # 最多允许5个线程同时进入 with semaphore: # 访问受限资源 pass
  • Condition(条件变量):用于复杂的线程间协调,一个线程等待某个条件成立,另一个线程在条件改变时通知等待的线程。典型生产者-消费者模型。
  • Event(事件):一个线程发出“事件”信号,一个或多个其他线程等待该事件。用于简单的线程间通知。
  • Barrier(屏障):让多个线程都到达某个点后再同时继续执行。

7. 进程间通信 (IPC)

进程拥有独立的内存空间,不能像线程那样直接共享变量。Python 的multiprocessing模块提供了多种安全的通信方式。

7.1 队列 (Queue)

队列是进程间通信最常用、最安全的方式之一,它实现了先进先出(FIFO)的机制。

import multiprocessing import time import random def producer(queue, items): """生产者进程:向队列中放入数据""" for item in items: print(f"生产者 放入: {item}") queue.put(item) # 放入队列 time.sleep(random.uniform(0.1, 0.5)) # 模拟生产耗时 # 放入结束信号 queue.put(None) print("生产者 完成") def consumer(queue, name): """消费者进程:从队列中取出并处理数据""" while True: item = queue.get() # 从队列获取,如果队列为空则会阻塞 if item is None: # 收到结束信号 queue.put(None) # 将结束信号放回,以便其他消费者也能结束 print(f"消费者 {name} 结束") break print(f"消费者 {name} 处理: {item}") time.sleep(random.uniform(0.2, 0.8)) # 模拟处理耗时 if __name__ == '__main__': # 创建进程间通信的队列 queue = multiprocessing.Queue(maxsize=3) # 设置队列最大容量为3 # 准备数据 data_to_produce = [f'产品-{i}' for i in range(10)] # 创建进程 prod = multiprocessing.Process(target=producer, args=(queue, data_to_produce)) cons1 = multiprocessing.Process(target=consumer, args=(queue, 'C1')) cons2 = multiprocessing.Process(target=consumer, args=(queue, 'C2')) # 启动进程 cons1.start() cons2.start() time.sleep(1) # 让消费者先启动并等待 prod.start() # 等待进程结束 prod.join() cons1.join() cons2.join() print("所有进程执行完毕")

7.2 管道 (Pipe) 与共享内存

管道 (Pipe)提供双向或单向的连接,适用于两个进程间的直接通信。

from multiprocessing import Process, Pipe def worker(conn): """子进程函数""" conn.send(['hello', 'from', 'worker']) # 发送数据到管道 received = conn.recv() # 从管道接收数据 print(f"Worker received: {received}") conn.close() if __name__ == '__main__': parent_conn, child_conn = Pipe() # 创建管道两端 p = Process(target=worker, args=(child_conn,)) p.start() print(f"Parent received: {parent_conn.recv()}") # 接收子进程数据 parent_conn.send('ack from parent') # 发送数据给子进程 p.join()

共享内存 (Value,Array)允许多个进程直接操作同一块内存,速度最快,但需要开发者自己处理同步问题。

from multiprocessing import Process, Value, Array def increment_shared_counter(counter): for _ in range(100000): counter.value += 1 # 这是一个竞态条件!实际应用需要加锁。 def safe_increment_shared_counter(counter, lock): for _ in range(100000): with lock: counter.value += 1 if __name__ == '__main__': # 使用 Value 共享一个整型 shared_counter = Value('i', 0) # 'i' 代表 C 语言的 int 类型 lock = multiprocessing.Lock() # 为共享变量创建一把锁 processes = [] for i in range(5): p = Process(target=safe_increment_shared_counter, args=(shared_counter, lock)) processes.append(p) p.start() for p in processes: p.join() print(f"最终计数器值 (应为 500000): {shared_counter.value}")

8. ThreadLocal:线程的私有储物柜

threading.local()或被称为ThreadLocal的机制,用于解决线程间数据隔离的问题。它为每个线程创建一个独立的命名空间,存储只属于该线程的数据,避免了在函数调用间传递参数的繁琐,也避免了全局变量被多线程污染的风险。

典型场景:Web 框架中,每个请求在一个独立线程中处理,需要存储当前请求的用户信息、数据库会话等。

import threading import time import random # 创建一个 ThreadLocal 对象 local_data = threading.local() def get_current_thread_name(): """获取当前线程名""" return threading.current_thread().name def process_user_request(user_id): """模拟处理用户请求,每个请求需要独立的用户上下文""" # 将用户ID存储到当前线程的 local_data 中 local_data.user_id = user_id local_data.start_time = time.time() # 模拟一些处理逻辑,这些逻辑可能调用其他函数 print(f"[{get_current_thread_name()}] 开始处理用户 {local_data.user_id} 的请求") time.sleep(random.uniform(0.5, 1.5)) # 模拟处理耗时 # 在其他地方(如下面的 log_request 函数)可以直接访问,无需传递参数 log_request() def log_request(): """记录请求日志,无需显式传递 user_id""" # 直接从当前线程的 local_data 中获取信息 duration = time.time() - local_data.start_time print(f"[{get_current_thread_name()}] 用户 {local_data.user_id} 的请求处理完毕,耗时 {duration:.2f} 秒") if __name__ == "__main__": # 模拟多个并发用户请求 user_ids = [101, 102, 103, 104, 105] threads = [] for uid in user_ids: t = threading.Thread(target=process_user_request, args=(uid,), name=f"Thread-{uid}") threads.append(t) t.start() for t in threads: t.join() print("所有请求处理完成")

输出示例

[Thread-101] 开始处理用户 101 的请求 [Thread-102] 开始处理用户 102 的请求 [Thread-103] 开始处理用户 103 的请求 [Thread-101] 用户 101 的请求处理完毕,耗时 1.23 秒 [Thread-102] 用户 102 的请求处理完毕,耗时 0.87 秒 ...

关键点local_data.user_id在每个线程中都是独立的。Thread-101无法访问到Thread-102存储的user_id。这完美解决了线程间数据隔离的问题。

9. 异步编程 (asyncio) 简要对比

虽然本文主题是线程和进程,但asyncio作为现代 Python 高并发的重要方案,有必要了解其定位。

  • 核心思想:单线程内通过协程(async/await)实现并发,在遇到 I/O 等待时主动让出控制权,去执行其他协程,从而在 I/O 密集型场景下获得极高的并发能力。
  • 与多线程对比
    • 优势:没有线程切换开销,并发数可轻松上万;无需处理复杂的线程同步问题(因为本质是单线程)。
    • 劣势:如果一个协程执行了 CPU 密集型计算且不主动await,会阻塞整个事件循环;所有代码必须是异步风格(async/await),改造现有同步代码有成本。
  • 如何选择
    • asyncio:适用于高并发 I/O 操作,如微服务、爬虫、实时消息推送。典型框架:FastAPI, aiohttp。
    • threading:适用于涉及阻塞 I/O(如某些同步数据库驱动、文件操作)且并发量不是特别巨大的场景,或需要维护响应性的 GUI 程序。
    • multiprocessing:适用于 CPU 密集型计算,需要利用多核。

10. 常见问题与排查方法

在并发编程实践中,你会遇到各种“坑”。下表汇总了常见问题及解决思路:

问题现象可能原因排查方式解决方案
程序运行结果不一致或数据错乱竞态条件:多个线程/进程同时修改共享数据未加锁。检查所有对共享变量(如全局变量、类属性)的写操作。使用日志或调试器观察执行顺序。使用Lock,RLock等同步原语保护临界区。优先使用队列 (Queue) 进行通信,而非共享内存。
程序卡死,无任何输出死锁:两个及以上线程/进程互相等待对方释放锁。检查锁的获取顺序。是否存在嵌套锁且顺序不一致?1. 使用with lock上下文管理器,避免忘记释放锁。
2. 设计固定的锁获取顺序。
3. 使用RLock或设置锁的超时参数lock.acquire(timeout=5)
多进程程序在 Windows 下报错或行为异常Windows 系统下,子进程会重新导入主模块 (__main__),导致代码重复执行。检查是否将所有进程启动代码放在了if __name__ == '__main__':保护块内。必须将进程启动逻辑放入if __name__ == '__main__':中。
创建大量线程/进程后程序变慢或崩溃1. 系统资源(内存、上下文切换开销)耗尽。
2. GIL 争抢严重(多线程)。
监控系统资源使用情况(任务管理器、topps)。观察 CPU 使用率是否真的上去了(多进程)。使用线程池 (ThreadPoolExecutor) 或进程池 (ProcessPoolExecutor) 限制并发数量。根据任务类型(I/O 或 CPU)选择正确模型。
使用multiprocessing.Queueputget阻塞队列已满(put阻塞)或队列为空(get阻塞)。默认情况下,队列大小无限制,put不会阻塞;若设置了maxsize则可能阻塞。检查生产者和消费者的速度是否匹配。是否设置了maxsize1. 使用put_nowait()get_nowait()(需捕获queue.Full/queue.Empty异常)。
2. 使用put(item, block=True, timeout=5)设置超时。
3. 调整生产/消费逻辑或队列大小。
子线程/进程中的异常未被主线程捕获子线程/进程的异常默认不会传播到父线程/进程。子线程中使用try...except捕获并记录日志。对于concurrent.futures,通过future.result()future.exception()获取异常。1. 为线程函数设置完善的异常处理。
2. 使用ThreadPoolExecutor时,迭代as_completed(futures)并调用future.result(),它会抛出子线程中的异常。
ThreadLocal数据在第一次访问时报AttributeError在当前线程的local对象上访问了一个尚未设置的属性。检查是否在所有可能执行到的线程路径上都初始化了该属性。1. 在访问前先判断属性是否存在 (if hasattr(local_data, 'user_id'))。
2. 为local对象设置默认值(例如,重载__getattr__方法或使用一个包装类)。

11. 最佳实践与使用建议

  1. 优先使用高层抽象:除非有特殊需求,否则优先使用concurrent.futures模块的ThreadPoolExecutorProcessPoolExecutor,而不是手动管理ThreadProcess对象。池化技术能自动管理生命周期,避免资源泄露。
  2. 理解 GIL,正确选型
    • I/O 密集型-> 多线程 (threading) 或异步 (asyncio)。
    • CPU 密集型-> 多进程 (multiprocessing)。
    • 不确定时,可以写一个小型测试程序对比两种方式的性能。
  3. 避免共享状态:多线程/进程编程的复杂性主要来源于共享状态。尽量设计无状态的函数,通过队列 (Queue) 传递消息和数据,这是最清晰、最安全的方式。
  4. 锁的粒度要细,持有时间要短:只锁住真正需要保护的共享数据,并且一旦操作完成立即释放锁。长时间持有锁会严重降低并发性能。
  5. 善用ThreadLocal:对于需要在线程内传递的上下文信息(如请求ID、用户会话、数据库连接),使用threading.local()是优雅的解决方案。
  6. 做好日志和错误处理:并发程序调试困难。务必为每个任务、线程、进程记录清晰的日志,包括开始、结束、关键步骤和异常。使用logging模块,并考虑为日志记录器添加线程/进程ID。
  7. 设置超时:对任何可能阻塞的操作(如lock.acquire(),queue.get(),thread.join())都设置合理的超时时间,防止程序因意外情况永久挂起。
  8. 资源清理:确保线程/进程结束时能正确释放其占用的资源(如文件句柄、网络连接、数据库连接)。使用with语句或try...finally块来保证。

掌握 Python 并发编程,意味着你能让程序真正“跑起来”,充分利用硬件资源。从理解 GIL 开始,分清 I/O 与 CPU 密集型任务的区别,然后熟练运用多线程处理等待,使用多进程榨干 CPU 性能,最后用同步机制和进程通信工具解决协作难题。记住,先从高层抽象的concurrent.futures用起,遇到复杂场景再深入threadingmultiprocessing的底层原语。多写、多测、多观察日志,是掌握这门技术的不二法门。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询