ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Python多线程编程实战:从基础到高级应用与性能优化

2026/8/13 8:15:01 拓冰建站 浏览量
Python多线程编程实战:从基础到高级应用与性能优化

1. 从“单打独斗”到“协同作战”:为什么需要并发与多线程?

如果你写过一些Python脚本,处理过文件、爬取过网页数据,或者做过简单的数据分析,你大概率已经习惯了程序“一条道走到黑”的执行方式。代码从上到下,一行接一行,一个任务做完再做下一个。这在处理小数据量、简单任务时完全没问题,运行起来也很快。但当你需要处理一个包含十万条记录的CSV文件,或者需要同时监控十几个传感器的实时数据流时,这种“单线程”的模式就会让你感到力不从心。程序会“卡”在某个耗时的操作上,比如等待网络请求返回、等待磁盘I/O完成,此时CPU只能干等着,什么也做不了,用户体验就是“程序未响应”。

这就是并发和多线程要解决的问题。想象一下你一个人在厨房做饭,如果按照“单线程”模式,你得先烧水,等水烧开,下面条,等面条煮熟,捞出来,然后再开始切菜、炒菜。整个过程耗时很长,大部分时间你都在“等待”。而“并发”就像是请了一个帮手(多线程),你可以在烧水的同时切菜,在煮面的同时调酱汁。虽然本质上还是你一个人在厨房(一个CPU核心)里忙活,但通过合理地安排任务和利用等待时间,整体效率大大提升。如果厨房够大,有多个灶台(多CPU核心),那甚至可以真正地“同时”进行多个任务,这就是“并行”。

在Python的世界里,多线程是实现并发编程最直观的方式之一。它允许你在一个程序内部创建多个“执行流”,这些线程共享程序的内存空间,可以“同时”执行不同的任务。对于I/O密集型任务(如网络请求、文件读写、数据库查询),多线程能显著提升程序的响应速度和吞吐量,因为当一个线程在等待I/O时,CPU可以立刻切换到另一个就绪的线程去执行计算。尽管Python因为GIL(全局解释器锁)的存在,在多核CPU上进行纯计算密集型任务时,多线程无法实现真正的并行加速,但这并不妨碍它成为处理I/O密集型并发场景的利器。理解并掌握多线程,是Python开发者从编写脚本迈向构建高效应用程序的关键一步。

2. 线程的创建与管理:从threading模块开始

Python标准库中的threading模块为我们提供了构建多线程程序所需的一切基础工具。与更低级的_thread模块相比,threading模块进行了更高层次的封装,功能更强大,接口也更友好。我们首先从最核心的线程对象Thread开始。

2.1 创建线程的两种核心方式

创建线程主要有两种方式:直接实例化Thread类,或者继承Thread类并重写run方法。第一种方式更为灵活和常用。

方式一:传入目标函数

这是最直接、最清晰的方式。你只需要定义一个普通的函数作为线程要执行的任务,然后将这个函数作为target参数传递给Thread构造函数。

import threading import time def download_file(filename): """模拟下载文件的任务""" print(f"[{threading.current_thread().name}] 开始下载 {filename}") time.sleep(2) # 模拟耗时的网络I/O print(f"[{threading.current_thread().name}] 完成下载 {filename}") if __name__ == "__main__": print(f"[主线程 {threading.current_thread().name}] 启动下载任务") # 创建线程对象,指定目标函数和参数 thread1 = threading.Thread(target=download_file, args=("movie.mp4",), name="下载线程-1") thread2 = threading.Thread(target=download_file, args=("music.zip",), name="下载线程-2") # 启动线程 thread1.start() thread2.start() print(f"[主线程 {threading.current_thread().name}] 已启动所有下载线程,继续处理其他事情...") # 主线程可以继续执行其他任务 time.sleep(1) print(f"[主线程 {threading.current_thread().name}] 其他事情处理完毕") # 等待子线程结束 thread1.join() thread2.join() print(f"[主线程 {threading.current_thread().name}] 所有下载任务完成")

关键点解析:

  • target: 指定线程要执行的函数对象。
  • args: 以元组形式传递给目标函数的参数。如果只有一个参数,记得加逗号,如args=("movie.mp4",),否则Python会将其视为一个普通括号表达式。
  • name: 为线程设置一个易于识别的名字,这在调试多线程程序时非常有用,日志输出会更清晰。
  • start(): 调用此方法后,线程进入“就绪”状态,由操作系统调度执行。注意,start()只能调用一次。
  • join([timeout]): 阻塞当前线程(通常是主线程),直到调用join()的线程执行完毕。timeout参数可以设置最长等待时间(秒),超时后join()方法返回,但线程可能仍在运行。在主线程中调用join()是为了防止主线程提前退出导致程序结束,子线程被强制终止。

方式二:继承Thread

这种方式更适合需要将线程与复杂对象状态绑定的场景,例如一个长期运行的后台服务线程。

class WorkerThread(threading.Thread): def __init__(self, task_queue): super().__init__() # 必须调用父类初始化 self.task_queue = task_queue self._stop_event = threading.Event() # 用于优雅停止线程 def run(self): """重写run方法,定义线程的主体逻辑""" print(f"[{self.name}] 工作者线程启动") while not self._stop_event.is_set(): try: # 从队列中获取任务,设置超时以定期检查停止事件 task = self.task_queue.get(timeout=0.5) print(f"[{self.name}] 处理任务: {task}") self.task_queue.task_done() # 通知队列任务已完成 except queue.Empty: continue # 队列为空,继续循环 print(f"[{self.name}] 工作者线程停止") def stop(self): """请求线程停止""" self._stop_event.set()

使用建议:对于大多数情况,推荐使用第一种方式(传入目标函数)。它更符合“组合优于继承”的原则,将线程执行逻辑与线程对象本身解耦,代码更清晰,也更容易进行单元测试。继承方式仅在需要高度定制线程行为或封装复杂状态时使用。

2.2 守护线程:那些“后台运行”的线程

守护线程(Daemon Thread)是一种特殊的线程,它的生命周期依赖于主线程(或非守护线程)。当程序中所有非守护线程都结束时,无论守护线程是否执行完毕,Python解释器都会强制退出,守护线程也随之终止。

def background_logger(): import datetime while True: # 模拟一个持续记录日志的后台任务 with open("app.log", "a") as f: f.write(f"{datetime.datetime.now()}: Heartbeat\n") time.sleep(5) if __name__ == "__main__": # 创建守护线程,daemon=True logger_thread = threading.Thread(target=background_logger, daemon=True) logger_thread.start() print("主程序开始运行...") time.sleep(12) # 主程序运行12秒 print("主程序运行结束,即将退出。此时守护线程会被强制终止。") # 程序退出,日志线程的循环被中断,不会执行完下一次sleep和写日志。

核心区别与选择:

  • 非守护线程(默认):主线程必须等待其结束。适用于必须完成的关键任务,如数据处理、结果保存。
  • 守护线程:主线程无需等待其结束。适用于非关键的后台服务,如心跳检测、缓存刷新、监控信息收集。即使这些任务被突然中断,也不会影响程序核心逻辑的正确性。

注意:守护线程在退出时不会执行finally子句,也不会正常地清理资源(如关闭文件、释放锁)。因此,如果线程持有锁或打开了需要关闭的资源,应避免将其设置为守护线程,或者实现明确的停止机制。

2.3 线程的生命周期与状态查询

一个线程从创建到销毁,会经历多个状态。threading模块提供了查询这些状态的方法,对于调试复杂的并发问题至关重要。

import threading import time def worker(): time.sleep(1) t = threading.Thread(target=worker, name="示例线程") print(f"线程创建后,是否存活? {t.is_alive()}") # False print(f"线程标识符: {t.ident}") # None,启动后才有 print(f"线程名称: {t.name}") t.start() print(f"线程启动后,是否存活? {t.is_alive()}") # True print(f"线程标识符: {t.ident}") # 一个非零整数 print(f"线程原生ID (可通过系统工具查看): {t.native_id}") # Python 3.8+ # 在主线程中,我们可以获取所有活跃线程的信息 for thread in threading.enumerate(): print(f"活跃线程: {thread.name} (ID: {thread.ident}, 是否守护: {thread.daemon})") t.join() print(f"线程结束后,是否存活? {t.is_alive()}") # False

状态解读:

  • 初始状态Thread对象创建后,处于“新建”状态,is_alive()返回False
  • 就绪/运行:调用start()后,线程进入“就绪”状态,由操作系统调度执行。此时is_alive()返回Trueident被赋值。
  • 阻塞:线程可能因为time.sleep()、等待I/O、等待锁而进入“阻塞”状态,但它仍然是存活的。
  • 终止:线程函数执行完毕或出现未处理异常后,线程进入“终止”状态,is_alive()返回False。一个终止的线程无法再次启动。

threading.enumerate()函数在调试时非常有用,它可以列出当前所有存活的线程对象,帮助你理解程序在某一时刻的并发结构。

3. 线程间的通信与协调:共享数据与同步原语

当多个线程需要访问或修改同一个资源(如一个变量、一个列表、一个文件)时,就会产生“竞态条件”。如果不加控制,程序的运行结果将变得不可预测。例如,一个经典的“丢失更新”问题:

import threading counter = 0 def increment(): global counter for _ in range(100000): counter += 1 # 这个操作不是原子的! threads = [] for i in range(10): t = threading.Thread(target=increment) threads.append(t) t.start() for t in threads: t.join() print(f"理论结果: 1000000, 实际结果: {counter}")

运行多次,你会发现结果几乎每次都小于1000000。这是因为counter += 1这个语句实际上包含了三个步骤:读取counter的值、将值加1、写回counter。两个线程可能同时读取到相同的值,然后各自加1后写回,导致其中一次增加“丢失”了。

3.1 互斥锁:保护共享资源的“门卫”

解决上述问题最直接的工具就是互斥锁(threading.Lock)。锁就像一个房间的钥匙,一次只允许一个线程进入“临界区”(访问共享资源的代码段)。

import threading counter = 0 counter_lock = threading.Lock() # 创建一把锁 def increment_with_lock(): global counter for _ in range(100000): # 进入临界区前获取锁 counter_lock.acquire() try: counter += 1 finally: # 无论是否发生异常,都必须释放锁,否则会导致死锁 counter_lock.release() threads = [] for i in range(10): t = threading.Thread(target=increment_with_lock) threads.append(t) t.start() for t in threads: t.join() print(f"使用锁后的结果: {counter}") # 正确输出 1000000

使用锁的最佳实践:

  1. 使用with语句(推荐)Lock对象支持上下文管理器协议,使用with可以自动获取和释放锁,代码更简洁,且能确保异常发生时锁也能被释放。
    def increment_with_lock_better(): global counter for _ in range(100000): with counter_lock: # 自动获取和释放锁 counter += 1
  2. 锁的粒度要适中:锁保护的范围(临界区)越小越好,只包含真正需要互斥访问的代码。锁的粒度过大会严重降低并发性能,因为其他线程需要等待更长时间。
  3. 避免嵌套锁与死锁:如果一个线程在持有锁A的情况下去请求锁B,而另一个线程持有锁B并请求锁A,就会发生死锁。设计时应尽量避免嵌套锁,或使用threading.RLock(可重入锁)允许同一线程多次获取同一把锁。

3.2 线程安全的数据结构:queue.Queue

对于生产者-消费者这类模型,使用queue.Queue是比手动操作锁更安全、更高效的选择。Queue是线程安全的,内部已经实现了所有必要的锁机制。

import threading import queue import time import random def producer(q, producer_id): """生产者线程:生成任务放入队列""" for i in range(5): item = f"产品-{producer_id}-{i}" time.sleep(random.uniform(0.1, 0.5)) # 模拟生产耗时 q.put(item) print(f"[生产者{producer_id}] 生产了 {item}") print(f"[生产者{producer_id}] 生产完毕") def consumer(q, consumer_id): """消费者线程:从队列取出任务处理""" while True: try: # block=True, timeout=1 表示最多阻塞1秒等待物品 item = q.get(block=True, timeout=1) except queue.Empty: # 超时后队列仍为空,认为所有生产者都已结束,消费者也退出 print(f"[消费者{consumer_id}] 等待超时,退出") break # 处理物品 time.sleep(random.uniform(0.2, 0.8)) # 模拟消费耗时 print(f"[消费者{consumer_id}] 处理了 {item}") q.task_done() # 非常重要!通知队列该任务已完成 if __name__ == "__main__": task_queue = queue.Queue(maxsize=3) # 设置队列最大容量为3 # 创建2个生产者,3个消费者 producers = [threading.Thread(target=producer, args=(task_queue, i)) for i in range(2)] consumers = [threading.Thread(target=consumer, args=(task_queue, i)) for i in range(3)] for p in producers: p.start() for c in consumers: c.start() # 等待所有生产者结束 for p in producers: p.join() # 等待队列中所有任务被消费者处理完 task_queue.join() # 阻塞,直到队列中每个item都调用了task_done() print("所有任务生产并消费完毕") # 此时消费者线程因get超时而陆续退出 for c in consumers: c.join() print("程序结束")

Queue的核心方法:

  • put(item, block=True, timeout=None): 放入项目。如果队列满,block=True时会阻塞直到有空位;timeout设置阻塞超时时间。
  • get(block=True, timeout=None): 取出项目。如果队列空,block=True时会阻塞直到有项目;timeout设置阻塞超时时间。
  • task_done(): 消费者处理完一个从get()得到的项目后调用。用于通知队列该任务已完成。
  • join(): 阻塞调用者,直到队列中所有项目都被处理(即每个put()进来的项目都对应调用了task_done())。这在协调生产者和消费者结束时非常有用。

使用Queue极大地简化了线程间通信的复杂度,你不再需要关心底层的锁和条件变量,只需关注业务逻辑。

3.3 条件变量与事件:更复杂的线程协调

当线程间的协作不仅仅是传递数据,还需要等待某个条件成立时,就需要用到threading.Conditionthreading.Event

Event:一次性通知机制Event对象内部有一个标志位。线程可以wait()这个事件,阻塞直到标志位被设置为True;其他线程可以set()这个事件,唤醒所有等待的线程。clear()可以将标志位重置为False

import threading import time # 模拟一个资源准备场景 resource_ready = threading.Event() def resource_loader(): print("[资源加载器] 开始加载资源...") time.sleep(3) # 模拟加载耗时 print("[资源加载器] 资源加载完毕!") resource_ready.set() # 发出“资源就绪”信号 def worker(worker_id): print(f"[工作者{worker_id}] 等待资源就绪...") resource_ready.wait() # 阻塞,直到事件被set print(f"[工作者{worker_id}] 检测到资源就绪,开始工作!") # ... 执行具体工作 if __name__ == "__main__": loader = threading.Thread(target=resource_loader) workers = [threading.Thread(target=worker, args=(i,)) for i in range(3)] loader.start() for w in workers: w.start() loader.join() for w in workers: w.join()

Condition:基于锁的复杂条件等待Condition通常与一个共享状态(如队列长度、某个标志)结合使用。它允许线程在某个条件不满足时释放锁并等待,直到被其他线程通知条件可能已改变。

import threading import time import collections class BoundedBuffer: """一个有限容量的缓冲区,生产者-消费者模型的另一种实现""" def __init__(self, capacity): self.capacity = capacity self.buffer = collections.deque(maxlen=capacity) self.lock = threading.Lock() self.not_full = threading.Condition(self.lock) # 条件:缓冲区未满 self.not_empty = threading.Condition(self.lock) # 条件:缓冲区非空 def put(self, item): with self.lock: # Condition内部也使用这把锁 # 等待“缓冲区未满”的条件成立 while len(self.buffer) == self.capacity: self.not_full.wait() # 释放锁并等待,被唤醒后重新获取锁 self.buffer.append(item) print(f"[生产者] 放入 {item}, 缓冲区大小: {len(self.buffer)}") self.not_empty.notify() # 通知可能正在等待“非空”的消费者 def get(self): with self.lock: # 等待“缓冲区非空”的条件成立 while len(self.buffer) == 0: self.not_empty.wait() item = self.buffer.popleft() print(f"[消费者] 取出 {item}, 缓冲区大小: {len(self.buffer)}") self.not_full.notify() # 通知可能正在等待“未满”的生产者 return item # 使用示例略,与Queue类似,但展示了Condition的底层用法。

选择建议:

  • 简单的“准备好了吗?”通知,用Event
  • 复杂的、与共享状态相关的等待/通知逻辑,用Condition
  • 绝大多数生产者-消费者场景,直接用queue.Queue,它内部就是用Condition实现的。

4. Python多线程的“阿喀琉斯之踵”:GIL与性能真相

谈到Python多线程,一个无法回避的话题就是GIL(Global Interpreter Lock,全局解释器锁)。这是CPython解释器(我们通常使用的Python)中的一个机制,它确保同一时刻只有一个线程在执行Python字节码。这意味着,即使在多核CPU上,一个Python进程中的多个线程也无法实现真正的并行计算。

4.1 GIL是如何工作的?

你可以把GIL想象成解释器的“话筒”。一个线程想要执行Python代码,必须先拿到这个“话筒”。执行一段时间后(基于ticks或时间片),它会释放话筒,然后由操作系统调度决定下一个拿到话筒的线程是谁。这个机制主要是为了简化CPython的内存管理,因为对象的引用计数操作需要保证原子性。

4.2 GIL对性能的真实影响:I/O密集型 vs CPU密集型

理解GIL的影响,关键在于区分任务类型:

I/O密集型任务:任务的大部分时间花在等待上,如网络请求、磁盘读写、数据库查询。当一个线程因I/O而阻塞时,它会自动释放GIL,其他线程就可以获得GIL并执行。因此,对于I/O密集型任务,多线程可以显著提升性能,因为线程在等待时可以切换,CPU利用率更高。

CPU密集型任务:任务的大部分时间花在计算上,如科学计算、图像处理、复杂算法。由于GIL的存在,多个线程无法同时利用多个CPU核心进行计算。它们会争抢GIL,线程切换本身还会带来开销。因此,对于纯CPU密集型任务,使用多线程通常不会带来加速,甚至可能因为锁竞争和切换开销而变慢

一个简单的测试:

import threading import time import math def cpu_bound_task(n): """一个模拟的CPU密集型任务:计算平方根""" count = 0 for i in range(n): math.sqrt(i) count += 1 return count def run_with_threads(num_threads, total_work): """使用多线程执行CPU密集型任务""" work_per_thread = total_work // num_threads threads = [] start_time = time.time() for _ in range(num_threads): t = threading.Thread(target=cpu_bound_task, args=(work_per_thread,)) threads.append(t) t.start() for t in threads: t.join() end_time = time.time() return end_time - start_time if __name__ == "__main__": total_work = 5_000_000 print("CPU密集型任务测试 (计算5百万次平方根):") for n in [1, 2, 4]: duration = run_with_threads(n, total_work) print(f" 使用 {n} 个线程: {duration:.2f} 秒") # 对比单线程 start = time.time() cpu_bound_task(total_work) single_thread_time = time.time() - start print(f" 单线程: {single_thread_time:.2f} 秒")

在我的测试环境(4核CPU)上,结果可能是:单线程最快,4个线程最慢。这清晰地展示了GIL对CPU密集型任务的限制。

4.3 突破GIL限制的实战方案

如果确实需要在Python中利用多核进行CPU密集型计算,有以下几种主流方案:

方案一:使用多进程(multiprocessing模块)每个进程有自己独立的Python解释器和内存空间,因此也有自己独立的GIL。多进程可以实现真正的并行计算。

import multiprocessing import time import math def cpu_bound_task(n): count = 0 for i in range(n): math.sqrt(i) count += 1 return count if __name__ == '__main__': # 多进程必须保护入口点 total_work = 5_000_000 num_processes = 4 work_per_process = total_work // num_processes start_time = time.time() with multiprocessing.Pool(processes=num_processes) as pool: # 将任务映射到多个进程 results = pool.map(cpu_bound_task, [work_per_process] * num_processes) end_time = time.time() print(f"使用 {num_processes} 个进程: {end_time - start_time:.2f} 秒") print(f"总计算量: {sum(results)}")

多进程的缺点是进程间通信(IPC)开销比线程间通信大,因为数据需要在不同内存空间之间传递(通常通过序列化)。

方案二:使用C扩展或利用释放GIL的库一些用C编写的底层库(如NumPy、SciPy、Pandas中的部分计算,以及concurrent.futuresThreadPoolExecutor在某些I/O操作中)在进行耗时运算时会主动释放GIL,从而允许其他线程运行。如果你在编写C扩展,也可以使用Py_BEGIN_ALLOW_THREADSPy_END_ALLOW_THREADS宏来临时释放GIL。

方案三:使用concurrent.futures模块这个高级模块提供了ThreadPoolExecutorProcessPoolExecutor,它们提供了统一的接口来执行并发任务,底层自动选择线程池或进程池。对于I/O密集型,用ThreadPoolExecutor;对于CPU密集型,用ProcessPoolExecutor

from concurrent.futures import ProcessPoolExecutor, as_completed import math def cpu_bound_task(n): count = 0 for i in range(n): math.sqrt(i) count += 1 return count if __name__ == '__main__': total_work = 5_000_000 num_workers = 4 work_per_worker = total_work // num_workers with ProcessPoolExecutor(max_workers=num_workers) as executor: # 提交任务 futures = [executor.submit(cpu_bound_task, work_per_worker) for _ in range(num_workers)] results = [] # 异步获取结果 for future in as_completed(futures): results.append(future.result()) print(f"总计算量: {sum(results)}")

实战选择建议:

  • Web服务器、爬虫、文件批处理等I/O密集型应用:大胆使用多线程,这是其主战场,能有效提升吞吐量。
  • 数据分析、机器学习训练、图像渲染等CPU密集型应用:优先考虑多进程(multiprocessingProcessPoolExecutor),或者直接使用专为科学计算设计的库(如NumPy、Numba、JAX),它们内部已做了并行优化。
  • 混合型任务:可以采用“多进程+多线程”的混合模式,例如用多进程利用多核,每个进程内再用多线程处理I/O。但这会大大增加程序的复杂度,需要谨慎设计。

5. 高级模式与实战避坑指南

掌握了基础之后,我们来看看如何在实际项目中更优雅、更安全地使用多线程,以及那些容易踩坑的地方。

5.1 线程池:管理线程的生命周期

频繁地创建和销毁线程是有开销的。线程池模式预先创建好一组线程,并将任务提交给池子,由池子分配空闲线程来执行,线程执行完任务后并不销毁,而是等待下一个任务。这避免了重复创建线程的开销。

Python中实现线程池主要有两种方式:

1. 使用concurrent.futures.ThreadPoolExecutor(推荐)这是现代Python中最简洁、最安全的方式。

from concurrent.futures import ThreadPoolExecutor, as_completed import urllib.request import time def download_url(url): """下载单个URL的内容""" try: with urllib.request.urlopen(url, timeout=5) as response: content = response.read() return f"{url}: 成功,长度 {len(content)} 字节" except Exception as e: return f"{url}: 失败 - {e}" urls = [ "https://www.python.org", "https://docs.python.org", "https://pypi.org", "https://www.example.com", # ... 更多URL ] def download_with_threadpool(urls, max_workers=3): """使用线程池并发下载""" results = [] with ThreadPoolExecutor(max_workers=max_workers) as executor: # 使用submit提交单个任务,返回Future对象 future_to_url = {executor.submit(download_url, url): url for url in urls} # 使用as_completed获取已完成的任务结果(按完成顺序) for future in as_completed(future_to_url): url = future_to_url[future] try: result = future.result() # 获取结果,如果任务抛出异常,这里会重新抛出 results.append(result) print(result) except Exception as exc: print(f'{url} 产生了异常: {exc}') return results # 或者使用map方法,更简洁,但结果顺序固定,且一个异常会导致整个map中断 def download_with_map(urls, max_workers=3): with ThreadPoolExecutor(max_workers=max_workers) as executor: # map会保持输入顺序和输出顺序一致 for result in executor.map(download_url, urls): print(result) if __name__ == "__main__": start = time.time() download_with_threadpool(urls, max_workers=3) print(f"耗时: {time.time() - start:.2f}秒")

2. 使用multiprocessing.pool.ThreadPoolmultiprocessing模块也提供了一个线程池,其接口与进程池类似。

from multiprocessing.pool import ThreadPool def task(x): return x * x with ThreadPool(processes=4) as pool: # 注意参数名是processes,但创建的是线程 results = pool.map(task, range(10)) print(results)

线程池大小设置经验:对于I/O密集型任务,线程池大小可以设置得较大,通常可以是CPU核心数的数倍(如10倍、20倍),具体取决于I/O等待时间与CPU计算时间的比例。一个粗略的公式是:线程数 = CPU核心数 * (1 + I/O等待时间 / CPU计算时间)。在实践中,可以通过压力测试找到一个最优值。对于CPU密集型任务,由于GIL,线程数设置超过CPU核心数通常无益。

5.2 线程局部数据:threading.local

有时,你需要一些数据只对某个线程可见,对其他线程不可见,比如数据库连接、请求上下文、用户会话等。threading.local()可以创建一个线程本地存储对象。

import threading import time # 创建一个线程本地存储对象 local_data = threading.local() def show_data(): """每个线程打印自己的数据""" try: value = local_data.value except AttributeError: print(f"[{threading.current_thread().name}] 还没有设置value") return print(f"[{threading.current_thread().name}] value = {value}") def worker(num): """每个线程设置自己独有的数据""" local_data.value = num # 这个value属性是线程独立的 time.sleep(0.1) # 模拟一些操作 show_data() threads = [] for i in range(3): t = threading.Thread(target=worker, args=(i,), name=f"Thread-{i}") threads.append(t) t.start() for t in threads: t.join() # 在主线程中访问 show_data() # 会输出“还没有设置value”,因为主线程的local_data没有value属性

threading.local的实现为每个线程维护了一个独立的字典。它非常适用于Web框架中为每个请求线程存储独立上下文的情况。

5.3 常见“坑”与最佳实践

坑1:忘记处理异常子线程中未捕获的异常会导致线程静默终止,可能不会打印任何错误信息,使得调试极其困难。

def buggy_worker(): raise ValueError("线程内部出错了!") t = threading.Thread(target=buggy_worker) t.start() t.join() print("主线程结束") # 程序会正常结束,你看不到错误信息!

解决方案:在线程函数内部用try...except捕获所有异常,并记录日志。

import traceback import logging logging.basicConfig(level=logging.INFO) def safe_worker(): try: # 你的业务逻辑 raise ValueError("出错了") except Exception as e: logging.error(f"线程 {threading.current_thread().name} 发生异常: {e}") logging.error(traceback.format_exc()) # 打印完整的堆栈跟踪

坑2:死锁两个或多个线程互相等待对方持有的锁,导致所有线程都无法继续执行。

lock_a = threading.Lock() lock_b = threading.Lock() def thread_1(): with lock_a: time.sleep(0.1) # 故意sleep,增加死锁概率 with lock_b: # 需要锁b,但锁b可能被thread_2持有 print("Thread 1 got both locks") def thread_2(): with lock_b: time.sleep(0.1) with lock_a: # 需要锁a,但锁a被thread_1持有 print("Thread 2 got both locks")

解决方案

  1. 避免嵌套锁:重新设计代码逻辑,尽量减少需要同时持有多个锁的情况。
  2. 固定锁的获取顺序:如果必须获取多个锁,确保所有线程都以相同的顺序获取它们(例如,总是先获取lock_a,再获取lock_b)。
  3. 使用带超时的锁lock.acquire(timeout=5),超时后可以执行回退逻辑。
  4. 使用高级抽象:尽可能使用queue.Queue等线程安全的数据结构,避免直接操作锁。

坑3:资源泄漏线程中打开文件、网络连接或数据库连接,如果线程异常终止,可能无法正确关闭。

def leaky_worker(): f = open('temp.txt', 'w') f.write('data') # 如果这里发生异常,文件句柄可能不会被关闭 f.close() # 正常情况应该关闭

解决方案:使用with语句(上下文管理器)来管理资源,确保即使发生异常资源也能被正确释放。

def safe_worker(): with open('temp.txt', 'w') as f: f.write('data') # 离开with块,文件自动关闭

最佳实践总结:

  1. 优先使用高层抽象:如concurrent.futures.ThreadPoolExecutorqueue.Queue,它们更安全,更不易出错。
  2. 明确线程的职责和生命周期:设计时就想好线程何时启动、何时结束、如何优雅停止(使用Event信号)。
  3. 所有共享数据都必须同步:对任何可能被多个线程修改的数据,都要考虑使用锁或线程安全数据结构。
  4. 保持简单:多线程代码本来就复杂,尽量让每个线程的逻辑简单、独立。复杂的交互尽量通过队列进行。
  5. 充分测试:多线程bug常常难以复现。需要进行压力测试、长时间运行测试,并仔细检查日志。

6. 从多线程到异步编程:asyncio的简要对比

当并发任务数量极大(成千上万)时,操作系统线程的创建、切换和内存开销会成为瓶颈。这时,异步编程模型(如Python的asyncio)就显示出其优势。它使用单线程(或少量线程)配合事件循环,通过“协程”在I/O等待时主动让出控制权,来实现高并发。

一个简单的asyncio示例:

import asyncio import aiohttp # 需要安装 aiohttp import time async def fetch_url(session, url): """异步获取URL""" async with session.get(url) as response: text = await response.text() return f"{url}: 状态码 {response.status}, 长度 {len(text)}" async def main(): urls = ["https://www.python.org", "https://docs.python.org", "https://pypi.org"] * 10 # 30个URL async with aiohttp.ClientSession() as session: tasks = [fetch_url(session, url) for url in urls] results = await asyncio.gather(*tasks) # 并发执行所有任务 for result in results[:3]: # 只打印前3个结果 print(result) if __name__ == "__main__": start = time.time() asyncio.run(main()) print(f"异步耗时: {time.time() - start:.2f}秒")

多线程 vs 异步 (asyncio) 如何选择?

特性多线程 (threading)异步 (asyncio)
编程模型基于操作系统线程,抢占式调度。基于协程,协作式调度(需要显式await让出控制)。
并发能力受限于操作系统线程数(通常数百到数千)。可轻松处理数万甚至数十万并发连接(如WebSocket服务器)。
适用场景I/O密集型任务,特别是涉及阻塞式I/O调用(如某些数据库驱动、文件操作)。高并发I/O密集型任务,尤其是网络服务。所有I/O操作都必须是异步的(使用async/await)。
CPU密集型受GIL限制,无法利用多核。同样受限于单线程,CPU密集型任务会阻塞事件循环。
调试难度较难,存在竞态条件、死锁。也难,但逻辑流更清晰(单线程),不过堆栈跟踪可能更复杂。
生态兼容兼容绝大多数同步库。需要库本身支持async/await,或者使用run_in_executor在线程池中运行同步代码。

简单决策树:

  • 如果你要处理成百上千的并发网络连接(如聊天服务器、爬虫),优先考虑asyncio
  • 如果你的任务主要是I/O密集型,但使用的库是传统的同步阻塞式API(如requests,psycopg2(同步模式)),那么多线程是更直接的选择。
  • 如果你的任务混合了CPU计算和I/O,可以考虑结合使用:用多进程处理CPU部分,用多线程或异步处理I/O部分。

多线程是Python并发编程工具箱中不可或缺的一件利器。尽管有GIL的限制,但在正确的场景(I/O密集型)下使用,它能极大地提升程序性能。理解其原理,掌握threadingqueueconcurrent.futures等核心模块,并牢记同步、通信、异常处理和资源管理的要点,你就能写出高效、健壮的多线程程序。当任务规模继续扩大,你自然会接触到asyncio和多进程等更高级的并发模型,而扎实的多线程基础将为理解它们铺平道路。