ARTICLE DETAIL

资讯详情

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

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

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(targetdownload_file, args(movie.mp4,), name下载线程-1) thread2 threading.Thread(targetdownload_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(timeout0.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__: # 创建守护线程daemonTrue logger_thread threading.Thread(targetbackground_logger, daemonTrue) 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(targetworker, 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(targetincrement) 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(targetincrement_with_lock) threads.append(t) t.start() for t in threads: t.join() print(f使用锁后的结果: {counter}) # 正确输出 1000000使用锁的最佳实践使用with语句推荐Lock对象支持上下文管理器协议使用with可以自动获取和释放锁代码更简洁且能确保异常发生时锁也能被释放。def increment_with_lock_better(): global counter for _ in range(100000): with counter_lock: # 自动获取和释放锁 counter 1锁的粒度要适中锁保护的范围临界区越小越好只包含真正需要互斥访问的代码。锁的粒度过大会严重降低并发性能因为其他线程需要等待更长时间。避免嵌套锁与死锁如果一个线程在持有锁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: # blockTrue, timeout1 表示最多阻塞1秒等待物品 item q.get(blockTrue, timeout1) 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(maxsize3) # 设置队列最大容量为3 # 创建2个生产者3个消费者 producers [threading.Thread(targetproducer, args(task_queue, i)) for i in range(2)] consumers [threading.Thread(targetconsumer, 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, blockTrue, timeoutNone): 放入项目。如果队列满blockTrue时会阻塞直到有空位timeout设置阻塞超时时间。get(blockTrue, timeoutNone): 取出项目。如果队列空blockTrue时会阻塞直到有项目timeout设置阻塞超时时间。task_done(): 消费者处理完一个从get()得到的项目后调用。用于通知队列该任务已完成。join(): 阻塞调用者直到队列中所有项目都被处理即每个put()进来的项目都对应调用了task_done()。这在协调生产者和消费者结束时非常有用。使用Queue极大地简化了线程间通信的复杂度你不再需要关心底层的锁和条件变量只需关注业务逻辑。3.3 条件变量与事件更复杂的线程协调当线程间的协作不仅仅是传递数据还需要等待某个条件成立时就需要用到threading.Condition或threading.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(targetresource_loader) workers [threading.Thread(targetworker, 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(maxlencapacity) 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多线程一个无法回避的话题就是GILGlobal 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(targetcpu_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(processesnum_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.futures的ThreadPoolExecutor在某些I/O操作中在进行耗时运算时会主动释放GIL从而允许其他线程运行。如果你在编写C扩展也可以使用Py_BEGIN_ALLOW_THREADS和Py_END_ALLOW_THREADS宏来临时释放GIL。方案三使用concurrent.futures模块这个高级模块提供了ThreadPoolExecutor和ProcessPoolExecutor它们提供了统一的接口来执行并发任务底层自动选择线程池或进程池。对于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_workersnum_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密集型应用优先考虑多进程multiprocessing或ProcessPoolExecutor或者直接使用专为科学计算设计的库如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, timeout5) 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_workers3): 使用线程池并发下载 results [] with ThreadPoolExecutor(max_workersmax_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_workers3): with ThreadPoolExecutor(max_workersmax_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_workers3) print(f耗时: {time.time() - start:.2f}秒)2. 使用multiprocessing.pool.ThreadPoolmultiprocessing模块也提供了一个线程池其接口与进程池类似。from multiprocessing.pool import ThreadPool def task(x): return x * x with ThreadPool(processes4) 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(targetworker, args(i,), namefThread-{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(targetbuggy_worker) t.start() t.join() print(主线程结束) # 程序会正常结束你看不到错误信息解决方案在线程函数内部用try...except捕获所有异常并记录日志。import traceback import logging logging.basicConfig(levellogging.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)解决方案避免嵌套锁重新设计代码逻辑尽量减少需要同时持有多个锁的情况。固定锁的获取顺序如果必须获取多个锁确保所有线程都以相同的顺序获取它们例如总是先获取lock_a再获取lock_b。使用带超时的锁lock.acquire(timeout5)超时后可以执行回退逻辑。使用高级抽象尽可能使用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块文件自动关闭最佳实践总结优先使用高层抽象如concurrent.futures.ThreadPoolExecutor和queue.Queue它们更安全更不易出错。明确线程的职责和生命周期设计时就想好线程何时启动、何时结束、如何优雅停止使用Event信号。所有共享数据都必须同步对任何可能被多个线程修改的数据都要考虑使用锁或线程安全数据结构。保持简单多线程代码本来就复杂尽量让每个线程的逻辑简单、独立。复杂的交互尽量通过队列进行。充分测试多线程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密集型下使用它能极大地提升程序性能。理解其原理掌握threading、queue、concurrent.futures等核心模块并牢记同步、通信、异常处理和资源管理的要点你就能写出高效、健壮的多线程程序。当任务规模继续扩大你自然会接触到asyncio和多进程等更高级的并发模型而扎实的多线程基础将为理解它们铺平道路。
返回列表