一、先搞懂几个基本概念
1.1 进程(Process)是什么?
进程 = 一个正在运行的程序。
比如你同时打开了:
每个程序都有自己独立的内存空间,互不干扰。
1.2 线程(Thread)是什么?
线程 = 进程内部的一条"执行路径"。
一个进程里可以有多条线程同时工作。
生活比喻:
🏠 把进程想象成一家餐厅
- 单线程:餐厅里只有1个服务员,客人来了点菜→做菜→上菜→收钱,全是一个人干,后面的客人只能排队等。
- 多线程:餐厅里有多个服务员,A服务员给1号桌点菜的同时,B服务员在给2号桌上菜,C服务员在收银。大家同时干活,效率大大提高。
1.3 为什么需要多线程?
核心原因:避免"等待"浪费时间。
没有多线程(串行):下载文件A(等3秒)→ 下载文件B(等3秒)→ 下载文件C(等3秒)总耗时:9秒 ❌有多线程(并行):下载文件A(等3秒)下载文件B(等3秒) ← 同时进行!下载文件C(等3秒)总耗时:3秒 ✅
适合多线程的场景(I/O密集型):
不适合多线程的场景(CPU密集型):
这些CPU密集型任务应该用多进程(后面会提到原因)。
二、Python 中的线程模块
Python 提供了两个线程相关的模块:
| |
|---|
_thread | |
threading | 高级模块,推荐使用 |
我们全程使用 threading。
三、创建线程的 3 种方式
方式1:用函数创建(最常用)
import threadingimport time# 定义一个普通函数(这就是线程要执行的任务)def download_file(file_name, seconds): """模拟下载文件""" print(f"⬇️ 开始下载:{file_name}") time.sleep(seconds) # 模拟耗时操作(比如网络等待) print(f"✅ 下载完成:{file_name}")# ============ 主程序 ============print("=== 程序开始 ===")start_time = time.time()# 创建线程对象# target: 线程要执行的函数# args: 传给函数的参数(必须是元组!)t1 = threading.Thread(target=download_file, args=("电影.mp4", 3))t2 = threading.Thread(target=download_file, args=("音乐.mp3", 2))t3 = threading.Thread(target=download_file, args=("文档.pdf", 1))# 启动线程(注意:是 start(),不是 run()!)t1.start()t2.start()t3.start()# 等待所有线程完成(主线程会在这里阻塞)t1.join()t2.join()t3.join()end_time = time.time()print(f"=== 全部完成,总耗时:{end_time - start_time:.2f}秒 ===")
运行结果:
=== 程序开始 ===⬇️ 开始下载:电影.mp4⬇️ 开始下载:音乐.mp3⬇️ 开始下载:文档.pdf✅ 下载完成:文档.pdf✅ 下载完成:音乐.mp3✅ 下载完成:电影.mp4=== 全部完成,总耗时:3.00秒 ===
📌 注意:总耗时是3秒(最慢的那个),而不是 3+2+1=6秒!因为三个下载是同时进行的。
关键方法解释
| | |
|---|
Thread(target=函数, args=参数) | | |
.start() | | |
.join() | | |
.is_alive() | | |
.name | | |
⚠️ 新手常犯错误
# ❌ 错误:直接调用 run(),这不会创建新线程!t1.run() # 这只是普通函数调用,还是串行的!# ✅ 正确:调用 start(),这才会创建新线程t1.start()# ❌ 错误:args 不是元组t1 = threading.Thread(target=download_file, args=("电影.mp4", 3)) # ✅ 正确t1 = threading.Thread(target=download_file, args=("电影.mp4", 3)) # ✅ 正确t1 = threading.Thread(target=download_file, args="电影.mp4") # ❌ 错误!字符串会被拆开# 只有一个参数时,注意逗号!t1 = threading.Thread(target=func, args=("hello",)) # ✅ 注意逗号t1 = threading.Thread(target=func, args=("hello")) # ❌ 这不是元组,是字符串!
方式2:继承 Thread 类
import threadingimport time# 自定义线程类,继承 threading.Threadclass DownloadThread(threading.Thread): def __init__(self, file_name, seconds): super().__init__() # 必须调用父类的 __init__ self.file_name = file_name self.seconds = seconds # 重写 run() 方法(线程启动后会自动执行这个方法) def run(self): print(f"⬇️ [{self.name}] 开始下载:{self.file_name}") time.sleep(self.seconds) print(f"✅ [{self.name}] 下载完成:{self.file_name}")# ============ 使用 ============t1 = DownloadThread("电影.mp4", 3)t2 = DownloadThread("音乐.mp3", 2)# 可以给线程起名字t1.name = "线程-1"t2.name = "线程-2"t1.start()t2.start()t1.join()t2.join()print("全部完成!")
什么时候用继承方式?
- 当你的线程逻辑比较复杂,需要维护状态(实例变量)时
方式3:用 lambda 表达式(适合简单任务)
import threading# 适合非常简单的任务t = threading.Thread(target=lambda: print("Hello from thread!"))t.start()t.join()# 带参数的 lambdat = threading.Thread(target=lambda x, y: print(f"{x} + {y} = {x+y}"), args=(3, 5))t.start()t.join()
四、线程的常用属性和方法
import threadingimport timedef worker(): print(f" 当前线程名:{threading.current_thread().name}") print(f" 当前线程ID:{threading.current_thread().ident}") time.sleep(2)# 创建线程t = threading.Thread(target=worker, name="我的工作线程")print(f"线程是否存活(启动前):{t.is_alive()}") # Falset.start()print(f"线程是否存活(启动后):{t.is_alive()}") # Truet.join()print(f"线程是否存活(结束后):{t.is_alive()}") # False# 查看当前有多少个线程在运行print(f"当前活跃线程数:{threading.active_count()}")# 列出所有活跃线程for thread in threading.enumerate(): print(f" - {thread.name}")
输出:
线程是否存活(启动前):False 当前线程名:我的工作线程 当前线程ID:140234567890线程是否存活(启动后):True线程是否存活(结束后):False当前活跃线程数:1 - MainThread
五、守护线程(Daemon Thread)
5.1 什么是守护线程?
守护线程 = 后台线程,当主线程结束时,守护线程会被强制终止。
生活比喻:
主线程 = 你(主人) 守护线程 = 你家的扫地机器人
你出门了(主线程结束),扫地机器人自动停止(守护线程被杀死)。 非守护线程 = 你请的保姆,你出门了她还会继续打扫完才走。
5.2 代码示例
import threadingimport timedef background_task(): """后台任务:每隔1秒打印一次""" while True: print(" 🤖 后台监控中...") time.sleep(1)def normal_task(): """普通任务""" time.sleep(3) print(" ✅ 普通任务完成")# 创建守护线程daemon_thread = threading.Thread(target=background_task, daemon=True)# 或者:daemon_thread.daemon = True(必须在start之前设置)# 创建普通线程normal_thread = threading.Thread(target=normal_task)daemon_thread.start()normal_thread.start()# 主线程只等普通线程normal_thread.join()print("主线程结束!守护线程会被自动杀死。")# 程序在这里退出,daemon_thread 被强制终止
输出:
🤖 后台监控中... 🤖 后台监控中... 🤖 后台监控中... ✅ 普通任务完成主线程结束!守护线程会被自动杀死。
注意:后台监控只打印了3次就被杀了,不会无限打印下去。
5.3 守护线程的典型用途
六、线程同步 —— 锁(Lock)
6.1 为什么需要锁?
问题场景:多个线程同时修改同一个变量,会出错!
import threading# 共享变量balance = 100 # 账户余额def withdraw(amount): """取钱""" global balance # 模拟操作需要时间(比如网络延迟) import time time.sleep(0.1) if balance >= amount: balance -= amount print(f" 取出 {amount} 元,余额:{balance}") else: print(f" ❌ 余额不足!当前余额:{balance}")# 两个线程同时取钱t1 = threading.Thread(target=withdraw, args=(80,))t2 = threading.Thread(target=withdraw, args=(80,))t1.start()t2.start()t1.join()t2.join()print(f"最终余额:{balance}")
可能的错误输出:
取出 80 元,余额:20 取出 80 元,余额:20 ← 错了!应该是余额不足!最终余额:20 ← 应该是 -60 或者只取一次
为什么会错?
时间线:线程1:读取 balance=100 → 判断 100>=80 ✓ → (等待0.1秒)→ balance = 100-80 = 20线程2:读取 balance=100 → 判断 100>=80 ✓ → (等待0.1秒)→ balance = 100-80 = 20 ↑ 两个线程同时读到了100!都以为够取!
这就是竞态条件(Race Condition)。
6.2 用锁解决问题
import threadingimport timebalance = 100lock = threading.Lock() # 创建一把锁 🔒def withdraw(amount): global balance # 获取锁(如果锁被别人拿着,就在这里等待) lock.acquire() try: # 这段代码同一时间只有一个线程能执行 time.sleep(0.1) # 模拟耗时操作 if balance >= amount: balance -= amount print(f" ✅ 取出 {amount} 元,余额:{balance}") else: print(f" ❌ 余额不足!当前余额:{balance}") finally: # 无论是否出错,都必须释放锁! lock.release()t1 = threading.Thread(target=withdraw, args=(80,))t2 = threading.Thread(target=withdraw, args=(80,))t1.start()t2.start()t1.join()t2.join()print(f"最终余额:{balance}")
正确输出:
✅ 取出 80 元,余额:20 ❌ 余额不足!当前余额:20最终余额:20
6.3 更优雅的写法:with 语句
import threadingimport timebalance = 100lock = threading.Lock()def withdraw(amount): global balance # with 语句会自动 acquire 和 release,即使出错也会释放! with lock: time.sleep(0.1) if balance >= amount: balance -= amount print(f" ✅ 取出 {amount} 元,余额:{balance}") else: print(f" ❌ 余额不足!当前余额:{balance}")# 使用方式不变t1 = threading.Thread(target=withdraw, args=(80,))t2 = threading.Thread(target=withdraw, args=(80,))t1.start()t2.start()t1.join()t2.join()print(f"最终余额:{balance}")
📌 强烈建议:永远用 with lock: 而不是手动 acquire()/release(),避免忘记释放锁导致死锁。
6.4 锁的工作原理图
没有锁:线程1 ──→ [读balance] ──→ [判断] ──→ [修改balance]线程2 ──→ [读balance] ──→ [判断] ──→ [修改balance] ← 同时执行,冲突!有锁:线程1 ──→ 🔒获取锁 ──→ [读] ──→ [判断] ──→ [修改] ──→ 🔓释放锁线程2 ──→ ⏳等待...等待...等待... ──→ 🔒获取锁 ──→ [读] ──→ [判断] ──→ [修改] ──→ 🔓释放锁
七、其他同步机制
7.1 RLock(可重入锁)
普通锁的问题:同一个线程不能重复获取同一把锁(会死锁)。
import threadinglock = threading.Lock()def func_a(): with lock: print("func_a 获取了锁") func_b() # 在锁里面调用另一个也需要锁的函数def func_b(): with lock: # ❌ 死锁!因为锁已经被 func_a 拿着了 print("func_b 获取了锁")# 解决:用 RLock(可重入锁)rlock = threading.RLock()def func_a(): with rlock: print("func_a 获取了锁") func_b() # ✅ 没问题!RLock 允许同一线程重复获取def func_b(): with rlock: # ✅ 同一线程可以再次获取 print("func_b 获取了锁")
7.2 Semaphore(信号量)
作用:限制同时访问某个资源的线程数量。
生活比喻:
停车场只有3个车位(信号量=3),最多同时停3辆车。第4辆车必须在外面等。
import threadingimport time# 最多允许3个线程同时执行semaphore = threading.Semaphore(3)def access_resource(thread_id): with semaphore: print(f" 🚗 线程{thread_id} 进入(剩余车位:{semaphore._value})") time.sleep(2) # 模拟占用资源 print(f" 🚙 线程{thread_id} 离开")# 创建6个线程threads = []for i in range(6): t = threading.Thread(target=access_resource, args=(i,)) threads.append(t) t.start()for t in threads: t.join()print("全部完成")
输出(每次最多3个同时进入):
🚗 线程0 进入(剩余车位:2) 🚗 线程1 进入(剩余车位:1) 🚗 线程2 进入(剩余车位:0) 🚙 线程0 离开 🚙 线程1 离开 🚗 线程3 进入(剩余车位:1) 🚗 线程4 进入(剩余车位:0) 🚙 线程2 离开 ...
7.3 Event(事件)
作用:一个线程发信号,其他线程收到信号后继续执行。
生活比喻:
老师(线程A)说"下课了"(set事件),所有学生(线程B、C、D)听到后离开教室。
import threadingimport time# 创建事件对象(初始状态:未触发)event = threading.Event()def teacher(): """老师:3秒后宣布下课""" print("👨🏫 老师:上课中...") time.sleep(3) print("👨🏫 老师:下课了!") event.set() # 触发事件,通知所有等待的线程def student(name): """学生:等待下课""" print(f" 🧑🎓 {name}:等待下课...") event.wait() # 阻塞,直到事件被触发 print(f" 🧑🎓 {name}:太好了,下课!收拾书包走人~")# 启动t_teacher = threading.Thread(target=teacher)t_student1 = threading.Thread(target=student, args=("小明",))t_student2 = threading.Thread(target=student, args=("小红",))t_student1.start()t_student2.start()t_teacher.start()t_teacher.join()t_student1.join()t_student2.join()
输出:
🧑🎓 小明:等待下课... 🧑🎓 小红:等待下课...👨🏫 老师:上课中...👨🏫 老师:下课了! 🧑🎓 小明:太好了,下课!收拾书包走人~ 🧑🎓 小红:太好了,下课!收拾书包走人~
7.4 Condition(条件变量)
比 Event 更灵活,可以等待某个"条件"满足。
import threadingimport timeimport random# 生产者-消费者模型queue = []condition = threading.Condition()def producer(): """生产者:不断生产数据""" for i in range(5): with condition: item = f"商品-{i}" queue.append(item) print(f" 🏭 生产了:{item},队列长度:{len(queue)}") condition.notify() # 通知消费者:有新数据了! time.sleep(random.uniform(0.5, 1.5))def consumer(): """消费者:不断消费数据""" while True: with condition: # 等待条件:队列不为空 while len(queue) == 0: print(" 🛒 队列为空,消费者等待中...") condition.wait() # 等待通知 item = queue.pop(0) print(f" 🛒 消费了:{item},队列长度:{len(queue)}")t1 = threading.Thread(target=producer)t2 = threading.Thread(target=consumer, daemon=True)t2.start()t1.start()t1.join()time.sleep(1) # 等消费者处理完print("完成!")
八、线程间通信 —— Queue(队列)
queue.Queue 是线程安全的队列,是线程间传递数据最方便的方式。
import threadingimport queueimport timeimport random# 创建队列(maxsize=5 表示最多放5个元素)task_queue = queue.Queue(maxsize=5)def producer(name): """生产者:往队列里放任务""" for i in range(5): task = f"{name}-任务{i}" task_queue.put(task) # 放入队列(如果满了会自动等待) print(f" 📤 [{name}] 放入:{task}(队列大小:{task_queue.qsize()})") time.sleep(random.uniform(0.1, 0.5)) print(f" 📤 [{name}] 生产完毕")def consumer(name): """消费者:从队列里取任务""" while True: try: # 从队列取任务(timeout=2:等2秒取不到就退出) task = task_queue.get(timeout=2) print(f" 📥 [{name}] 取出:{task}(队列大小:{task_queue.qsize()})") time.sleep(random.uniform(0.3, 1.0)) # 模拟处理任务 task_queue.task_done() # 标记任务完成 except queue.Empty: print(f" 📥 [{name}] 队列为空,退出") break# 创建线程producers = [ threading.Thread(target=producer, args=(f"生产者{i}",)) for i in range(2)]consumers = [ threading.Thread(target=consumer, args=(f"消费者{i}",)) for i in range(3)]# 启动所有线程for t in producers + consumers: t.start()# 等待所有生产者完成for t in producers: t.join()# 等待队列中所有任务被处理完task_queue.join()print("✅ 所有任务处理完毕!")
Queue 的常用方法
| |
|---|
q.put(item) | |
q.get() | |
q.get(timeout=2) | |
q.qsize() | |
q.empty() | |
q.full() | |
q.task_done() | |
q.join() | |
九、线程池(ThreadPoolExecutor)
9.1 为什么需要线程池?
频繁创建/销毁线程有开销。线程池预先创建好一批线程,重复使用。
生活比喻:
不用线程池 = 每来一个客人就临时招一个服务员,客人走了就辞退 用线程池 = 提前招好10个服务员,有活就干,没活就待命
9.2 基本用法
from concurrent.futures import ThreadPoolExecutor, as_completedimport timedef download(url): """模拟下载""" print(f" ⬇️ 开始下载:{url}") time.sleep(2) # 模拟网络耗时 return f"{url} 的内容({len(url)*10}字节)"# 要下载的URL列表urls = [ "https://example.com/file1", "https://example.com/file2", "https://example.com/file3", "https://example.com/file4", "https://example.com/file5",]# ============ 方式1:submit + as_completed ============print("=== 方式1:逐个获取结果 ===")start = time.time()with ThreadPoolExecutor(max_workers=3) as executor: # 提交所有任务(返回 Future 对象) future_to_url = { executor.submit(download, 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() # 获取返回值 print(f" ✅ 完成:{result}") except Exception as e: print(f" ❌ {url} 出错:{e}")print(f"总耗时:{time.time() - start:.2f}秒")# 5个任务,3个线程,每个2秒 → 大约 4秒(而不是10秒)
# ============ 方式2:map(更简洁) ============print("\n=== 方式2:map 批量处理 ===")start = time.time()with ThreadPoolExecutor(max_workers=3) as executor: # map 会按顺序返回结果(和输入顺序一致) results = executor.map(download, urls) for url, result in zip(urls, results): print(f" ✅ {result}")print(f"总耗时:{time.time() - start:.2f}秒")
9.3 submit vs map 的区别
9.4 实战:多线程爬虫
from concurrent.futures import ThreadPoolExecutor, as_completedimport urllib.requestimport timedef fetch_url(url): """抓取网页""" try: start = time.time() response = urllib.request.urlopen(url, timeout=10) html = response.read().decode("utf-8", errors="ignore") elapsed = time.time() - start return {"url": url, "length": len(html), "time": elapsed, "status": "成功"} except Exception as e: return {"url": url, "error": str(e), "status": "失败"}urls = [ "https://www.baidu.com", "https://www.qq.com", "https://www.taobao.com", "https://www.jd.com", "https://www.zhihu.com", "https://www.bilibili.com", "https://www.douban.com", "https://www.weibo.com",]print(f"要抓取 {len(urls)} 个网页\n")start_time = time.time()# 使用线程池,最多5个线程同时工作with ThreadPoolExecutor(max_workers=5) as executor: futures = {executor.submit(fetch_url, url): url for url in urls} for future in as_completed(futures): result = future.result() if result["status"] == "成功": print(f" ✅ {result['url']}") print(f" 大小:{result['length']} 字节,耗时:{result['time']:.2f}秒") else: print(f" ❌ {result['url']}") print(f" 错误:{result['error']}")total_time = time.time() - start_timeprint(f"\n总耗时:{total_time:.2f}秒")print(f"如果串行执行,大约需要:{sum(r.get('time', 2) for r in [])}秒")
十、GIL —— Python 多线程的"天花板"
10.1 什么是 GIL?
GIL(Global Interpreter Lock,全局解释器锁)是 CPython 解释器的一把全局锁。
规则:同一时刻,只有一个线程能执行 Python 字节码。
10.2 生活比喻
想象一个厨房(CPU)只有一个灶台(GIL):
- 即使你请了10个厨师(10个线程),同一时间也只能有1个人在用灶台炒菜。
- 但是!如果某个厨师在等外卖食材送到(I/O等待),他会让出灶台给别人用。
所以:
- I/O密集型(等网络、等文件):多线程有效 ✅(等待时让出GIL)
- CPU密集型
10.3 验证 GIL 的影响
import threadingimport time# ============ CPU密集型任务 ============def cpu_task(): """纯计算任务""" count = 0 for i in range(10_000_000): count += i return count# 单线程start = time.time()cpu_task()cpu_task()single_time = time.time() - startprint(f"单线程 CPU 任务:{single_time:.2f}秒")# 多线程start = time.time()t1 = threading.Thread(target=cpu_task)t2 = threading.Thread(target=cpu_task)t1.start()t2.start()t1.join()t2.join()multi_time = time.time() - startprint(f"多线程 CPU 任务:{multi_time:.2f}秒")# 结果:多线程不会更快,甚至可能更慢!# 因为 GIL 导致两个线程轮流执行,并没有真正并行
10.4 CPU密集型怎么办?用多进程!
from multiprocessing import Poolimport timedef cpu_task(n): count = 0 for i in range(n): count += i return count# 多进程(绕过GIL,每个进程有自己的解释器)start = time.time()with Pool(processes=2) as pool: results = pool.map(cpu_task, [10_000_000, 10_000_000])print(f"多进程 CPU 任务:{time.time() - start:.2f}秒")# 这次真的会快接近2倍!
10.5 总结:什么时候用什么?
┌─────────────────────────────────────────────────────┐│ 你的任务是什么类型? │├─────────────────────────────────────────────────────┤│ ││ I/O密集型(网络请求、文件读写、数据库) ││ → 用多线程 threading / ThreadPoolExecutor ✅ ││ ││ CPU密集型(数学计算、图像处理、加密) ││ → 用多进程 multiprocessing / ProcessPoolExecutor ✅││ ││ 两者都有 ││ → 多进程 + 每个进程内多线程 ││ │└─────────────────────────────────────────────────────┘
十一、Timer(定时器线程)
import threadingdef say_hello(): print("👋 你好!3秒到了!")# 3秒后执行(只执行一次)timer = threading.Timer(3.0, say_hello)timer.start()print("定时器已设置,等待中...")# 主线程继续做其他事...# 如果想取消定时器(在触发之前):# timer.cancel()
重复定时器:
import threadingdef repeat_task(): print("⏰ 定时任务执行!") # 每次执行完后再创建新的定时器(实现重复) timer = threading.Timer(2.0, repeat_task) timer.daemon = True # 设为守护线程,主程序退出时自动停止 timer.start()# 启动timer = threading.Timer(2.0, repeat_task)timer.daemon = Truetimer.start()# 主线程做其他事import timetime.sleep(10)print("主程序结束")
十二、线程局部变量(threading.local)
问题
多个线程共享全局变量会冲突,但有时候每个线程需要自己的"私有数据"。
import threading# 创建线程局部存储local_data = threading.local()def worker(name): # 每个线程设置自己的值(互不影响) local_data.username = name local_data.count = 0 for i in range(3): local_data.count += 1 print(f" [{name}] count = {local_data.count}") print(f" [{name}] 最终 count = {local_data.count}")t1 = threading.Thread(target=worker, args=("线程A",))t2 = threading.Thread(target=worker, args=("线程B",))t1.start()t2.start()t1.join()t2.join()
输出(每个线程的 count 独立计数):
[线程A] count = 1 [线程B] count = 1 [线程A] count = 2 [线程B] count = 2 [线程A] count = 3 [线程B] count = 3 [线程A] 最终 count = 3 [线程B] 最终 count = 3
如果用普通全局变量,两个线程会互相覆盖。threading.local() 让每个线程有独立的副本。
十三、死锁(Deadlock)
13.1 什么是死锁?
两个线程互相等待对方释放锁,永远等下去。
生活比喻:
两个人在窄路上相遇:
13.2 死锁示例
import threadingimport timelock_a = threading.Lock()lock_b = threading.Lock()def thread_1(): with lock_a: print("线程1:拿到了锁A,想要锁B...") time.sleep(1) # 给线程2时间去拿锁B with lock_b: # 等待锁B(但线程2拿着) print("线程1:拿到了锁B")def thread_2(): with lock_b: print("线程2:拿到了锁B,想要锁A...") time.sleep(1) # 给线程1时间去拿锁A with lock_a: # 等待锁A(但线程1拿着) print("线程2:拿到了锁A")t1 = threading.Thread(target=thread_1)t2 = threading.Thread(target=thread_2)t1.start()t2.start()# 程序会永远卡住!
13.3 如何避免死锁?
# 方法1:总是按相同顺序获取锁# 规定:必须先获取 lock_a,再获取 lock_bdef thread_1(): with lock_a: with lock_b: print("线程1完成")def thread_2(): with lock_a: # 也是先A后B(顺序一致!) with lock_b: print("线程2完成")# 方法2:使用超时if lock_b.acquire(timeout=5): # 最多等5秒 try: # 做事 pass finally: lock_b.release()else: print("获取锁超时,放弃")
十四、实战案例汇总
14.1 多线程文件下载器
import threadingimport urllib.requestimport osimport timedef download_file(url, save_dir, results, index): """下载单个文件""" try: filename = url.split("/")[-1] or f"file_{index}" filepath = os.path.join(save_dir, filename) urllib.request.urlretrieve(url, filepath) size = os.path.getsize(filepath) results[index] = {"file": filename, "size": size, "status": "成功"} print(f" ✅ {filename}({size/1024:.1f} KB)") except Exception as e: results[index] = {"file": url, "error": str(e), "status": "失败"} print(f" ❌ {url}:{e}")def batch_download(urls, save_dir="./downloads", max_threads=5): """批量下载""" os.makedirs(save_dir, exist_ok=True) results = [None] * len(urls) threads = [] print(f"开始下载 {len(urls)} 个文件(最多{max_threads}个线程)\n") start = time.time() # 用信号量控制并发数 semaphore = threading.Semaphore(max_threads) def limited_download(url, idx): with semaphore: download_file(url, save_dir, results, idx) for i, url in enumerate(urls): t = threading.Thread(target=limited_download, args=(url, i)) threads.append(t) t.start() for t in threads: t.join() elapsed = time.time() - start success = sum(1 for r in results if r and r["status"] == "成功") print(f"\n完成!成功 {success}/{len(urls)},耗时 {elapsed:.2f}秒")# 使用urls = [ "https://www.python.org/static/img/python-logo.png", "https://www.baidu.com/img/PCtm_d9c8750bed0b3c7d089fa7d55720d6cf.png",]batch_download(urls)
14.2 多线程进度条
import threadingimport timeimport sysdef progress_bar(total, interval=0.1): """显示进度条(在子线程中运行)""" for i in range(total + 1): percent = i / total * 100 bar_length = 40 filled = int(bar_length * i / total) bar = "█" * filled + "░" * (bar_length - filled) sys.stdout.write(f"\r [{bar}] {percent:.0f}% ({i}/{total})") sys.stdout.flush() time.sleep(interval) print() # 换行def do_work(): """模拟耗时工作""" time.sleep(5) print(" 工作完成!")# 主线程做工作,子线程显示进度progress_thread = threading.Thread(target=progress_bar, args=(50, 0.1))progress_thread.start()do_work()progress_thread.join()
14.3 多线程 + 日志记录
import threadingimport loggingimport timeimport random# 配置日志logging.basicConfig( level=logging.INFO, format="%(asctime)s [%(threadName)s] %(message)s", datefmt="%H:%M:%S")logger = logging.getLogger(__name__)def task(task_id): logger.info(f"任务{task_id} 开始") time.sleep(random.uniform(1, 3)) # 模拟随机失败 if random.random() < 0.3: logger.error(f"任务{task_id} 失败!") raise Exception(f"任务{task_id} 出错了") logger.info(f"任务{task_id} 成功完成") return f"结果-{task_id}"# 运行threads = []for i in range(5): t = threading.Thread(target=task, args=(i,), name=f"Worker-{i}") threads.append(t) t.start()for t in threads: t.join()logger.info("所有任务结束")
十五、最佳实践总结
✅ 应该做的
# 1. 使用 with 语句管理锁with lock: # 操作共享资源 pass# 2. 使用线程池而不是手动创建大量线程from concurrent.futures import ThreadPoolExecutorwith ThreadPoolExecutor(max_workers=10) as executor: results = list(executor.map(func, data_list))# 3. 给线程起有意义的名字t = threading.Thread(target=func, name="数据下载线程")# 4. 用 Queue 做线程间通信(而不是共享变量+锁)import queueq = queue.Queue()# 5. 设置合理的超时lock.acquire(timeout=5)q.get(timeout=3)# 6. 守护线程用于后台任务t.daemon = True
❌ 不应该做的
# 1. 不要创建过多线程(几百个线程会拖慢系统)# 一般 I/O 密集型:10~100 个足够# CPU 密集型:用多进程,线程数 = CPU核心数# 2. 不要在锁里面做 I/O 操作(会阻塞其他线程)with lock: time.sleep(10) # ❌ 其他线程要等10秒!# 3. 不要忘记 join()(否则主线程可能提前退出)t.start()# 忘记 t.join() → 程序可能在子线程完成前就退出了# 4. 不要在线程里使用 print 做关键逻辑(print不是线程安全的)# 用 logging 或 Queue 代替# 5. 不要嵌套太多锁(容易死锁)
十六、速查表
创建线程: threading.Thread(target=func, args=(...), name="名字")启动线程: t.start()等待线程: t.join()守护线程: t.daemon = True(start之前设置)线程锁: lock = threading.Lock() → with lock:可重入锁: lock = threading.RLock()信号量: sem = threading.Semaphore(3)事件: event = threading.Event() → event.wait() / event.set()队列: q = queue.Queue(maxsize=10)线程池: ThreadPoolExecutor(max_workers=5)定时器: threading.Timer(秒数, 函数)线程局部: local = threading.local()当前线程: threading.current_thread()活跃线程数: threading.active_count()