🚀 25. 并发编程
本章概述
学习 Python 并发编程:多线程、多进程、协程,提升程序效率。 预计学习时间:60 分钟
25.1 并发 vs 并行
概念
- 并发(Concurrency):多个任务交替执行,看起来像同时进行
- 并行(Parallelism):多个任务真正同时执行(需要多核CPU)
Python 的并发方式
| 方式 | 适用场景 | 特点 |
|---|---|---|
| 多线程 | I/O 密集型任务 | 有GIL限制,不能真正并行 |
| 多进程 | CPU 密集型任务 | 可以真正并行,开销大 |
| 协程 | I/O 密集型任务 | 单线程内切换,开销最小 |
GIL(全局解释器锁)
Python 有 GIL(Global Interpreter Lock),同一时刻只能有一个线程执行 Python 字节码。 所以多线程不能真正并行执行 CPU 密集型任务,但 I/O 密集型任务可以受益。
25.2 多线程 threading
基本用法
import threading
import time
def task(name):
print(f"任务 {name} 开始")
time.sleep(2) # 模拟耗时操作
print(f"任务 {name} 结束")
# 创建线程
t1 = threading.Thread(target=task, args=("A",))
t2 = threading.Thread(target=task, args=("B",))
# 启动线程
t1.start()
t2.start()
# 等待线程结束
t1.join()
t2.join()
print("所有任务完成")继承 Thread 类
import threading
import time
class MyThread(threading.Thread):
def __init__(self, name):
super().__init__()
self.name = name
def run(self): # 线程执行的代码
print(f"线程 {self.name} 开始")
time.sleep(2)
print(f"线程 {self.name} 结束")
t1 = MyThread("A")
t2 = MyThread("B")
t1.start()
t2.start()
t1.join()
t2.join()线程安全问题
多个线程同时修改共享数据,会出现问题。
import threading
count = 0
def increment():
global count
for _ in range(100000):
count += 1
t1 = threading.Thread(target=increment)
t2 = threading.Thread(target=increment)
t1.start()
t2.start()
t1.join()
t2.join()
print(count) # 不是200000!因为线程不安全锁 Lock
import threading
count = 0
lock = threading.Lock()
def increment():
global count
for _ in range(100000):
lock.acquire() # 获取锁
count += 1
lock.release() # 释放锁
# 或者用 with 语句(推荐)
def increment():
global count
for _ in range(100000):
with lock: # 自动获取和释放
count += 1
t1 = threading.Thread(target=increment)
t2 = threading.Thread(target=increment)
t1.start()
t2.start()
t1.join()
t2.join()
print(count) # 200000(正确了)常用方法
# 当前线程
threading.current_thread()
# 活跃线程数
threading.active_count()
# 所有线程列表
threading.enumerate()
# 线程是否存活
t.is_alive()
# 设置守护线程(主线程结束时自动结束)
t = threading.Thread(target=task, daemon=True)25.3 多进程 multiprocessing
基本用法
import multiprocessing
import time
def task(name):
print(f"进程 {name} 开始")
time.sleep(2)
print(f"进程 {name} 结束")
if __name__ == "__main__":
# 创建进程
p1 = multiprocessing.Process(target=task, args=("A",))
p2 = multiprocessing.Process(target=task, args=("B",))
# 启动进程
p1.start()
p2.start()
# 等待进程结束
p1.join()
p2.join()
print("所有进程完成")注意
Windows 下多进程代码必须放在
if __name__ == "__main__":里面,否则会报错。
进程池 Pool
import multiprocessing
import time
def task(n):
time.sleep(1)
return n * n
if __name__ == "__main__":
# 创建进程池,大小为CPU核心数
with multiprocessing.Pool() as pool:
# map:批量执行
results = pool.map(task, [1, 2, 3, 4, 5])
print(results) # [1, 4, 9, 16, 25]
# apply_async:异步执行
result = pool.apply_async(task, (10,))
print(result.get()) # 100进程间通信
Queue 队列
import multiprocessing
def producer(q):
for i in range(5):
q.put(i)
print(f"生产:{i}")
def consumer(q):
while True:
item = q.get()
if item is None:
break
print(f"消费:{item}")
if __name__ == "__main__":
q = multiprocessing.Queue()
p1 = multiprocessing.Process(target=producer, args=(q,))
p2 = multiprocessing.Process(target=consumer, args=(q,))
p1.start()
p2.start()
p1.join()
q.put(None) # 发送结束信号
p2.join()25.4 协程 asyncio
协程是比线程更轻量级的并发方式,在单线程内切换。
基本用法
import asyncio
async def task(name, seconds):
print(f"任务 {name} 开始")
await asyncio.sleep(seconds) # 异步等待
print(f"任务 {name} 结束")
return f"{name}完成"
async def main():
# 方式1:顺序执行(慢)
# await task("A", 2)
# await task("B", 1)
# 方式2:并发执行(快)
task1 = task("A", 2)
task2 = task("B", 1)
results = await asyncio.gather(task1, task2)
print(results)
# 运行协程
asyncio.run(main())创建任务
import asyncio
async def task(name, seconds):
print(f"任务 {name} 开始")
await asyncio.sleep(seconds)
print(f"任务 {name} 结束")
return name
async def main():
# 创建任务
t1 = asyncio.create_task(task("A", 2))
t2 = asyncio.create_task(task("B", 1))
# 等待任务完成
result1 = await t1
result2 = await t2
print(result1, result2)
asyncio.run(main())asyncio.gather
并发运行多个协程:
import asyncio
async def task(name, seconds):
await asyncio.sleep(seconds)
return name
async def main():
results = await asyncio.gather(
task("A", 2),
task("B", 1),
task("C", 3)
)
print(results) # ['A', 'B', 'C']
asyncio.run(main())25.5 并发编程对比
什么时候用什么?
| 场景 | 推荐方式 | 原因 |
|---|---|---|
| I/O 密集型(网络请求、文件读写) | 协程 / 多线程 | 等待时可以切换,效率高 |
| CPU 密集型(计算、数据处理) | 多进程 | 绕过GIL,真正并行 |
| 简单任务,代码简单 | 多线程 | 简单易用 |
| 高并发 I/O 任务 | 协程 | 开销最小,效率最高 |
对比总结
| 维度 | 多线程 | 多进程 | 协程 |
|---|---|---|---|
| 切换开销 | 中 | 大 | 小 |
| 数据共享 | 容易(注意锁) | 困难(需要IPC) | 容易 |
| 并行能力 | 不能(GIL) | 能 | 不能 |
| 适用场景 | I/O密集型 | CPU密集型 | I/O密集型 |
| 代码复杂度 | 中 | 高 | 中 |
25.6 常用并发库
concurrent.futures
统一的线程池和进程池接口:
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
import time
def task(n):
time.sleep(1)
return n * n
# 线程池
with ThreadPoolExecutor(max_workers=4) as executor:
results = list(executor.map(task, [1, 2, 3, 4, 5]))
print(results)
# 进程池
with ProcessPoolExecutor(max_workers=4) as executor:
results = list(executor.map(task, [1, 2, 3, 4, 5]))
print(results)submit 和 as_completed
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
def task(name, seconds):
time.sleep(seconds)
return f"{name}完成"
with ThreadPoolExecutor() as executor:
futures = [
executor.submit(task, "A", 3),
executor.submit(task, "B", 1),
executor.submit(task, "C", 2)
]
# 按完成顺序获取结果
for future in as_completed(futures):
print(future.result())🔗 相关章节
📝 我的笔记
在这里记录你的理解和练习代码
# 你的练习代码
✅ 本章检查清单
- 理解并发和并行的概念
- 了解 GIL 的影响
- 会使用多线程 threading
- 了解线程安全和锁
- 会使用多进程 multiprocessing
- 了解进程池的使用
- 了解协程 asyncio 的基本用法
- 知道不同并发方式的适用场景
- 了解 concurrent.futures