ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Python threading模块详解:多线程编程从入门到实战

Python threading模块详解:多线程编程从入门到实战 各位学习 Python 的朋友们大家好在日常编写 Python 程序时我们经常会遇到这样的场景一个爬虫需要下载 1000 张图片一个脚本需要批量读取 500 个文件一个接口需要同时向多个第三方系统发起请求。如果按照传统的方式顺序执行总耗时将是无数次 IO 等待的简单相加。这时候你肯定思考过一个问题有没有办法让多个任务“同时”进行答案是肯定的Python 中提供了threading模块来支持多线程编程。但很多初学者在第一次接触多线程时往往会被“线程同步”、“锁”、“GIL”、“死锁”这些概念吓退。网上教程虽然多但大多要么过于理论化看着昏昏欲睡要么代码片段不完整照着敲还报错。今天我们就围绕 Python 的threading模块展开一次系统整理从底层概念讲到代码实战从基础创建讲到生产环境避坑。本文会通过大量可复制的代码示例让还没有接触过多线程的读者也能平稳入门对于有一定开发经验的朋友可以重点关注线程安全和异步切异步的性能边界。考虑到电脑 CPU 的核心数正在不断增加合理使用多线程能够显著提升 IO 密集型任务的执行效率掌握这项技能在求职面试和实际项目中都很有必要。1. 多线程背景与核心概念在正式写代码之前必须先弄清楚几个可能困扰你很久的基础问题。如果连“线程”和“进程”的区别都没搞懂后面的锁机制代码就会像天书一样。1.1 进程、线程与并行进程是操作系统分配资源的最小单位。可以简单地把进程理解为“一个正在运行的程序实例”。每启动一个 Python 脚本操作系统就帮你创建了一个进程。而线程则隶属于进程是 CPU 调度的最小单位。一个进程内部可以包含多个线程它们共享进程的代码段、数据段和堆空间。为了便于理解我们用生活场景做个类比进程一家正在营业的餐厅。线程餐厅里负责不同工作的服务员。餐厅的厨房、大厅、收银台资源被所有服务员共享。如果一个服务员在传菜时被堵住其他服务员依然可以继续工作。在 Python 多线程中最重要的一个底层概念就是GILGlobal Interpreter Lock全局解释器锁。这是 CPython 解释器为了内存安全而引入的机制。GIL 的存在意味着在同一时刻一个 Python 进程中只有一个线程可以执行 Python 字节码。很多人在了解了 GIL 后会有一个疑问既然同一时刻只能跑一个线程那我们学多线程还有什么用这是一个非常关键的认知误区。GIL 主要限制的是 CPU 密集型任务的并行执行速度。对于存在大量输入输出操作的IO 密集型任务比如网络请求、文件读写、数据库交互在线程陷入等待时GIL 会被释放其他线程就能拿起“锁”继续执行。因此多线程对于 IO 密集型任务有着非常显著的效率提升。1.2 threading 模块解决什么问题threading模块是 Python 标准库的一部分无需额外安装。它提供了比原始_thread模块更高级的线程接口主要解决三类问题创建和管理线程通过Thread类实例化线程对象。线程同步提供Lock、RLock、Condition、Semaphore等同步原语解决多个线程争抢共享数据导致的并发安全问题。线程间通信配合queue.Queue实现线程安全的数据交换。1.3 单线程和多线程的直观对比我们先不急着谈复杂的代码看一个最简单的对比。假设我们要写一个下载文件的模拟函数每次下载耗时 1 秒共需要下载 3 个文件。先看单线程模式import time def download(name): print(f开始下载{name}) time.sleep(1) print(f下载完成{name}) start time.time() for i in range(3): download(ffile-{i}.zip) print(f单线程总耗时{time.time() - start:.2f} 秒)如果使用多线程来执行同样的任务import threading import time def download(name): print(f开始下载{name}) time.sleep(1) print(f下载完成{name}) start time.time() threads [] for i in range(3): t threading.Thread(targetdownload, args(ffile-{i}.zip,)) t.start() threads.append(t) for t in threads: t.join() print(f多线程总耗时{time.time() - start:.2f} 秒)在 IO 等待型任务中多线程版本的耗时接近 1 秒而单线程版本需要 3 秒。这就是我们的直观收益。2. 环境准备与版本说明threading模块属于 Python 标准库因此环境准备相对简单。2.1 Python 环境建议使用 Python 3.8 及以上的版本虽然threading在 Python 2.x 时代就已经存在但新版本在语法和功能上更友好。本文示例基于 Python 3.10 编写代码在 Python 3.8 的环境下均可正常运行。如果你还没有安装 Python可以前往 Python 官方网站选择对应操作系统的安装包。安装时记得勾选Add Python to PATH这能避免后续命令行输入python不生效的问题。查看版本信息python --version2.2 IDE 选择Visual Studio Code 和 PyCharm 是首选。VS Code 需要在扩展市场安装 Python 扩展PyCharm 则开箱即用。无论选择哪个工具只需要能创建.py文件并成功运行即可不涉及特殊插件依赖。在写多线程代码时建议开启 IDE 的“运行当前文件”能力方便通过多次执行来观察不同线程调度顺序的变化。2.3 关于示例环境为了让例子避免和你本机环境发生冲突所有代码以独立文件的方式提供。你只需要创建一个名为thread_demo的文件夹把每一节的代码存成独立的.py文件即可。3. threading 模块核心知识拆解这一节会逐一拆解threading模块的常用 API。不管你是初学者还是老手都建议从头看一遍因为很多隐蔽的问题都出在 API 使用不恰当上。3.1 Thread 类创建线程最常用的方式是实例化threading.Thread。关键参数target线程要执行的函数名。args以元组形式传入目标函数的参数。name当前线程的自定义名称方便日志定位。daemon是否为守护线程。先看基础的多线程写法import threading import time def job(name, delay): time.sleep(delay) print(f线程 {name} 结束了当前活动线程数{threading.active_count()}) if __name__ __main__: print(主线程开始) start time.time() t1 threading.Thread(targetjob, args(A, 1)) t2 threading.Thread(targetjob, args(B, 2)) t1.start() t2.start() # join 方法会阻塞主线程直到子线程结束 t1.join() t2.join() print(主线程结束总耗时, time.time() - start)运行这段程序你会发现主线程会一直等待t1和t2都执行完才结束。如果不调用join()主线程可能先执行完并退出但子线程并不会立刻被销毁而是继续在后台运行。3.2 自定义线程类除了通过函数方式创建线程我们还可以继承Thread类并重写run方法。import threading import time class DownloadThread(threading.Thread): def __init__(self, name, file_url): super().__init__(namename) self.file_url file_url def run(self): print(f{self.name} 正在处理 {self.file_url}) time.sleep(2) print(f{self.name} 处理完成) if __name__ __main__: t1 DownloadThread(线程1, https://example.com/a.zip) t2 DownloadThread(线程2, https://example.com/b.zip) t1.start() t2.start() t1.join() t2.join() print(全部任务结束)需要强调直接调用run()方法并不会创建新线程。run()只是一个普通方法调用在当前主线程里顺序执行。只有调用start()Python 才会真正创建系统级线程并自动执行run()方法。3.3 守护线程与主线程的关系守护线程也叫后台线程当主线程结束时守护线程会被强制终止不管它执行到哪个位置。import threading import time def watch(): while True: print(守护线程正在监控...) time.sleep(0.5) if __name__ __main__: d threading.Thread(targetwatch, daemonTrue) d.start() time.sleep(2) print(主线程结束守护线程随之退出)如果不设置daemonTrue像上述while True循环程序将永远无法正常退出。这在写后台监控任务时是极其有用的特性。3.4 线程同步的必要性多线程的核心难点在于当多个线程同时修改同一个全局变量时会出现数据竞争问题。看下面这个经典例子import threading total 0 def add(): global total for _ in range(100000): total 1 def sub(): global total for _ in range(100000): total - 1 if __name__ __main__: t1 threading.Thread(targetadd) t2 threading.Thread(targetsub) t1.start() t2.start() t1.join() t2.join() print(f最终 total 值{total}预期结果0)多次运行你会发现total的最终结果往往不是 0而是一个随机数。原因在于total 1这一行字节码是由多个底层指令组成的CPU 在执行线程 A 的中间状态时可能切到线程 B进而导致数据不一致。解决这个问题需要引入锁。3.5 Lock互斥锁Lock是最基础的同步原语它有两个核心方法acquire()获取锁release()释放锁。使用锁后同一时刻只有一个线程能够修改共享数据。import threading total 0 lock threading.Lock() def add(): global total for _ in range(100000): lock.acquire() total 1 lock.release() def sub(): global total for _ in range(100000): with lock: total - 1 if __name__ __main__: t1 threading.Thread(targetadd) t2 threading.Thread(targetsub) t1.start() t2.start() t1.join() t2.join() print(f最终 total 值{total})推荐使用with lock:语法代替手写acquire和release因为即使在代码抛出异常时锁也能被自动释放避免死锁风险。3.6 RLock可重入锁Lock有一个限制同一个线程无法多次acquire同一个普通锁否则会进入死锁状态。考虑递归调用的场景就需要使用RLock。import threading rlock threading.RLock() def factorial(n): with rlock: if n 1: return 1 return n * factorial(n - 1) if __name__ __main__: print(factorial(5))RLock允许同一个线程多次获取锁每次获取都要对应一次释放。更稳妥的选择是直接用RLock替代大部分普通Lock场景。3.7 Semaphore信号量Semaphore用于控制同时访问某个资源的线程数量。比如限制数据库连接池最多只有 5 个连接import threading import time semaphore threading.Semaphore(3) def worker(name): with semaphore: print(f线程 {name} 进入) time.sleep(1) print(f线程 {name} 离开) if __name__ __main__: threads [threading.Thread(targetworker, args(i,)) for i in range(6)] for t in threads: t.start() for t in threads: t.join()运行时会发现每次最多只有 3 个线程同时处于“进入”状态。3.8 Event事件通知Event用于线程间的简单通信。一个线程等待某个事件另一个线程设置事件来唤醒等待者。import threading import time event threading.Event() def waiter(): print(等待事件...) event.wait() print(收到事件继续执行) def setter(): time.sleep(2) print(准备触发事件) event.set() if __name__ __main__: threading.Thread(targetwaiter).start() threading.Thread(targetsetter).start()这是实现“启动信号”或“平滑停机”常用的模式。3.9 Queue线程安全队列虽然Event能做通知但在多个线程之间传递数据最佳实践是使用queue.Queue。import queue import threading import time q queue.Queue() def producer(): for i in range(5): q.put(f数据{i}) time.sleep(0.2) q.put(None) # 结束信号 def consumer(): while True: item q.get() if item is None: break print(f正在处理{item}) q.task_done() if __name__ __main__: threading.Thread(targetproducer).start() threading.Thread(targetconsumer).start()Queue内部本身就加了锁因此多线程环境下读写是线程安全的。4. 完整实战案例多线程网页内容采集器前面的知识点比较分散这一节我们通过一个贴近真实开发的案例把Thread、Queue、锁和join完整串联起来。需求如下给定一批 URL 列表模拟 HTTP 请求并采集页面标题最终将结果写入一个共享的 CSV 文件。为了避免对单条请求等待过久需要并发处理为了防止打印和写入文件时数据错乱需要加锁。4.1 创建项目结构thread_spider/ ├── spider.py └── requirements.txt本文不使用第三方库因此requirements.txt可以留空。4.2 编写核心代码# spider.py import csv import random import threading import time from queue import Queue # 全局结果列表会被多个线程同时写入 results [] # 用于保护 results 写入过程的锁 file_lock threading.Lock() # 线程安全任务队列 task_queue Queue() def fetch_title(url): 模拟从网页中获取标题。 实际项目中通常使用 requests 等库发起网络请求。 time.sleep(random.uniform(0.1, 0.5)) return f标题来自{url} def worker(worker_id): 消费者线程从任务队列取 URL抓取标题并保存结果。 while True: try: # 队列为空时会抛出 Empty 异常 url task_queue.get(timeout1) except Exception: break print(f工作线程 [{worker_id}] 正在处理{url}) # 模拟网络异常概率 try: title fetch_title(url) except Exception as e: title f抓取失败{e} # 写入共享列表时加锁防止多线程同时 append 导致数据异常 with file_lock: results.append((url, title)) # 标记该任务已完成 task_queue.task_done() if __name__ __main__: # 准备 URL 任务 urls [fhttps://example.com/page/{i} for i in range(20)] for url in urls: task_queue.put(url) # 启动 5 个消费者线程 thread_pool [] for wid in range(5): t threading.Thread(targetworker, args(wid,), daemonTrue) t.start() thread_pool.append(t) # 主线程等待队列中所有任务被消费完成 task_queue.join() print(全部任务消费完成准备将结果写入 CSV) # 将结果写入 CSV 文件 with open(output.csv, w, newline, encodingutf-8) as f: writer csv.writer(f) writer.writerow([URL, TITLE]) writer.writerows(results) print(f已完成共采集 {len(results)} 个 URL)4.3 运行与验证在终端中执行cd thread_spider python spider.py预期会在当前目录生成output.csv文件内容类似URL,TITLE https://example.com/page/0,标题来自https://example.com/page/0 https://example.com/page/1,标题来自https://example.com/page/14.4 设计说明为什么使用task_queue.join()而不是直接让主线程sleep因为join()会阻塞主线程直到所有任务都被task_done()标记为完成这是一种优雅的线程协同等待机制。如果简单地轮询队列的empty()状态在多线程场景下容易出现“队列没来得及填满消费者就终止”的问题。5. 多线程写文件订单汇总实战上面的案例侧重“任务分发”和“结果聚合”。这一节我们换一个更贴近业务的视角多个流水线线程把各自的处理结果分批次写入同一个日志文件。实际业务中多个线程把日志写入同一个文件时容易出现内容穿插错行。通过自定义文件锁和行缓冲可以保证每一行日志都是完整的。import threading import time class OrderProcessor: def __init__(self, log_path): self.log_path log_path self.log_lock threading.Lock() def write_log(self, message): with self.log_lock: with open(self.log_path, a, encodingutf-8) as f: f.write(f{time.strftime(%H:%M:%S)} - {message}\n) def process_order(self, order_id): for step in range(1, 4): time.sleep(0.1) self.write_log(f订单 {order_id} 执行到第 {step} 步) self.write_log(f订单 {order_id} 处理完成) if __name__ __main__: processor OrderProcessor(orders.log) threads [] for order_id in range(1, 6): t threading.Thread(targetprocessor.process_order, args(order_id,)) threads.append(t) t.start() for t in threads: t.join() print(所有订单处理结束)由于每次写入都使用了独立的with open以及全局锁因此文件不会出现半行错乱。虽然频繁开关文件有效率损耗但在日志量可控的场景中这种写法更安全可靠不用考虑文件指针偏移的边界问题。6. 线程池concurrent.futures 快速上手从 Python 3.2 开始标准库新增了concurrent.futures。虽然本文核心是threading但在真实项目中直接手动管理线程往往不够优雅。线程池支持复用线程、限制最大并发数、获取返回值以及捕获异常是进阶的必备技能。import concurrent.futures import time def fetch_data(api_id): time.sleep(0.3) if api_id 3: raise ValueError(id3 请求失败) return fAPI-{api_id} 返回数据 if __name__ __main__: with concurrent.futures.ThreadPoolExecutor(max_workers3) as executor: future_to_id {executor.submit(fetch_data, i): i for i in range(6)} for future in concurrent.futures.as_completed(future_to_id): api_id future_to_id[future] try: result future.result() print(f{api_id} - {result}) except Exception as exc: print(f{api_id} - 发生异常{exc})submit方法返回一个Future对象它代表一个可查询的异步任务。as_completed会按任务完成顺序返回迭代结果这在处理耗时差异大的任务时非常高效。7. 常见问题与排查思路多线程代码的故障往往不像普通语法错误那样一眼可见。下面整理几个高频问题。问题现象常见原因解决思路变量值神秘变化和预期不符多个线程同时修改全局共享变量产生数据竞争用Lock包裹修改操作优先考虑使用Queue传递数据程序运行到最后无法退出存在非守护线程正在执行死循环或有acquire未释放检查所有的Thread是否设置daemonTrue使用with lock语法调用join()后仍然出现主线程先退出在多线程分析或调试工具中join()挂在已启动线程变量上但线程引用的函数内部又有新线程检查线程内部是否又启动了子线程确保所有子线程也被正确 join打印的日志顺序错乱多个线程竞争标准输出print并不是原子操作使用自定义日志函数加锁或者通过Queue汇集后统一写入打印使用了time.sleep还是频繁切换任务线程切换由操作系统决定sleep只是主动让出 CPU 的一种方式不要依赖sleep做线程同步改用Event、Lock或Condition遇到RuntimeError: threads can only be started once对同一个Thread对象再次调用了start()一个线程对象只能启动一次需要再次执行就不重新创建线程对象高并发下通过 requests 请求外部接口频繁超时单线程并发数过高触发了服务端的限流用Semaphore限制并发数并增加重试机制排查多线程问题建议按下面的静态清单来复现问题时先把线程数降到 1如果问题消失说明与并发密切相关。在关键变量修改处加上临时print观察线程名与时间戳。检查是否存在多个线程共同修改同一个 Python 列表或字典且未加锁。尝试把Lock换成RLock看是否消除递归死锁问题。如果代码中大量使用while True消费队列确认队列信号None是否被正确放置防止消费者永久阻塞。8. 最佳实践与工程建议在项目中使用 Python 多线程有几个工程层面的建议值得收藏。8.1 分清适用场景threading不是万能的。IO 密集型任务网络爬虫、API 调用、文件读写推荐使用多线程能明显缩短等待时间。CPU 密集型任务大量数值计算、图像像素遍历、复杂的加密解密多线程受 GIL 限制通常不如多进程。建议考虑multiprocessing或直接调用concurrent.futures.ProcessPoolExecutor。8.2 优先使用现成的高层封装手动管理线程的创建和销毁非常容易出错。如果你的任务只是“让一组函数并发执行并收集返回值”应优先使用from concurrent.futures import ThreadPoolExecutor这样既避免了每次创建新线程的开销也能通过max_workers控制并发上限避免对下游服务造成压力。8.3 共享资源访问最小化多线程程序最难以调试的位置就是全局变量。尽量避免编写多个线程直接修改同一个全局列表。推荐模式是生产者把任务放入Queue。消费者从Queue中取任务处理。处理结果再次放入结果Queue。由主线程统一从结果队列中取数据再写入文件或数据库。8.4 锁的粒度锁太多会损失性能甚至产生死锁锁太少则会出现线程安全问题。实际工程中需要在两者之间做平衡。一个经验是锁的粒度应该控制在临界区代码范围内。尽量不要把整个函数体用一个大锁包住那样本质上就是串行执行失去了多线程加速的效果。正确的思路是把共享变量的一次“读-改-写”操作作为一个不可分割的原子操作只锁住这一部分。# 不推荐锁住范围过大 def bad_worker(data): with lock: time.sleep(0.1) # 睡眠并不需要持有锁 parse_data(data) save_result(data) # 推荐锁住真正需要保护的共享写区域 def good_worker(data): time.sleep(0.1) # IO 等待不占锁 result parse_data(data) with lock: save_result_to_shared_file(result)8.5 异常捕获要足够细致线程中一旦抛出未捕获的异常整个线程会退出但程序主流程却不一定感知得到这是最容易隐藏 Bug 的地方。在线程入口函数内部使用try/except记录日志并把异常信息放入队列或日志文件是非常推荐的做法。import logging logging.basicConfig(levellogging.INFO, format%(asctime)s - %(threadName)s - %(message)s) def safe_worker(url): try: resp requests.get(url, timeout3) logging.info(请求成功 %s状态码 %s, url, resp.status_code) except Exception as exc: logging.exception(请求失败 %s原因 %s, url, exc)8.6 程序退出前的清理顺序在大型项目中主线程退出前需要完成几件确定的事停止接收新任务。等待队列中已有任务处理完。转移或写入最终缓冲结果。关闭数据库连接和文件句柄。通知守护线程退出。使用Event来发出“退出通知”比直接调用os._exit安全得多。8.7 性能测量与对比不要凭直觉判断多线程是否真的更快。建议使用装饰器记录函数的耗时在同等条件下对比单线程、多线程、协程的执行结果。只有基于数据说话才能找到最适合项目的方案。import time from functools import wraps def timeit(func): wraps(func) def wrapper(*args, **kwargs): start time.perf_counter() result func(*args, **kwargs) print(f{func.__name__} 耗时 {time.perf_counter() - start:.3f} 秒) return result return wrapper9. 小结与下一步学习方向到这里Pythonthreading模块的核心内容算是走完了一遍。我们从进程线程的基本概念出发分析了 GIL 对多线程的影响动手创建了线程对象梳理了锁、信号量、队列等同步机制最后通过实战案例完整展示了“线程池 任务队列 共享结果聚合”的开发模式。虽然 Python 多线程并不是万能的性能银弹但threading是必知必会的并发编程底层工具理解它能帮助你后面理解协程、理解多进程以及分布式任务系统。如果你接下来要继续深入并行编程可以从这几个方向依次推进熟悉concurrent.futures中线程池的完整 API练习批量并发请求场景。学习asyncio协程体会“单线程异步 IO”在超高并发场景中的优势。学习multiprocessing多进程解决 CPU 密集型任务需要利用多核能力的问题。阅读 Python 官方文档中关于 GIL 的 FAQ了解它是如何被释放和重新获取的。多线程相关的面试题无论是在校招还是社招中出现的频率都极高。下一次面试官再问到“Python 多线程能利用多核吗”的时候希望这篇文章里的知识点能帮到你而你也能基于实测数据说出多线程 IO 密集型任务到底快在哪里。真正理解多线程的程序员写出来的代码会更稳健定位并发问题的效率也会更高。
返回列表