Python多线程编程
发布于
多进程解决了”算”的问题,这一章解决”等”的问题。
上一章用多进程跑满多核,但每个进程独立内存、通信靠序列化,
对 I/O 密集型任务来说过于沉重。网络请求、数据库查询这类场景,
CPU 大部分时间在空转,瓶颈不在算力,在等待。真正需要的不是更多核心,
而是”谁在等,就让别人先干”。多线程正是为此而生:共享内存、创建轻量、切换高效。
不过线程并不是终点——当成百上千个请求同时到来,线程本身的开销
也会成为新的瓶颈。但在讨论那个更轻量的方案之前,我们先打好基础。
前置知识
- Python 函数与类的基础用法
核心内容
线程基础概念
进程 VS 线程
-
进程:操作系统资源分配的最小单位。每个进程拥有独立的内存地址空间、 文件描述符表和系统资源,进程间通信需要通过管道、共享内存等机制, 开销较大。
-
线程:进程内 CPU 调度的最小单位。同一进程内的线程共享内存空间和 文件描述符等资源,创建和切换的开销远小于进程,又称轻量级进程。
-
关系:一个进程至少包含一个主线程,可动态创建多个子线程。 同一进程内的所有线程共享全局变量和堆内存,但各自拥有独立的栈和 程序计数器。
| 进程 | 线程 | |
|---|---|---|
| 资源分配 | 独立内存、文件描述符表 | 共享所属进程资源 |
| 创建开销 | 重(MB 级) | 轻(KB 级) |
| 并行能力 | 真并行,绕过 GIL | 受 GIL 限制 |
| 适合场景 | CPU 密集型 | I/O 密集型 |
GIL(Global Interpreter Lock)
GIL是CPython 中的一把全局锁,保证同一时刻只有一个线程执行 Python 字节码。 这意味着多线程无法真正利用多核 CPU 进行并行计算。所以可以:
时间线(两个线程共享一个 GIL)
──────────────────────────────────────────
线程A: [计算] [发起网络请求 → 让出GIL] [计算] [完成]
线程B: [ 获得GIL → 趁机计算 ]
──────────────────────────────────────────
创建线程的两种方式
方式1:函数式
import threading
import time
def task(name, delay):
"""模拟耗时任务"""
print(f"线程 {name} 启动")
time.sleep(delay)
print(f"线程 {name} 执行完毕")
if __name__ == "__main__":
start = time.perf_counter()
# 创建两个线程
t1 = threading.Thread(target=task, args=("A", 2))
t2 = threading.Thread(target=task, kwargs={"name": "B", "delay": 1})
# 启动线程
t1.start()
t2.start()
# 等待两个线程都结束
t1.join()
t2.join()
elapsed = time.perf_counter() - start
print(f"所有子线程执行完成,耗时 {elapsed:.1f}s")
# 串行需要 3s,并行只需 2s
方式2:类继承 Thread
import threading
import time
class MyThread(threading.Thread):
def __init__(self, name, delay):
super().__init__(name=name)
self.delay = delay
def run(self):
print(f"线程 {self.name} 启动")
time.sleep(self.delay)
print(f"线程 {self.name} 结束")
if __name__ == "__main__":
t1 = MyThread("线程1", 1)
t2 = MyThread("线程2", 2)
t1.start()
t2.start()
t1.join()
t2.join()
print("程序结束") # 逐个等待,全部跑完再继续
join()不能省
- 主线程/主进程的代码执行完毕后,Python 解释器不会立即退出。
- 解释器会检查是否还有非守护的子线程和子进程在运行。
- 有 → 解释器继续等待;
没有 → 解释器退出,所有守护线程和守护子进程被强制销毁。简记:所有非守护的子线程和子进程都结束,程序才真正退出。
join() 阻塞调用线程(通常是主线程),直到被 join 的线程终止,或超时。
❌ 没有 join —— 可能在子线程结束前就用了结果;
✅ 有 join —— 确保拿到结果。
补充:守护线程
import threading
import time
def daemon_task():
while True:
print('守护线程正在运行中...')
time.sleep(1)
def normal_task():
time.sleep(3)
print('普通子线程结束')
if __name__ == '__main__':
# daemon=True,说明 t1 是守护线程,主线程结束,t1 直接被回收
t1 = threading.Thread(target=daemon_task, daemon=True)
# t2 不是守护线程,主线程要等待 t2 结束之后才能结束
t2 = threading.Thread(target=normal_task)
t1.start()
t2.start()
t2.join()
print('主线程执行完毕,守护线程直接销毁...')
运行逻辑:3 秒后普通线程结束,主线程结束,无限循环的守护线程直接终止。
daemon必须在start()前设置,启动后修改无效
守护线程不能执行文件、数据库落盘操作,主线程退出会强制杀死,数据丢失。
补充:线程常用属性与方法
| 方法 / 属性 | 作用 |
|---|---|
start() | 启动线程,调用 run () |
run() | 线程执行逻辑入口 |
join(timeout=None) | 阻塞主线程等待子线程,timeout 设置最大等待秒数 |
is_alive() | 返回布尔值,判断线程是否正在运行 |
getName() / name | 获取 / 设置线程名称 |
threading.current_thread() | 获取当前运行线程对象 |
threading.enumerate() | 返回当前所有存活线程列表 |
threading.active_count() | 获取当前活跃线程总 |
线程同步工具
多个线程同时改同一个变量会发生竞态:
import threading
import time
num = 0
def add():
global num
for _ in range(100_000):
current = num
time.sleep(0) # 不睡,但让出 GIL,立刻重新排队
num = current + 1 # 非原子操作,多线程互相覆盖 → 结果偏小
if __name__ == "__main__":
t1 = threading.Thread(target=add)
t2 = threading.Thread(target=add)
t1.start()
t2.start()
t1.join()
t2.join()
print("预期200_000,实际结果:", num) # 期望 200_000,实际通常不是
互斥锁 Lock
同一时间仅一个线程获取锁,保证代码块原子操作。
lock = threading.Lock()
def add():
global num
for _ in range(100_000):
with lock: # 等价于 acquire() + try/finally + release()
current = num
time.sleep(0)
num = current + 1
锁的使用原则
with lock:包裹的代码(锁的范围)越小越好,但不能小到漏掉需要保护的操作。锁保护的是数据,不是代码。
共享资源要锁,独立计算不锁。
检查和修改必须在同一把锁内完成。
可重入锁 RLock
函数嵌套调用时需要重复获取锁,普通 Lock 同一线程重复 acquire 会死锁,RLock 支持同一线程多次加锁,需要同等次数释放:
import threading
rlock = threading.RLock()
def transfer(account_from, account_to, amount):
with rlock:
# 获取锁后,调用内部方法
deduct(account_from, amount)
deposit(account_to, amount)
def deduct(account, amount):
with rlock: # 同一线程再次获取,RLock 允许
account.balance -= amount
def deposit(account, amount):
with rlock: # 同一线程再次获取,RLock 允许
account.balance += amount
死锁规避方案
- 统一所有线程获取锁的顺序;
- 设置 acquire 超时
lock.acquire(timeout=3);- 减少嵌套锁;
- 使用 RLock 替代普通 Lock。
信号量 Semaphore:限制最大并发线程
控制同时运行线程数量:
'''
Lock → 独占,一次 1 个 → 保护共享资源
Semaphore(n) → 共享,一次 n 个 → 限制并发数量
RLock → 可重入,同一线程多次获取 → 函数嵌套调用
'''
import threading
import time
# 同时允许 3 个线程访问
sem = threading.Semaphore(3)
def access_resource(thread_id):
print(f"线程{thread_id}: 等待进入...")
with sem:
print(f"线程{thread_id}: ✅ 进入,正在使用资源")
time.sleep(2)
print(f"线程{thread_id}: ❌ 离开")
if __name__ == "__main__":
threads = []
for i in range(8):
t = threading.Thread(target=access_resource, args=(i,))
threads.append(t)
t.start()
for t in threads:
t.join()
print("全部完成")
事件 Event:线程间通知机制
核心方法:
event.set():设置标志为 True,唤醒等待线程event.wait(timeout):阻塞等待标志为 Trueevent.clear():重置标志为 False
import threading
import time
event = threading.Event()
def wait_task():
print('子线程等待信号...')
event.wait() # 阻塞,直到 event 被 set
print('收到信号,并且开始执行任务')
def send_signal():
time.sleep(3)
print('发送通知信号')
event.set() # 通知所有等待的线程
if __name__ == '__main__':
t1 = threading.Thread(target=wait_task)
t2 = threading.Thread(target=send_signal)
t1.start()
t2.start()
t1.join()
t2.join()
条件变量 Condition:复杂生产消费模型
Condition = 锁 + 等待队列。线程在条件不满足时
wait()等待,条件满足时被notify_all()唤醒。
适合「生产者-消费者」这类一写多读的复杂模型。
wait() / notify_all() 机制
wait() 三步曲:释放锁 → 进入等待队列 → 阻塞
notify_all():唤醒等待队列 → 回到锁池重新抢锁被唤醒 ≠ 立即执行:抢到锁的线程才从 wait() 的下一行继续
锁池 vs 等待队列
锁池 :等待拿锁的线程,参与锁竞争
等待队列 :已 wait() 的线程,不参与竞争
wait() → 进等待队列;notify_all() → 回到锁池抢锁
完整可运行示例:
import threading
import time
import random
BUFFER = [] # 共享仓库
MAX, TOTAL = 5, 20 # 仓库容量 / 总任务数
produced = consumed = 0
done = False
cond = threading.Condition() # 锁 + 等待队列
def producer(name):
global produced, done
while True:
with cond:
if produced >= TOTAL:
done = True
cond.notify_all()
return
while len(BUFFER) >= MAX and not done:
cond.wait() # 仓库满 → 等待
if done: # 第二次检查,防止过度生产
cond.notify_all()
return
produced += 1
BUFFER.append(f"{name}-P{produced}")
cond.notify_all() # 唤醒消费者
time.sleep(random.uniform(0.1, 0.3))
def consumer(name):
global consumed, done
while True:
with cond:
while not BUFFER and not done:
cond.wait() # 仓库空 → 等待
if not BUFFER: # 没货且生产结束 → 退出
done = True
cond.notify_all()
return
item = BUFFER.pop(0)
consumed += 1
if consumed >= TOTAL:
done = True
cond.notify_all()
return
cond.notify_all()
time.sleep(random.uniform(0.1, 0.3))
if __name__ == "__main__":
ts = [threading.Thread(target=producer, args=(f"P{i+1}",)) for i in range(2)] + \
[threading.Thread(target=consumer, args=(f"C{i+1}",)) for i in range(3)]
for t in ts:
t.start()
for t in ts:
t.join()
print(f"完成!总生产 {produced},总消费 {consumed}")
为什么用
while不是if?
被唤醒 ≠ 条件一定满足:可能有虚假唤醒,或货已被其他消费者抢走。必须while重新检查,否则会越界pop。惊群效应:
notify_all()唤醒全部线程,但只有抢到锁的那个能干活,其余又回去等待,浪费 CPU。单个唤醒用notify(),但可能造成线程饥饿。
线程本地存储 threading.local ()
多线程共享全局变量会引发数据竞争。threading.local() 为每个线程
创建独立的变量副本,线程之间互不干扰,无需加锁。
import threading
# 共享变量 —— 会竞争
counter = 0
# 线程私有变量 —— 互不干扰
local = threading.local()
def task(name):
global counter
# 每个线程有自己的 local.value
local.value = 0
for _ in range(10000):
local.value += 1
print(f"{name} 的私有值: {local.value}")
# 共享变量 —— 有竞争
for _ in range(10000):
counter += 1
if __name__ == "__main__":
threads = [threading.Thread(target=task, args=(f"线程{i}",)) for i in range(3)]
for t in threads:
t.start()
for t in threads:
t.join()
print(f"\n私有变量: 各线程互不影响,结果都是 10000")
print(f"共享变量: counter = {counter}(预期 30000,实际不足)")
补充:线程间通信方式汇总表
| 机制 | 作用 | 适合场景 |
|---|---|---|
| Lock | 互斥,一次一个线程 | 保护共享资源 |
| RLock | 可重入锁 | 函数嵌套调用同一把锁 |
| Semaphore | 限制并发数 | 连接池、限流 |
| Event | 一对多通知 | 启动信号、状态广播 |
| Condition | 精准唤醒 | 生产消费模型 |
| Queue | 线程安全的队列 | 生产消费(推荐替代 Condition) |
| local | 线程私有变量 | 避免共享、无需加锁 |
线程池:ThreadPoolExecutor
手写线程需要管理
start、join、线程生命周期,代码冗余且容易出错。
ThreadPoolExecutor提供了线程池封装,自动管理线程的创建、复用和回收。
| 类型 | 导入 | 适用场景 | 并发机制 |
|---|---|---|---|
| ThreadPoolExecutor | concurrent.futures.ThreadPoolExecutor | I/O 密集型(网络、文件) | 多线程,受 GIL 限制 |
| ProcessPoolExecutor | concurrent.futures.ProcessPoolExecutor | CPU 密集型(计算、数据) | 多进程,绕过 GIL |
基本用法
from concurrent.futures import ThreadPoolExecutor
import time
def task(num):
print(f"任务{num}开始")
time.sleep(1)
return f"任务{num}完成"
if __name__ == "__main__":
with ThreadPoolExecutor(max_workers=4) as pool: # max_workers:最大并发数
# 写法一:submit —— 返回 Future,用 .result() 取值
futures = [pool.submit(task, i) for i in range(6)]
for f in futures:
print(f.result())
# 写法二:map —— 直接返回结果迭代器
for r in pool.map(task, range(6)):
print(r)
ThreadPoolExecutor 提供两种提交任务的方式:
| 写法 | 返回值 | 取结果 | 适合场景 |
|---|---|---|---|
submit | Future 对象 | .result() | 需要精细控制(取消、回调、超时) |
map | 结果迭代器 | 直接遍历 | 同一函数批量调用 |
两种写法执行流程完全一致,只是封装层级不同:submit 手动创建每个 Future 再 .result() 取值;map 内部自动完成创建与收集。
等待机制:结果什么时候到
submit 拿到 Future 后,取结果有三种方式,差别在「顺序 vs 及时性 vs 精确控制」:
| 方式 | 结果顺序 | 一句话 | 适用场景 |
|---|---|---|---|
遍历 .result() | 按提交顺序 | 逐个阻塞等 | 任务少、不在乎等待 |
as_completed() | 按完成顺序 | 先完成先处理 | 大量 I/O,想尽快消费先到的结果 |
wait() | 不取结果 | 只等一个「里程碑」 | 阶段同步 |
容易踩的坑:for f in futures: f.result() 严格按提交顺序等 —— 如果第 1 个任务最慢,后面全完成了也得干等它。要「先完成先处理」,用 as_completed:
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def task(num):
time.sleep(3 - num % 3) # 让不同任务耗时不同
return f"任务{num}"
if __name__ == "__main__":
with ThreadPoolExecutor(max_workers=4) as pool:
futures = [pool.submit(task, i) for i in range(6)]
for f in as_completed(futures): # 谁先完成先输出谁
print(f.result())
wait() 不直接返回结果,而是等待一组任务达到指定状态后,
把”已完成”和”未完成”分开返回,由你决定下一步怎么处理。
| 条件 | 场景 | 思路 |
|---|---|---|
| FIRST_COMPLETED | 竞速请求、搜索建议 | 谁快用谁 |
| ALL_COMPLETED | 批量处理、分页聚合 | 全部拿到再继续 |
| FIRST_EXCEPTION | 批量执行、快速失败 | 有一个错就停 |
| 用法参考: |
from concurrent.futures import ThreadPoolExecutor, wait, FIRST_COMPLETED, ALL_COMPLETED, FIRST_EXCEPTION
import time
import random
# ============================================================
# FIRST_COMPLETED:谁先完成就先返回
# ============================================================
def fast_search(source):
time.sleep(random.uniform(0.5, 2))
return f"来自 {source} 的结果"
print("=== FIRST_COMPLETED ===")
with ThreadPoolExecutor(3) as pool:
futures = [pool.submit(fast_search, s) for s in ["镜像A", "镜像B", "镜像C"]]
done, pending = wait(futures, return_when=FIRST_COMPLETED)
result = list(done)[0].result()
print(f"最快返回:{result}")
print(f"还在跑:{len(pending)} 个")
# ============================================================
# ALL_COMPLETED:全部完成才返回
# ============================================================
def process(order):
time.sleep(random.uniform(0.5, 1.5))
return f"订单 {order} 已处理"
print("\n=== ALL_COMPLETED ===")
with ThreadPoolExecutor(3) as pool:
futures = [pool.submit(process, o) for o in ["001", "002", "003"]]
done, pending = wait(futures, return_when=ALL_COMPLETED)
print(f"全部完成,共 {len(done)} 个结果:")
for f in done:
print(f" {f.result()}")
# ============================================================
# FIRST_EXCEPTION:有异常就立刻返回
# ============================================================
def risky_task(i):
time.sleep(0.5)
if i == 2:
raise ValueError(f"任务 {i} 失败了!")
return f"任务 {i} 成功"
print("\n=== FIRST_EXCEPTION ===")
with ThreadPoolExecutor(4) as pool:
futures = [pool.submit(risky_task, i) for i in range(5)]
done, pending = wait(futures, return_when=FIRST_EXCEPTION)
for f in done:
if f.exception():
print(f"发现异常:{f.exception()}")
for p in pending:
p.cancel()
break
else:
print(f"已完成:{f.result()}")
print(f"被取消:{len(pending)} 个")
异常捕获 + 任务超时
Future.result (timeout = 秒数) 超时抛 TimeoutError:
from concurrent.futures import ThreadPoolExecutor, TimeoutError
import time
def err_task():
time.sleep(2)
raise ValueError("任务内部报错")
with ThreadPoolExecutor(2) as pool:
f = pool.submit(err_task)
try:
res = f.result(timeout=1)
except TimeoutError:
print("任务执行超时")
except Exception as e:
print("任务异常:", e)
常见问题及补充
什么时候用?
| 场景 | 结论 | 原因 |
|---|---|---|
| 大量网络请求 / 文件读写 | 多线程 | CPU 大部分时间在等 I/O,线程等待时不阻塞其他线程 |
| 纯计算(循环、图像处理) | 多进程 | 绕过 GIL,真正利用多核 CPU 并行计算 |
| 海量连接的网络服务 | asyncio | 单线程协程,万级连接只需极少内存开销 |
| 任务少且很快 | 单线程 | 并发本身有开销,任务太轻反而更慢 |
| I/O + 计算混合 | 多线程 + 多进程 | I/O 用线程,计算提交给进程池 |
死锁
四个条件:互斥、持有并等待、不可抢占、循环等待。
预防方法:
- 统一加锁顺序;
- 使用 RLock(解决同一线程嵌套获取);
- 减少嵌套锁;
- 设置超时(lock.acquire(timeout=5));
实战示例:多线程网络请求对比测试
import time
import urllib.request
from concurrent.futures import ThreadPoolExecutor
URLS = [
"https://httpbin.org/delay/2",
"https://httpbin.org/delay/3",
"https://httpbin.org/delay/1",
]
def fetch(url):
try:
with urllib.request.urlopen(url, timeout=5) as resp:
data = resp.read()
return url, len(data), None
except Exception as e:
return url, 0, str(e)
def show_results(label, results):
for url, size, err in results:
status = f"{size} bytes" if not err else f"失败: {err}"
print(f" {url} → {status}")
if __name__ == "__main__":
# 顺序执行
start = time.perf_counter()
seq_results = [fetch(u) for u in URLS]
seq_time = time.perf_counter() - start
# 多线程并发
start = time.perf_counter()
with ThreadPoolExecutor(max_workers=3) as pool:
con_results = list(pool.map(fetch, URLS))
con_time = time.perf_counter() - start
print(f"顺序执行 ({seq_time:.2f}s):")
show_results("顺序", seq_results)
print(f"\n并发执行 ({con_time:.2f}s):")
show_results("并发", con_results)
print(f"\n加速比: {seq_time / con_time:.2f}x")
结语
多线程擅长 I/O 密集型任务,多进程擅长 CPU 密集型任务, 但两者都有一个共同的限制:当并发量上升到成百上千时, 线程和进程本身的开销也会成为瓶颈。
下一章我们将学习协程——一种更轻量的并发方案, 用单线程实现万级并发。