预计阅读时间:45 分钟
Python 并发与异步编程深度解析
一、理解并发:从单线程到多任务
在 Python 的世界里,并发编程常被误解和畏惧。但事实上,理解其核心概念后,你会发现它是一个强大而优雅的工具箱。让我们从最基本的问题开始:为什么需要并发?
1.1 I/O 密集型 vs CPU 密集型
import time
import math
def io_bound_task():
"""模拟 I/O 密集型任务:等待外部资源"""
time.sleep(1) # 模拟网络请求、文件读取等
return "I/O 任务完成"
def cpu_bound_task():
"""模拟 CPU 密集型任务:大量计算"""
result = 0
for i in range(10_000_000):
result += math.sqrt(i)
return f"计算结果: {result:.2f}"
# 串行执行
start = time.perf_counter()
io_bound_task()
io_bound_task()
cpu_bound_task()
cpu_bound_task()
serial_time = time.perf_counter() - start
print(f"串行执行耗时: {serial_time:.2f}秒")
print("\n关键洞察:")
print("- I/O 密集型:CPU 大部分时间在等待,适合并发")
print("- CPU 密集型:CPU 持续工作,需要并行才能真正加速")
1.2 Python 的 GIL:必须面对的现实
"""
GIL (Global Interpreter Lock) 深度解析:
GIL 是 CPython 解释器的一个互斥锁,它确保同一时刻只有一个线程执行 Python 字节码。
影响:
- CPU 密集型多线程:GIL 导致无法利用多核,甚至可能更慢
- I/O 密集型多线程:GIL 在 I/O 操作时释放,多线程有效
解决方案:
1. 多进程 (multiprocessing):绕过 GIL,利用多核
2. 异步编程 (asyncio):单线程内协作式多任务
3. C 扩展:在 C 代码中释放 GIL
4. 其他 Python 实现:Jython, IronPython 无 GIL
"""
import threading
import multiprocessing
def cpu_intensive(n):
"""CPU 密集型计算"""
while n > 0:
n -= 1
# 单线程基准
def benchmark_single_thread():
start = time.perf_counter()
cpu_intensive(50_000_000)
cpu_intensive(50_000_000)
return time.perf_counter() - start
# 多线程
def benchmark_multithread():
t1 = threading.Thread(target=cpu_intensive, args=(50_000_000,))
t2 = threading.Thread(target=cpu_intensive, args=(50_000_000,))
start = time.perf_counter()
t1.start()
t2.start()
t1.join()
t2.join()
return time.perf_counter() - start
# 多进程
def benchmark_multiprocess():
p1 = multiprocessing.Process(target=cpu_intensive, args=(50_000_000,))
p2 = multiprocessing.Process(target=cpu_intensive, args=(50_000_000,))
start = time.perf_counter()
p1.start()
p2.start()
p1.join()
p2.join()
return time.perf_counter() - start
print("\n=== GIL 的影响演示 ===")
print(f"单线程: {benchmark_single_thread():.2f}秒")
print(f"多线程 (受 GIL 影响): {benchmark_multithread():.2f}秒")
print(f"多进程 (绕过 GIL): {benchmark_multiprocess():.2f}秒")
二、多线程编程
2.1 threading 模块基础
import threading
import time
from concurrent.futures import ThreadPoolExecutor
print("=== threading 模块基础 ===\n")
# 1. 创建和启动线程
def worker(name, delay):
"""工作函数"""
print(f"线程 {name} 开始工作")
time.sleep(delay)
print(f"线程 {name} 完成工作")
return f"{name} 的结果"
# 创建线程
thread1 = threading.Thread(target=worker, args=("Thread-1", 2))
thread2 = threading.Thread(target=worker, kwargs={"name": "Thread-2", "delay": 1})
# 启动线程
print("启动线程...")
thread1.start()
thread2.start()
# 等待线程完成
thread1.join()
thread2.join()
print("所有线程完成\n")
# 2. 线程同步:Lock
print("=== 线程同步:Lock ===")
counter = 0
counter_lock = threading.Lock()
def increment_without_lock():
global counter
for _ in range(100000):
counter += 1
def increment_with_lock():
global counter
for _ in range(100000):
with counter_lock: # 使用上下文管理器
counter += 1
# 不加锁
counter = 0
threads = [threading.Thread(target=increment_without_lock) for _ in range(5)]
for t in threads:
t.start()
for t in threads:
t.join()
print(f"不加锁结果: {counter} (期望: 500000)")
# 加锁
counter = 0
threads = [threading.Thread(target=increment_with_lock) for _ in range(5)]
for t in threads:
t.start()
for t in threads:
t.join()
print(f"加锁结果: {counter} (期望: 500000)\n")
# 3. 线程同步:RLock (可重入锁)
print("=== RLock 可重入锁 ===")
rlock = threading.RLock()
def recursive_function(n):
"""演示可重入锁"""
with rlock:
if n > 0:
print(f"递归层级 {n}")
recursive_function(n - 1)
recursive_function(3)
print()
# 4. 线程同步:Semaphore (信号量)
print("=== Semaphore 信号量 ===")
semaphore = threading.Semaphore(2) # 最多允许2个线程同时访问
def limited_access(thread_id):
with semaphore:
print(f"线程 {thread_id} 获得访问权限")
time.sleep(1)
print(f"线程 {thread_id} 释放访问权限")
threads = [threading.Thread(target=limited_access, args=(i,)) for i in range(5)]
for t in threads:
t.start()
for t in threads:
t.join()
print()
# 5. 线程同步:Event (事件)
print("=== Event 事件 ===")
event = threading.Event()
def waiter():
print("等待事件...")
event.wait() # 阻塞直到事件被设置
print("事件已触发,继续执行")
def setter():
time.sleep(2)
print("设置事件")
event.set()
threading.Thread(target=waiter).start()
threading.Thread(target=setter).start()
time.sleep(3)
print()
# 6. 线程同步:Condition (条件变量)
print("=== Condition 条件变量 ===")
condition = threading.Condition()
items = []
def producer():
for i in range(5):
with condition:
items.append(f"item-{i}")
print(f"生产: item-{i}")
condition.notify() # 通知等待的消费者
time.sleep(0.5)
def consumer():
for _ in range(5):
with condition:
while not items:
condition.wait() # 等待生产者通知
item = items.pop(0)
print(f"消费: {item}")
threading.Thread(target=producer).start()
threading.Thread(target=consumer).start()
time.sleep(3)
2.2 线程池与 concurrent.futures
from concurrent.futures import ThreadPoolExecutor, as_completed, ProcessPoolExecutor
import time
import random
print("=== 线程池 ThreadPoolExecutor ===\n")
# 1. 基本使用
def task(n):
"""模拟任务"""
time.sleep(random.uniform(0.1, 0.5))
return f"任务 {n} 完成"
# 使用线程池
with ThreadPoolExecutor(max_workers=3) as executor:
# submit 提交单个任务
future1 = executor.submit(task, 1)
future2 = executor.submit(task, 2)
print(f"Future 结果: {future1.result()}")
print(f"Future 结果: {future2.result()}")
# map 批量提交任务
futures = executor.map(task, range(3, 8))
for result in futures:
print(f"Map 结果: {result}")
print()
# 2. as_completed 处理完成的任务
print("=== as_completed 按完成顺序处理 ===")
def slow_task(n):
time.sleep(n * 0.2)
return f"慢任务 {n} (延迟 {n*0.2}s)"
with ThreadPoolExecutor(max_workers=5) as executor:
futures = {executor.submit(slow_task, i): i for i in range(5, 0, -1)}
for future in as_completed(futures):
task_id = futures[future]
try:
result = future.result()
print(f"任务 {task_id}: {result}")
except Exception as e:
print(f"任务 {task_id} 失败: {e}")
print()
# 3. 实际应用:并发下载模拟
print("=== 实际应用:并发下载 ===")
def download_file(url):
"""模拟文件下载"""
print(f"开始下载: {url}")
time.sleep(random.uniform(0.5, 2)) # 模拟下载时间
size = random.randint(100, 1000)
return {"url": url, "size": size, "status": "success"}
urls = [
"http://example.com/file1",
"http://example.com/file2",
"http://example.com/file3",
"http://example.com/file4",
"http://example.com/file5",
]
with ThreadPoolExecutor(max_workers=3) as executor:
future_to_url = {executor.submit(download_file, url): url for url in urls}
total_size = 0
for future in as_completed(future_to_url):
url = future_to_url[future]
try:
result = future.result()
total_size += result['size']
print(f"✓ {url}: {result['size']}KB")
except Exception as e:
print(f"✗ {url}: {e}")
print(f"总下载大小: {total_size}KB\n")
# 4. 线程池 vs 进程池
print("=== 线程池 vs 进程池 ===")
def cpu_task(n):
"""CPU 密集型任务"""
return sum(i * i for i in range(n))
def io_task(n):
"""I/O 密集型任务"""
time.sleep(0.1)
return n
# CPU 密集型
n = 1_000_000
start = time.perf_counter()
with ThreadPoolExecutor(max_workers=4) as executor:
list(executor.map(cpu_task, [n] * 4))
thread_time = time.perf_counter() - start
start = time.perf_counter()
with ProcessPoolExecutor(max_workers=4) as executor:
list(executor.map(cpu_task, [n] * 4))
process_time = time.perf_counter() - start
print(f"CPU 密集型任务:")
print(f" 线程池: {thread_time:.2f}秒")
print(f" 进程池: {process_time:.2f}秒")
print(f" 进程池快 {thread_time/process_time:.1f}倍")
2.3 线程安全的数据结构
import threading
import queue
from collections import deque
print("=== 线程安全的数据结构 ===\n")
# 1. Queue - 线程安全的队列
print("=== Queue 队列 ===")
# 创建队列
q = queue.Queue(maxsize=5)
def producer():
"""生产者"""
for i in range(10):
item = f"产品-{i}"
q.put(item) # 如果队列满,会阻塞
print(f"生产: {item}")
time.sleep(0.1)
def consumer():
"""消费者"""
while True:
try:
item = q.get(timeout=1) # 超时后抛出异常
print(f"消费: {item}")
q.task_done() # 标记任务完成
time.sleep(0.2)
except queue.Empty:
print("队列为空,消费者退出")
break
# 启动生产者和消费者
threading.Thread(target=producer).start()
threading.Thread(target=consumer).start()
time.sleep(3)
print()
# 2. PriorityQueue - 优先级队列
print("=== PriorityQueue 优先级队列 ===")
pq = queue.PriorityQueue()
# 添加任务 (优先级, 任务)
tasks = [(3, "低优先级任务"), (1, "高优先级任务"), (2, "中优先级任务")]
for priority, task in tasks:
pq.put((priority, task))
while not pq.empty():
priority, task = pq.get()
print(f"处理: {task} (优先级 {priority})")
print()
# 3. LifoQueue - 后进先出队列 (栈)
print("=== LifoQueue 栈 ===")
lifo_q = queue.LifoQueue()
for i in range(3):
lifo_q.put(f"任务-{i}")
while not lifo_q.empty():
print(f"LIFO 取出: {lifo_q.get()}")
print()
# 4. 自定义线程安全容器
print("=== 自定义线程安全计数器 ===")
class ThreadSafeCounter:
"""线程安全的计数器"""
def __init__(self):
self._value = 0
self._lock = threading.Lock()
def increment(self):
with self._lock:
self._value += 1
return self._value
def decrement(self):
with self._lock:
self._value -= 1
return self._value
@property
def value(self):
with self._lock:
return self._value
counter = ThreadSafeCounter()
def worker():
for _ in range(1000):
counter.increment()
threads = [threading.Thread(target=worker) for _ in range(10)]
for t in threads:
t.start()
for t in threads:
t.join()
print(f"计数器最终值: {counter.value} (期望: 10000)")
三、多进程编程
3.1 multiprocessing 模块
import multiprocessing as mp
import os
import time
print("=== multiprocessing 模块 ===\n")
# 1. 基本进程创建
def worker_process(name, num):
"""工作进程函数"""
print(f"进程 {name} (PID: {os.getpid()}) 开始")
result = 0
for i in range(num):
result += i * i
print(f"进程 {name} 完成,结果: {result}")
return result
if __name__ == '__main__':
print(f"主进程 PID: {os.getpid()}")
# 创建进程
p1 = mp.Process(target=worker_process, args=("A", 1000))
p2 = mp.Process(target=worker_process, kwargs={"name": "B", "num": 2000})
p1.start()
p2.start()
p1.join()
p2.join()
print("所有进程完成\n")
# 2. 进程间通信:Queue
print("=== 进程间通信:Queue ===")
def producer_process(q):
"""生产者进程"""
for i in range(5):
q.put(f"数据-{i}")
print(f"生产者: 放入 数据-{i}")
time.sleep(0.1)
def consumer_process(q):
"""消费者进程"""
while True:
try:
data = q.get(timeout=1)
print(f"消费者: 取出 {data}")
except:
print("消费者: 队列为空,退出")
break
if __name__ == '__main__':
q = mp.Queue()
p_producer = mp.Process(target=producer_process, args=(q,))
p_consumer = mp.Process(target=consumer_process, args=(q,))
p_producer.start()
p_consumer.start()
p_producer.join()
p_consumer.join()
print()
# 3. 进程间通信:Pipe
print("=== 进程间通信:Pipe ===")
def pipe_worker(conn, name):
"""通过管道通信的工作进程"""
# 接收消息
msg = conn.recv()
print(f"{name} 收到: {msg}")
# 发送回复
conn.send(f"{name} 的回复")
conn.close()
if __name__ == '__main__':
parent_conn, child_conn = mp.Pipe()
p = mp.Process(target=pipe_worker, args=(child_conn, "Worker"))
p.start()
parent_conn.send("主进程的消息")
reply = parent_conn.recv()
print(f"主进程收到: {reply}")
p.join()
print()
# 4. 共享内存:Value 和 Array
print("=== 共享内存:Value 和 Array ===")
def shared_memory_worker(val, arr):
"""操作共享内存的工作进程"""
val.value += 1
for i in range(len(arr)):
arr[i] = arr[i] * 2
if __name__ == '__main__':
# 创建共享变量和数组
shared_val = mp.Value('i', 0) # 'i' 表示整数
shared_arr = mp.Array('i', [1, 2, 3, 4, 5])
print(f"初始值: {shared_val.value}")
print(f"初始数组: {shared_arr[:]}")
processes = []
for _ in range(3):
p = mp.Process(target=shared_memory_worker, args=(shared_val, shared_arr))
processes.append(p)
p.start()
for p in processes:
p.join()
print(f"最终值: {shared_val.value}")
print(f"最终数组: {shared_arr[:]}\n")
# 5. 进程池
print("=== 进程池 Pool ===")
def pool_worker(x):
"""进程池工作函数"""
time.sleep(0.1)
return x * x
if __name__ == '__main__':
with mp.Pool(processes=4) as pool:
# map - 保持顺序
results_map = pool.map(pool_worker, range(10))
print(f"map 结果: {results_map}")
# apply_async - 异步执行
async_result = pool.apply_async(pool_worker, (5,))
print(f"async 结果: {async_result.get()}")
# imap - 迭代器,按顺序返回
for result in pool.imap(pool_worker, range(3)):
print(f"imap 结果: {result}")
3.2 进程间数据共享与同步
import multiprocessing as mp
import time
print("=== 进程同步 ===\n")
# 1. Lock - 互斥锁
def lock_worker(lock, counter, name):
"""使用锁保护共享资源"""
for _ in range(100):
with lock:
counter.value += 1
# 模拟其他工作
time.sleep(0.001)
if __name__ == '__main__':
counter = mp.Value('i', 0)
lock = mp.Lock()
processes = [mp.Process(target=lock_worker, args=(lock, counter, f"P{i}"))
for i in range(5)]
for p in processes:
p.start()
for p in processes:
p.join()
print(f"带锁计数: {counter.value} (期望: 500)\n")
# 2. Event - 事件同步
def waiter_process(event):
"""等待事件"""
print("等待事件触发...")
event.wait()
print("事件已触发,继续执行")
def setter_process(event):
"""触发事件"""
time.sleep(1)
print("触发事件")
event.set()
if __name__ == '__main__':
event = mp.Event()
p1 = mp.Process(target=waiter_process, args=(event,))
p2 = mp.Process(target=setter_process, args=(event,))
p1.start()
p2.start()
p1.join()
p2.join()
print()
# 3. Barrier - 屏障同步
def barrier_worker(barrier, worker_id):
"""屏障同步工作函数"""
print(f"Worker {worker_id} 到达屏障前")
time.sleep(worker_id * 0.5)
print(f"Worker {worker_id} 等待在屏障处")
barrier.wait()
print(f"Worker {worker_id} 通过屏障")
if __name__ == '__main__':
barrier = mp.Barrier(3) # 等待3个进程
processes = [mp.Process(target=barrier_worker, args=(barrier, i))
for i in range(3)]
for p in processes:
p.start()
for p in processes:
p.join()
print()
# 4. Manager - 高级共享对象
def manager_worker(shared_dict, shared_list, name):
"""操作 Manager 共享对象"""
shared_dict[name] = name * 2
shared_list.append(name)
if __name__ == '__main__':
with mp.Manager() as manager:
# 创建共享字典和列表
shared_dict = manager.dict()
shared_list = manager.list()
processes = [mp.Process(target=manager_worker,
args=(shared_dict, shared_list, f"Worker-{i}"))
for i in range(3)]
for p in processes:
p.start()
for p in processes:
p.join()
print(f"共享字典: {dict(shared_dict)}")
print(f"共享列表: {list(shared_list)}")
四、异步编程:asyncio
4.1 协程基础
import asyncio
import time
from typing import Awaitable
print("=== 协程基础 ===\n")
# 1. 定义和运行协程
async def hello():
"""最简单的协程"""
print("Hello")
await asyncio.sleep(1) # 模拟异步等待
print("World")
return "Done"
# 运行协程的三种方式
async def main():
# 方式1:直接 await
result = await hello()
print(f"结果: {result}")
# 方式2:创建任务
task = asyncio.create_task(hello())
result = await task
print(f"任务结果: {result}")
# asyncio.run(main())
print("协程关键概念:")
print("- async def: 定义协程函数")
print("- await: 等待另一个协程完成")
print("- asyncio.sleep(): 异步睡眠,不阻塞事件循环\n")
# 2. 并发执行多个协程
async def fetch_data(name, delay):
"""模拟异步数据获取"""
print(f"开始获取 {name}")
await asyncio.sleep(delay)
print(f"完成获取 {name}")
return f"{name} 的数据"
async def concurrent_demo():
"""演示并发执行"""
# 并发执行多个协程
start = time.perf_counter()
# gather - 并发运行,保持顺序
results = await asyncio.gather(
fetch_data("A", 2),
fetch_data("B", 1),
fetch_data("C", 1.5)
)
elapsed = time.perf_counter() - start
print(f"结果: {results}")
print(f"总耗时: {elapsed:.2f}秒 (如果是串行需要 4.5秒)")
# asyncio.run(concurrent_demo())
print()
# 3. 任务管理
async def task_management_demo():
"""演示任务管理"""
# 创建任务(立即调度)
task1 = asyncio.create_task(fetch_data("Task1", 2))
task2 = asyncio.create_task(fetch_data("Task2", 1))
print(f"Task1 完成状态: {task1.done()}")
# 取消任务
task3 = asyncio.create_task(fetch_data("Task3", 10))
await asyncio.sleep(0.1)
task3.cancel()
try:
await task3
except asyncio.CancelledError:
print("Task3 被取消")
# 等待任务完成
result1 = await task1
result2 = await task2
print(f"结果: {result1}, {result2}")
# asyncio.run(task_management_demo())
print()
# 4. 超时控制
async def timeout_demo():
"""演示超时控制"""
try:
# 设置2秒超时
result = await asyncio.wait_for(
fetch_data("SlowTask", 5),
timeout=2.0
)
print(f"结果: {result}")
except asyncio.TimeoutError:
print("任务超时!")
# 使用 wait 实现更灵活的超时控制
task = asyncio.create_task(fetch_data("Flexible", 3))
done, pending = await asyncio.wait([task], timeout=1.0)
if task in pending:
print("任务未在1秒内完成,取消它")
task.cancel()
else:
result = task.result()
print(f"任务完成: {result}")
# asyncio.run(timeout_demo())
4.2 异步上下文管理器和迭代器
import asyncio
from contextlib import asynccontextmanager
print("=== 异步高级特性 ===\n")
# 1. 异步上下文管理器
class AsyncResource:
"""异步资源管理器"""
async def __aenter__(self):
print("获取异步资源")
await asyncio.sleep(0.1) # 模拟异步初始化
return self
async def __aexit__(self, exc_type, exc_val, exc_tb):
print("释放异步资源")
await asyncio.sleep(0.1) # 模拟异步清理
async def do_work(self):
print("执行工作")
await asyncio.sleep(0.5)
return "工作完成"
async def use_async_resource():
async with AsyncResource() as resource:
result = await resource.do_work()
print(f"结果: {result}")
# asyncio.run(use_async_resource())
print()
# 2. 使用装饰器创建异步上下文管理器
@asynccontextmanager
async def managed_connection(host):
"""使用装饰器创建的异步上下文管理器"""
print(f"连接到 {host}")
await asyncio.sleep(0.1) # 模拟连接
try:
yield f"connection-{host}"
finally:
print(f"关闭 {host} 连接")
await asyncio.sleep(0.1) # 模拟关闭
async def use_managed_connection():
async with managed_connection("example.com") as conn:
print(f"使用 {conn} 发送请求")
await asyncio.sleep(0.5)
# asyncio.run(use_managed_connection())
print()
# 3. 异步迭代器
class AsyncCounter:
"""异步迭代器"""
def __init__(self, start, end):
self.start = start
self.end = end
def __aiter__(self):
return self
async def __anext__(self):
if self.start >= self.end:
raise StopAsyncIteration
await asyncio.sleep(0.1) # 模拟异步操作
value = self.start
self.start += 1
return value
async def use_async_iterator():
print("异步迭代:")
async for num in AsyncCounter(0, 5):
print(f" 得到: {num}")
# asyncio.run(use_async_iterator())
print()
# 4. 异步生成器
async def async_generator():
"""异步生成器"""
for i in range(3):
await asyncio.sleep(0.1)
yield f"item-{i}"
async def use_async_generator():
print("异步生成器:")
async for item in async_generator():
print(f" 生成: {item}")
# 使用列表推导式(需要 async for)
items = [item async for item in async_generator()]
print(f" 收集到: {items}")
# asyncio.run(use_async_generator())
4.3 实际应用:异步 Web 爬虫
import asyncio
import aiohttp
import time
from typing import List, Dict
from urllib.parse import urlparse
print("=== 异步 Web 爬虫示例 ===\n")
class AsyncWebCrawler:
"""异步网页爬虫"""
def __init__(self, max_concurrent=5, timeout=10):
self.max_concurrent = max_concurrent
self.timeout = timeout
self.session = None
self.results = []
async def __aenter__(self):
# 创建带连接池的会话
connector = aiohttp.TCPConnector(limit=self.max_concurrent)
timeout = aiohttp.ClientTimeout(total=self.timeout)
self.session = aiohttp.ClientSession(
connector=connector,
timeout=timeout
)
return self
async def __aexit__(self, *args):
await self.session.close()
async def fetch_url(self, url: str) -> Dict:
"""获取单个 URL"""
try:
async with self.session.get(url) as response:
content = await response.text()
return {
'url': url,
'status': response.status,
'size': len(content),
'content_type': response.headers.get('content-type', ''),
'success': response.status == 200
}
except asyncio.TimeoutError:
return {'url': url, 'error': 'timeout', 'success': False}
except Exception as e:
return {'url': url, 'error': str(e), 'success': False}
async def crawl(self, urls: List[str]) -> List[Dict]:
"""爬取多个 URL"""
print(f"开始爬取 {len(urls)} 个 URL...")
start = time.perf_counter()
# 创建所有任务
tasks = [self.fetch_url(url) for url in urls]
# 并发执行
results = await asyncio.gather(*tasks, return_exceptions=True)
# 处理结果
processed_results = []
for result in results:
if isinstance(result, Exception):
processed_results.append({'error': str(result), 'success': False})
else:
processed_results.append(result)
elapsed = time.perf_counter() - start
print(f"爬取完成,耗时: {elapsed:.2f}秒")
return processed_results
async def crawl_with_semaphore(self, urls: List[str]) -> List[Dict]:
"""使用信号量控制并发数"""
semaphore = asyncio.Semaphore(self.max_concurrent)
async def fetch_with_semaphore(url):
async with semaphore:
return await self.fetch_url(url)
tasks = [fetch_with_semaphore(url) for url in urls]
return await asyncio.gather(*tasks)
async def crawl_streaming(self, urls: List[str]):
"""流式处理结果(边完成边处理)"""
tasks = [asyncio.create_task(self.fetch_url(url)) for url in urls]
for task in asyncio.as_completed(tasks):
result = await task
yield result
# 模拟 URL 列表(实际使用时替换为真实 URL)
test_urls = [
"https://httpbin.org/delay/1",
"https://httpbin.org/delay/2",
"https://httpbin.org/delay/1",
"https://httpbin.org/status/200",
"https://httpbin.org/status/404",
]
async def demo_crawler():
"""演示爬虫"""
async with AsyncWebCrawler(max_concurrent=3) as crawler:
# 方式1:并发爬取
results = await crawler.crawl(test_urls)
print("\n爬取结果:")
for result in results:
if result.get('success'):
print(f" ✓ {result['url']}: {result['status']} ({result['size']} bytes)")
else:
print(f" ✗ {result.get('url', 'unknown')}: {result.get('error', 'failed')}")
# 方式2:流式处理
print("\n流式处理结果:")
async for result in crawler.crawl_streaming(test_urls[:3]):
if result.get('success'):
print(f" → {result['url']} 完成")
# asyncio.run(demo_crawler())
print("\n注意:需要安装 aiohttp: pip install aiohttp\n")
# 4. 同步 vs 异步性能对比
def sync_fetch_simulate():
"""模拟同步爬虫"""
urls = [f"http://example.com/{i}" for i in range(10)]
results = []
for url in urls:
time.sleep(0.1) # 模拟网络延迟
results.append(f"{url} - done")
return results
async def async_fetch_simulate():
"""模拟异步爬虫"""
urls = [f"http://example.com/{i}" for i in range(10)]
tasks = [asyncio.sleep(0.1) for _ in urls] # 模拟并发网络请求
await asyncio.gather(*tasks)
return [f"{url} - done" for url in urls]
print("=== 同步 vs 异步性能对比 ===")
# 同步
start = time.perf_counter()
sync_fetch_simulate()
sync_time = time.perf_counter() - start
# 异步
start = time.perf_counter()
# asyncio.run(async_fetch_simulate())
# 手动模拟异步时间
async_time = 0.1 # 理论上的异步时间
print(f"同步方式 (10个请求): {sync_time:.2f}秒")
print(f"异步方式 (10个请求): ~{async_time:.2f}秒")
print(f"异步快约 {sync_time/async_time:.0f}倍")
五、并发模式与最佳实践
5.1 生产者-消费者模式
import asyncio
import queue
import threading
import time
import random
from typing import List
print("=== 生产者-消费者模式 ===\n")
# 1. 多线程版本
class ThreadedProducerConsumer:
"""多线程生产者-消费者"""
def __init__(self, queue_size=10, num_producers=2, num_consumers=3):
self.queue = queue.Queue(maxsize=queue_size)
self.num_producers = num_producers
self.num_consumers = num_consumers
self.stop_event = threading.Event()
def producer(self, producer_id):
"""生产者"""
while not self.stop_event.is_set():
try:
item = f"Producer-{producer_id}-Item-{random.randint(1, 100)}"
self.queue.put(item, timeout=1)
print(f"[生产] {item} (队列大小: {self.queue.qsize()})")
time.sleep(random.uniform(0.1, 0.5))
except queue.Full:
continue
def consumer(self, consumer_id):
"""消费者"""
while not self.stop_event.is_set():
try:
item = self.queue.get(timeout=1)
print(f"[消费] Consumer-{consumer_id} 处理: {item}")
time.sleep(random.uniform(0.2, 0.8))
self.queue.task_done()
except queue.Empty:
continue
def start(self, duration=5):
"""启动系统"""
print(f"启动生产者-消费者系统...")
print(f" 生产者: {self.num_producers}, 消费者: {self.num_consumers}")
# 创建并启动生产者
producers = []
for i in range(self.num_producers):
t = threading.Thread(target=self.producer, args=(i,))
t.start()
producers.append(t)
# 创建并启动消费者
consumers = []
for i in range(self.num_consumers):
t = threading.Thread(target=self.consumer, args=(i,))
t.start()
consumers.append(t)
# 运行指定时间后停止
time.sleep(duration)
self.stop_event.set()
# 等待所有线程结束
for t in producers + consumers:
t.join(timeout=2)
print(f"系统停止,队列剩余: {self.queue.qsize()}")
return self
# 运行示例
# system = ThreadedProducerConsumer(queue_size=5, num_producers=2, num_consumers=2)
# system.start(duration=3)
# 2. 异步版本
class AsyncProducerConsumer:
"""异步生产者-消费者"""
def __init__(self, queue_size=10):
self.queue = asyncio.Queue(maxsize=queue_size)
async def producer(self, producer_id, num_items=5):
"""异步生产者"""
for i in range(num_items):
item = f"AsyncProducer-{producer_id}-Item-{i}"
await self.queue.put(item)
print(f"[异步生产] {item}")
await asyncio.sleep(random.uniform(0.1, 0.3))
async def consumer(self, consumer_id):
"""异步消费者"""
while True:
try:
# 使用超时避免无限等待
item = await asyncio.wait_for(self.queue.get(), timeout=1.0)
print(f"[异步消费] Consumer-{consumer_id} 处理: {item}")
await asyncio.sleep(random.uniform(0.2, 0.5))
self.queue.task_done()
except asyncio.TimeoutError:
if self.queue.empty():
break
async def run(self, num_producers=2, num_consumers=2):
"""运行异步系统"""
print("启动异步生产者-消费者系统...")
# 创建生产者任务
producer_tasks = [
asyncio.create_task(self.producer(i, 5))
for i in range(num_producers)
]
# 等待所有生产者完成
await asyncio.gather(*producer_tasks)
# 创建消费者任务
consumer_tasks = [
asyncio.create_task(self.consumer(i))
for i in range(num_consumers)
]
# 等待队列清空
await self.queue.join()
# 取消消费者
for task in consumer_tasks:
task.cancel()
print("异步系统完成")
# asyncio.run(AsyncProducerConsumer(queue_size=10).run())
5.2 工作池模式
import asyncio
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
import time
import random
print("=== 工作池模式 ===\n")
# 1. 通用工作池
class WorkPool:
"""通用工作池"""
def __init__(self, num_workers=4, use_process=False):
self.num_workers = num_workers
if use_process:
self.executor = ProcessPoolExecutor(max_workers=num_workers)
else:
self.executor = ThreadPoolExecutor(max_workers=num_workers)
def submit(self, func, *args, **kwargs):
"""提交任务"""
return self.executor.submit(func, *args, **kwargs)
def map(self, func, iterable):
"""批量提交"""
return self.executor.map(func, iterable)
def shutdown(self):
"""关闭工作池"""
self.executor.shutdown()
def __enter__(self):
return self
def __exit__(self, *args):
self.shutdown()
def work_function(task_id, complexity):
"""工作函数"""
print(f"任务 {task_id} 开始 (复杂度: {complexity})")
time.sleep(complexity * 0.1)
result = f"Task-{task_id}-Result"
print(f"任务 {task_id} 完成")
return result
# 使用工作池
print("线程工作池:")
with WorkPool(num_workers=3) as pool:
futures = []
for i in range(5):
complexity = random.randint(1, 5)
future = pool.submit(work_function, i, complexity)
futures.append((i, future))
for task_id, future in futures:
result = future.result()
print(f" 任务 {task_id} 结果: {result}")
print()
# 2. 混合工作池 (异步 + 线程池)
class HybridWorkPool:
"""混合工作池:异步接口 + 线程池执行"""
def __init__(self, max_workers=4):
self.executor = ThreadPoolExecutor(max_workers=max_workers)
self.loop = asyncio.get_event_loop()
async def submit_async(self, func, *args, **kwargs):
"""异步提交任务"""
return await self.loop.run_in_executor(
self.executor,
lambda: func(*args, **kwargs)
)
async def map_async(self, func, items):
"""异步批量提交"""
tasks = [self.submit_async(func, item) for item in items]
return await asyncio.gather(*tasks)
def shutdown(self):
self.executor.shutdown()
def blocking_io_task(item):
"""模拟阻塞 I/O 任务"""
time.sleep(random.uniform(0.1, 0.5))
return f"Processed-{item}"
async def hybrid_demo():
pool = HybridWorkPool(max_workers=5)
items = list(range(10))
print(f"处理 {len(items)} 个项目...")
start = time.perf_counter()
results = await pool.map_async(blocking_io_task, items)
elapsed = time.perf_counter() - start
print(f"结果: {results[:3]}...")
print(f"耗时: {elapsed:.2f}秒")
pool.shutdown()
# asyncio.run(hybrid_demo())
5.3 发布-订阅模式
import asyncio
from typing import Callable, Dict, List, Any
from dataclasses import dataclass
from datetime import datetime
print("=== 发布-订阅模式 ===\n")
# 1. 简单事件总线
class EventBus:
"""事件总线"""
def __init__(self):
self._subscribers: Dict[str, List[Callable]] = {}
def subscribe(self, event_type: str, callback: Callable):
"""订阅事件"""
if event_type not in self._subscribers:
self._subscribers[event_type] = []
self._subscribers[event_type].append(callback)
print(f"订阅 {event_type}: {callback.__name__}")
def unsubscribe(self, event_type: str, callback: Callable):
"""取消订阅"""
if event_type in self._subscribers:
self._subscribers[event_type].remove(callback)
def publish(self, event_type: str, data: Any = None):
"""发布事件"""
if event_type in self._subscribers:
print(f"\n发布事件: {event_type}")
for callback in self._subscribers[event_type]:
callback(data)
# 订阅者回调
def email_handler(data):
print(f" [邮件] 发送邮件: {data}")
def sms_handler(data):
print(f" [短信] 发送短信: {data}")
def log_handler(data):
print(f" [日志] 记录事件: {data}")
# 使用事件总线
bus = EventBus()
bus.subscribe("user_registered", email_handler)
bus.subscribe("user_registered", log_handler)
bus.subscribe("order_created", sms_handler)
bus.subscribe("order_created", email_handler)
bus.publish("user_registered", {"user": "alice", "email": "alice@example.com"})
bus.publish("order_created", {"order_id": 123, "amount": 99.99})
print()
# 2. 异步事件总线
@dataclass
class Event:
"""事件"""
type: str
data: Any
timestamp: datetime = None
def __post_init__(self):
if self.timestamp is None:
self.timestamp = datetime.now()
class AsyncEventBus:
"""异步事件总线"""
def __init__(self):
self._subscribers: Dict[str, List[Callable]] = {}
self._queue = asyncio.Queue()
def subscribe(self, event_type: str, callback: Callable):
"""订阅事件"""
if event_type not in self._subscribers:
self._subscribers[event_type] = []
self._subscribers[event_type].append(callback)
async def publish(self, event: Event):
"""发布事件"""
await self._queue.put(event)
async def _process_events(self):
"""处理事件循环"""
while True:
event = await self._queue.get()
if event.type in self._subscribers:
tasks = []
for callback in self._subscribers[event.type]:
if asyncio.iscoroutinefunction(callback):
tasks.append(callback(event))
else:
callback(event)
if tasks:
await asyncio.gather(*tasks)
async def start(self):
"""启动事件处理"""
self._processor_task = asyncio.create_task(self._process_events())
async def stop(self):
"""停止事件处理"""
self._processor_task.cancel()
async def async_email_handler(event: Event):
await asyncio.sleep(0.1) # 模拟异步操作
print(f" [异步邮件] 发送给 {event.data.get('email')}")
async def async_log_handler(event: Event):
print(f" [异步日志] {event.timestamp}: {event.type} - {event.data}")
async def async_event_demo():
bus = AsyncEventBus()
bus.subscribe("user_login", async_email_handler)
bus.subscribe("user_login", async_log_handler)
await bus.start()
await bus.publish(Event("user_login", {"user": "bob", "email": "bob@example.com"}))
await bus.publish(Event("user_login", {"user": "charlie", "email": "charlie@example.com"}))
await asyncio.sleep(1)
await bus.stop()
# asyncio.run(async_event_demo())
六、调试与性能分析
6.1 并发调试技巧
import threading
import asyncio
import logging
import time
import sys
print("=== 并发调试技巧 ===\n")
# 1. 配置日志记录
logging.basicConfig(
level=logging.DEBUG,
format='%(asctime)s [%(threadName)s] %(levelname)s: %(message)s',
handlers=[
logging.StreamHandler(sys.stdout)
]
)
def debug_thread_function():
"""带日志的线程函数"""
logger = logging.getLogger(__name__)
logger.info("线程开始执行")
time.sleep(0.5)
logger.info("线程执行完成")
print("带日志的线程:")
thread = threading.Thread(target=debug_thread_function, name="DebugThread")
thread.start()
thread.join()
print()
# 2. 线程状态检查
def check_thread_state():
"""检查线程状态"""
print(f"当前线程: {threading.current_thread().name}")
print(f"活跃线程数: {threading.active_count()}")
print("所有线程:")
for t in threading.enumerate():
print(f" {t.name}: alive={t.is_alive()}, daemon={t.daemon}")
check_thread_state()
print()
# 3. 死锁检测示例
def deadlock_detection_demo():
"""演示死锁情况"""
lock1 = threading.Lock()
lock2 = threading.Lock()
def worker1():
with lock1:
time.sleep(0.1)
with lock2:
print("Worker1 完成")
def worker2():
with lock2:
time.sleep(0.1)
with lock1:
print("Worker2 完成")
# 使用超时避免死锁
def safe_worker1():
if lock1.acquire(timeout=1):
try:
time.sleep(0.1)
if lock2.acquire(timeout=1):
try:
print("SafeWorker1 完成")
finally:
lock2.release()
finally:
lock1.release()
else:
print("SafeWorker1 获取锁超时")
def safe_worker2():
if lock2.acquire(timeout=1):
try:
time.sleep(0.1)
if lock1.acquire(timeout=1):
try:
print("SafeWorker2 完成")
finally:
lock1.release()
finally:
lock2.release()
else:
print("SafeWorker2 获取锁超时")
print("安全版本(带超时):")
t1 = threading.Thread(target=safe_worker1)
t2 = threading.Thread(target=safe_worker2)
t1.start()
t2.start()
t1.join()
t2.join()
deadlock_detection_demo()
print()
# 4. 异步调试
async def async_debug_demo():
"""异步调试演示"""
# 启用 asyncio 调试
import asyncio
# 设置调试模式
asyncio.get_event_loop().set_debug(True)
async def slow_coroutine():
await asyncio.sleep(1)
return "Done"
# 创建任务但忘记 await(调试模式会警告)
task = asyncio.create_task(slow_coroutine())
# 正确做法:await task
print("异步调试模式已启用(会显示警告)")
# asyncio.run(async_debug_demo())
6.2 性能基准测试
import time
import asyncio
import threading
import multiprocessing as mp
from concurrent.futures import ThreadPoolExecutor, ProcessPoolExecutor
from typing import Callable, List
import functools
print("=== 性能基准测试 ===\n")
class Benchmark:
"""性能基准测试"""
def __init__(self, name: str):
self.name = name
self.results = {}
def measure(self, func: Callable, *args, **kwargs) -> float:
"""测量函数执行时间"""
start = time.perf_counter()
func(*args, **kwargs)
elapsed = time.perf_counter() - start
return elapsed
def compare_io_bound(self, num_tasks: int = 10, task_duration: float = 0.1):
"""比较 I/O 密集型任务的性能"""
print(f"\nI/O 密集型任务 ({num_tasks} 个任务,每个 {task_duration}秒)")
def io_task():
time.sleep(task_duration)
async def async_io_task():
await asyncio.sleep(task_duration)
# 串行
def serial():
for _ in range(num_tasks):
io_task()
serial_time = self.measure(serial)
print(f" 串行: {serial_time:.2f}秒")
# 多线程
def multithread():
with ThreadPoolExecutor(max_workers=num_tasks) as executor:
list(executor.map(lambda _: io_task(), range(num_tasks)))
thread_time = self.measure(multithread)
print(f" 多线程: {thread_time:.2f}秒 (提升 {serial_time/thread_time:.1f}x)")
# 异步
async def async_main():
tasks = [async_io_task() for _ in range(num_tasks)]
await asyncio.gather(*tasks)
async_time = self.measure(lambda: asyncio.run(async_main()))
print(f" 异步: {async_time:.2f}秒 (提升 {serial_time/async_time:.1f}x)")
def compare_cpu_bound(self, num_tasks: int = 4, complexity: int = 5_000_000):
"""比较 CPU 密集型任务的性能"""
print(f"\nCPU 密集型任务 ({num_tasks} 个任务,复杂度 {complexity:,})")
def cpu_task():
total = 0
for i in range(complexity):
total += i * i
return total
# 串行
def serial():
for _ in range(num_tasks):
cpu_task()
serial_time = self.measure(serial)
print(f" 串行: {serial_time:.2f}秒")
# 多线程
def multithread():
with ThreadPoolExecutor(max_workers=num_tasks) as executor:
list(executor.map(lambda _: cpu_task(), range(num_tasks)))
thread_time = self.measure(multithread)
print(f" 多线程: {thread_time:.2f}秒 (提升 {serial_time/thread_time:.1f}x)")
print(f" 注意:受 GIL 影响,多线程可能没有提升")
# 多进程
def multiprocess():
with ProcessPoolExecutor(max_workers=num_tasks) as executor:
list(executor.map(lambda _: cpu_task(), range(num_tasks)))
process_time = self.measure(multiprocess)
print(f" 多进程: {process_time:.2f}秒 (提升 {serial_time/process_time:.1f}x)")
# 运行基准测试
benchmark = Benchmark("Concurrency Benchmark")
benchmark.compare_io_bound(num_tasks=10, task_duration=0.1)
benchmark.compare_cpu_bound(num_tasks=4, complexity=3_000_000)
七、最佳实践与总结
7.1 并发编程决策树
"""
并发编程决策指南:
1. 任务类型判断
├── I/O 密集型 (网络请求、文件读写、数据库操作)
│ ├── 简单并发 → asyncio (单线程异步)
│ ├── 需要线程安全 → threading (多线程)
│ └── 大量并发连接 → asyncio + aiohttp
│
└── CPU 密集型 (数值计算、图像处理、加密)
├── 利用多核 → multiprocessing (多进程)
├── 需要共享内存 → multiprocessing.Manager
└── 分布式计算 → concurrent.futures + 进程池
2. 具体场景选择
├── Web 爬虫 → asyncio + aiohttp
├── Web 服务器 → asyncio (FastAPI, Sanic) 或 threading (Flask, Django)
├── 数据处理 → multiprocessing.Pool
├── GUI 应用 → threading (避免阻塞主线程)
├── 实时系统 → asyncio (低延迟)
└── 遗留代码集成 → threading (简单易用)
3. 性能考虑
├── 少量长任务 → threading / multiprocessing
├── 大量短任务 → asyncio
├── 混合 I/O 和 CPU → 线程池 + 异步
└── 需要超时控制 → asyncio.wait_for / ThreadPoolExecutor
"""
def concurrency_decision_helper(task_type: str, num_tasks: int, need_shared_state: bool):
"""并发方案推荐"""
if task_type == "io":
if num_tasks > 100:
return "推荐: asyncio (处理大量并发连接)"
elif need_shared_state:
return "推荐: threading (需要共享状态)"
else:
return "推荐: asyncio 或 ThreadPoolExecutor"
elif task_type == "cpu":
if need_shared_state:
return "推荐: multiprocessing.Manager (需要进程间共享)"
else:
return "推荐: ProcessPoolExecutor (充分利用多核)"
else: # mixed
return "推荐: 异步 + run_in_executor (混合模式)"
print("并发方案推荐示例:")
print(f" I/O, 1000任务, 不需要共享: {concurrency_decision_helper('io', 1000, False)}")
print(f" CPU, 8任务, 不需要共享: {concurrency_decision_helper('cpu', 8, False)}")
7.2 最佳实践清单
"""
并发编程最佳实践清单:
1. 线程安全
✓ 使用 queue.Queue 进行线程间通信
✓ 使用 threading.Lock 保护共享资源
✓ 避免在持有锁时调用外部代码
✓ 使用 threading.local 存储线程本地数据
2. 异步编程
✓ 使用 asyncio.run() 运行主协程
✓ 使用 asyncio.create_task() 创建后台任务
✓ 使用 asyncio.gather() 并发执行多个协程
✓ 避免在协程中使用阻塞调用(使用 run_in_executor)
✓ 设置合理的超时时间
3. 多进程
✓ 将主代码放在 if __name__ == '__main__': 中
✓ 使用 Pool 管理进程池
✓ 使用 Manager 共享复杂数据结构
✓ 注意进程间通信的开销
4. 错误处理
✓ 始终处理并发任务的异常
✓ 使用 asyncio.wait_for 设置超时
✓ 实现优雅关闭机制
✓ 记录并发相关的日志
5. 测试与调试
✓ 编写并发测试用例
✓ 使用 timeouts 避免测试挂起
✓ 启用 asyncio 调试模式
✓ 使用线程转储分析死锁
6. 性能优化
✓ 选择合适的并发模型
✓ 调优线程池/进程池大小
✓ 避免过度同步
✓ 使用连接池复用资源
"""
# 示例:优雅关闭
class GracefulShutdown:
"""优雅关闭示例"""
def __init__(self):
self.shutdown_event = asyncio.Event()
self.tasks = []
async def worker(self, name):
"""工作协程"""
while not self.shutdown_event.is_set():
print(f"{name} 工作中...")
try:
await asyncio.wait_for(
self.shutdown_event.wait(),
timeout=1.0
)
except asyncio.TimeoutError:
continue
print(f"{name} 收到关闭信号,清理资源...")
async def run(self):
"""运行系统"""
# 启动工作协程
self.tasks = [
asyncio.create_task(self.worker(f"Worker-{i}"))
for i in range(3)
]
# 运行一段时间
await asyncio.sleep(3)
# 发送关闭信号
print("\n发送关闭信号...")
self.shutdown_event.set()
# 等待所有任务完成
await asyncio.gather(*self.tasks)
print("系统已优雅关闭")
# asyncio.run(GracefulShutdown().run())
7.3 总结
核心概念回顾: - 并发 vs 并行:并发是任务交替执行,并行是同时执行 - GIL:CPython 的全局解释器锁,限制了多线程的 CPU 并行 - 协程:轻量级的用户态线程,通过 async/await 实现协作式多任务
三种并发模型的比较:
| 特性 | threading | multiprocessing | asyncio |
|---|---|---|---|
| 适用场景 | I/O 密集型 | CPU 密集型 | 高并发 I/O |
| 内存开销 | 低 | 高 | 极低 |
| 通信方式 | 共享内存 | IPC | 共享状态 |
| GIL 影响 | 受限 | 无影响 | 无影响 |
| 学习曲线 | 中等 | 中等 | 较陡 |
选择建议: - 默认选择:对于大多数 I/O 密集型应用,asyncio 是最佳选择 - 简单优先:对于简单的并发需求,ThreadPoolExecutor 足够 - CPU 密集:multiprocessing 是唯一能利用多核的选择 - 混合场景:asyncio + run_in_executor 组合使用
记住:并发编程的目标不是让代码更复杂,而是让程序更高效、响应更快。选择合适的工具,理解其原理,你就能写出优雅而高效的并发代码。
本文由 尚先生 原创,转载请注明出处。
评论
0