Python 并发与异步 Python Concurrency And Asynchrony
并发 VS 并行
并发(Concurrency):一个 CPU 核心交替执行多个任务,宏观上看起来像同时,但微观上快速地在多个任务间切换串行地执行。
并行(Parallelism):多个 CPU 核心同时执行多个任务,真正的同时进行。

CPU 密集型和 IO 密集型
CPU密集型任务:任务大部分时间在占用 CPU 做运算,几乎没有等待。如矩阵运算、AI 模型推理、视频编解码等。
IO密集型任务:任务大部分时间在等待外部资源,CPU 多数时间闲着。
- 网络 IO:HTTP 请求、爬虫、接口调用、Redis / MQ 网络通信。
进程 VS 线程 VS 协程
进程 Process:资源分配的最小单位。程序跑起来后就是进程。一个进程拥有独立内存空间、文件句柄、CPU时间、堆、环境变量。
- 进程切换开销巨大:操作系统要刷新页表、缓存、保存全套上下文。
线程 Thread:CPU 调度执行的最小单位。线程寄生在进程内部,一个进程至少有 1 个主线程。
协程 Coroutine:用户态的轻量级 "微线程"。
- 协程不由操作系统内核调度,完全由应用程序代码自己调度,内核完全不知道协程的存在。
- 调度规则:主动让出 CPU (yield / await),协程不是抢占式,是协作调度式。
- 如果某个协程死循环让不出 CPU,其他协程将会一直等待。
GIL - Global - Interpreter Lock
全局解释器锁:是 CPython 解释器中的一个互斥锁,它保证同一进程内、同一时刻只有一个线程能执行 Python 字节码。单进程下,无论硬件设备有多少个CPU核心,Python 的多个线程并不能真正并行执行 Python 代码。不同进程拥有独立解释器、独立 GIL,多进程场景下,多个进程的线程可以跨 CPU 核心真正并行执行 Python 代码。
| | |
|---|
| 核心数量 | | |
| 执行方式 | | 同一时刻同时执行 |
| 本质 | | |
| 目标 | 充分利用空闲等待资源(IO 阻塞时调度其他任务),减少等待耗时 | |
| Python实现 | 线程、协程 (asyncio)、多进程也可以并发(所有并行程序一定并发) | |
| 受GIL影响 | 协程运行在单线程内,同一时刻只有一个协程在执行,因此不存在多线程争抢 GIL 的问题,也就没有线程上下文切换和锁竞争的开销。但协程本质上仍然受 GIL 约束,因为其底层依托的线程受 GIL 限制。 | |
Python 多线程编程
Python 的多线程编程适合 IO密集型操作而非 CPU密集型操作。
多线程基本创建方式
方式一:直接创建 Thread 对象。
import threadingimport timedefworker(name, delay):"""线程要执行的函数""" print(f"线程 {name} 开始工作") time.sleep(delay) # 模拟耗时操作 print(f"线程 {name} 完成工作")# 创建线程t1 = threading.Thread(target=worker, args=("A", 2))t2 = threading.Thread(target=worker, args=("B", 1))# 启动线程t1.start()t2.start()# 等待线程结束t1.join()t2.join()print("所有线程完成")
threading.Thread() 类,传入线程要执行的函数和参数表。join() 方法用于将已有的线程加入等待,使主线程阻塞,,等待所有子线程结束后再执行余下的代码。如果没有 join() 操作,那么主线程就会立即执行剩下的逻辑,打印 print("所有线程完成")。
方式二:继承 Thread 类,添加更多自定义的逻辑。
import threadingimport timeclassWorkerThread(threading.Thread):"""自定义线程类"""def__init__(self, name, delay): super().__init__() self.name = name self.delay = delay self.result = None# 添加一个 result 成员变量,存储线程执行结果defrun(self):"""重写 run 方法,线程执行的代码""" print(f"线程 {self.name} 开始工作") time.sleep(self.delay) self.result = f"线程 {self.name} 的结果" print(f"线程 {self.name} 完成工作")return self.result # 返回值无法被直接获取# 使用t1 = WorkerThread("A", 2)t1.start() # start() 方法线程启动,并自动调用 run() 方法# t1.run() 不创建新线程,直接在当前线程执行 run() 方法里面的逻辑,run() 方法执行完毕才能执行后面的代码t2 = WorkerThread("B", 3)t2.start()t1.join()t2.join()print(t1.result, t2.result) # 通过实例属性获取结果
线程启动时,start() 方法会自动调用 run() 方法,如果不用 start() 方法来启动线程直接调用 run() ,那么相当于在当前线程直接调用了一个函数。
线程生命周期管理

Python 多线程生命周期状态图简洁版(不包含守护线程):

守护线程
创建线程时,传入 daemon=True 参数,线程会变为守护线程。当程序中所有非守护线程执行完毕,Python 解释器直接退出,所有还在运行的守护线程会被强制终止,不会等待任务完成。
import threadingimport timedefdaemon_worker():"""写一个后台运行,永不结束的子线程""" count = 0whileTrue: count += 1 print(f"守护线程运行中... {count}") time.sleep(1)# 创建守护线程(daemon=True)daemon = threading.Thread(target=daemon_worker, daemon=True)daemon.start()# 主线程等待3秒time.sleep(3)print("主线程结束")# 守护线程会随主线程结束而自动终止# 输出:大约3-4次后程序退出
守护线程适用于:
等,不适合在守护线程中写带有资源释放的操作。
线程同步机制
Lock 互斥锁 Mutex - Mutual Exclusion
- 一个线程使用完毕忘记释放锁,或者多个线程互相占用对方想要使用的资源,都在等待彼此释放,会造成死锁。
import threadingcounter = 0# 多个线程中都要使用的 Global 变量lock = threading.Lock()defsafe_increment():global counterfor _ in range(100000):with lock: # 使用上下文管理器,自动获取和释放 counter += 1# 或者手动方式defsafe_increment_manual():global counterfor _ in range(100000): lock.acquire()try: counter += 1finally: lock.release() # 必须释放!# 测试threads = []for _ in range(10): t = threading.Thread(target=safe_increment) threads.append(t) t.start()for t in threads: t.join()print(f"期望结果: 1000000")print(f"实际结果: {counter}") # 1000000
RLock 可重入锁
允许同一个线程多次获取锁(递归锁),适用于递归函数或嵌套函数。
import threading# 使用普通 Lock(会死锁)# lock = threading.Lock()# 使用 RLock(不会死锁)rlock = threading.RLock()defrecursive_function(n):if n <= 0:returnwith rlock: # 同一个线程可以重复获取 print(f"递归深度: {n}") recursive_function(n - 1) # 再次获取同一个锁# 递归调用10层,每层都获取锁recursive_function(10)
Semaphore 信号量
限制同时访问资源的线程数量。
import threadingimport time# 最多允许3个线程同时访问semaphore = threading.Semaphore(3)deflimited_worker(name): print(f"{name}: 等待获取信号量")with semaphore: # 获取信号量 print(f"{name}: 获取信号量,开始工作") time.sleep(2) # 模拟耗时操作 print(f"{name}: 工作完成,释放信号量")# 启动10个线程,只有3个能同时执行threads = []for i in range(10): t = threading.Thread(target=limited_worker, args=(f"线程{i}",)) threads.append(t) t.start()for t in threads: t.join()
Event 事件
一个线程等待另一个线程发出信号。
import threadingimport timeevent = threading.Event()data = Nonedefproducer():"""生产者:准备数据后发出信号"""global data print("生产者: 正在准备数据...") time.sleep(3) data = "重要的数据" print("生产者: 数据已准备好,发出通知") event.set() # 发出信号defconsumer():"""消费者:等待信号""" print("消费者: 等待数据...") event.wait() # 阻塞等待信号 print(f"消费者: 收到数据 {data}")# 启动t1 = threading.Thread(target=consumer) # 先启动消费者t2 = threading.Thread(target=producer)t1.start()t2.start()t1.join()t2.join()
Event 的方法:
| |
|---|
set() | |
wait() | |
clear() | |
is_set() | |
Condition 条件变量
更复杂的线程协调:等待条件满足,然后通知其他线程。
import threadingimport timefrom collections import deque # 双端队列classThreadSafeQueue:"""线程安全的队列(使用 Condition)"""def__init__(self, max_size=5): self.queue = deque() self.max_size = max_size self.condition = threading.Condition()defput(self, item):with self.condition:# 队列满了就等待while len(self.queue) >= self.max_size: print(f"队列已满,生产者等待...") self.condition.wait() self.queue.append(item) print(f"生产: {item},队列大小: {len(self.queue)}") self.condition.notify() # 通知消费者defget(self):with self.condition:# 队列空了就等待while len(self.queue) == 0: print(f"队列为空,消费者等待...") self.condition.wait() item = self.queue.popleft() print(f"消费: {item},队列大小: {len(self.queue)}") self.condition.notify() # 通知生产者return item# 测试queue = ThreadSafeQueue(max_size=3)defproducer():for i in range(10): queue.put(f"Item-{i}") time.sleep(0.5)defconsumer():for _ in range(10): item = queue.get() time.sleep(1)# 启动t1 = threading.Thread(target=producer)t2 = threading.Thread(target=consumer)t1.start()t2.start()t1.join()t2.join()
threading.Condition() :带"等待室"的锁,它允许线程在某个条件不满足时主动"等待",并在条件满足时被"唤醒"。self.queue 一个双端队列,用于存放生产者产生的数据。put() 方法将生产者产生的数据放入队列中。如果队列满了则生产者等待。get() 方法用于让消费者从队列中获取数据。如果队列为空则消费者等待。
使用 queue.Queue 线程安全队列进行多线程通信
queue.Queue 线程安全的先进先出 FIFO 队列,专门用于在多线程之间安全地传递数据。其内部已经实现了所有必要的锁机制,多个线程可以安全地 put() 和 get() 。
其常用方法有:
| | |
|---|
put(item) | | |
put_nowait(item) | | |
get() | | |
get_nowait() | | |
task_done() | | |
join() | | 阻塞直到 task_done() 被调用的次数等于 get() 的次数 |
qsize() | | |
empty() | | |
full() | | |
用于生产者-消费者模型:

queue.Queue.join() 工作流程图:

import threadingimport timefrom queue import Queue# 创建队列task_queue = Queue(maxsize=10)defproducer():"""生产者:产生任务"""for i in range(20): task_queue.put(f"任务-{i}") print(f"生产: 任务-{i}") time.sleep(0.2)# 发送结束信号for _ in range(3): task_queue.put(None) print("生产者完成")defconsumer(name):"""消费者:处理任务"""whileTrue: task = task_queue.get()if task isNone: task_queue.task_done() print(f"消费者 {name} 退出")break print(f"消费者 {name} 处理: {task}") time.sleep(0.5) task_queue.task_done() # 标记任务完成,每一个 get() 必须对应一个 task_done()producer_thread = threading.Thread(target=producer)consumers = []for i in range(3): t = threading.Thread(target=consumer, args=(i, )) consumers.append(t) t.start()producer_thread.start()task_queue.join() # 阻塞,直到所有的 task_done() 被调用for t in consumers: t.join()producer_thread.join()print("所有任务完成")
需要注意的是,这种生产消费者模型适用于 IO 密集型任务,利用多线程可以充分利用 CPU 时间。Python 的多线程是并发,而不是并行,它是靠时间片轮转来实现的,用于 CPU 密集型任务可能反而降低效率。

线程池
手动创建大量线程管理起来特别麻烦,Python 有一个内置的高级 API 可以让线程管理特别方便简洁,这个 API 就是线程池 concurrent.futures.ThreadPoolExecutor 。
线程池核心目的:预先创建一批线程复用,限制最大并发数,统一调度任务。
concurrent.futures.ThreadPoolExecutor 底层实现原理就是 queue.Queue + threading 。
导入:
from concurrent.futures import ThreadPoolExecutorimport time
关键参数:
ThreadPoolExecutor(max_workers=8)
max_workers = 8 最大并发线程数量。
即便是 IO 密集型任务,也不是线程越多越好,超过阈值后操作系统线程切换开销会抵消收益。
提交任务
单个任务使用 submit 提交,返回 future 对象:
deftask(name): time.sleep(1) # 模拟IO等待returnf"任务 {name} 完成"if __name__ == "__main__":with ThreadPoolExecutor(max_workers=3) as pool: future = pool.submit(task, "A")# 获取结果,阻塞等待任务完成 res = future.result() print(res)
future.result(timeout=5):获取返回值,可设置超时future.exception() 获取任务抛出的异常
多任务使用 map 批量执行,输入顺序和输出顺序保持一致:
deftask(x): time.sleep(0.5)return x * 2if __name__ == "__main__":with ThreadPoolExecutor(max_workers=3) as pool: data = [1,2,3,4,5,6] results = pool.map(task, data)for r in results: print(r)
在需要保证响应结果有序的网络请求中,可以用这种方式。
with ThreadPoolExecutor() as pool: 会在代码块结束后,自动调用 shutdown(wait=True) ,等待所有线程执行完毕,再释放资源。
也可以手动调用 pool.shutdown(wait=True) 关闭线程,wait=True:主线程阻塞,等所有任务跑完;wait=False:不等待,线程池进入关闭状态,不再接收新任务,主线程直接往下走。
from concurrent.futures import ThreadPoolExecutor, as_completedimport timedeftask(name, delay):"""模拟任务""" print(f"任务 {name} 开始") time.sleep(delay) print(f"任务 {name} 完成")returnf"结果-{name}"# ==================== 方式1:使用 map(保持顺序) ====================defwith_map(): print("\n" + "="*50) print("方式1: executor.map() - 保持提交顺序") print("="*50)with ThreadPoolExecutor(max_workers=3) as executor:# 提交多个任务,结果按提交顺序返回 results = executor.map(task, ['A', 'B', 'C'], [2, 1, 3])for result in results: print(result)# ==================== 方式2:使用 submit + as_completed(按完成顺序) ====================defwith_submit(): print("\n" + "="*50) print("方式2: submit() + as_completed() - 按完成顺序") print("="*50)with ThreadPoolExecutor(max_workers=3) as executor:# 提交任务,获得 Future 对象 futures = { executor.submit(task, name, delay): namefor name, delay in [('A', 2), ('B', 1), ('C', 3)] } # 字典推导式,key 是 Future 对象,value 是 name# 按完成顺序获取结果,as_completed 按完成顺序返回 Futurefor future in as_completed(futures): name = futures[future]try: result = future.result() print(f"{name} 的结果: {result}")except Exception as e: print(f"{name} 出错了: {e}")# ==================== 方式3:错误处理 ====================deferror_task(x):if x < 0:raise ValueError(f"负数: {x}")return x ** 2defwith_error_handling(): print("\n" + "="*50) print("方式3: 错误处理") print("="*50)with ThreadPoolExecutor(max_workers=2) as executor: futures = [executor.submit(error_task, i) for i in [-1, 2, -3, 4]]for future in futures:try: result = future.result() print(f"结果: {result}")except ValueError as e: print(f"出错: {e}")if __name__ == "__main__": with_map() with_submit() with_error_handling()
多进程编程
想要使用 Python 真正实现并行,就需要通过多进程来实现。
我们平时所说的 CPU 多核多线程,比如 8核 12线程,这里的核指的是物理核心数,12 线程利用超线程技术在一个物理核心上实现多个逻辑线程,CPU 这里所说的线程就是可以单独跑操作系统进程的逻辑单元。


使用多线程+多进程可以充分发挥 CPU 的性能。
多进程编程是非常灵活的,它是语言无关的,任何语言都可以创建进程;每个进程都是互相隔离的,一个进程崩溃不影响其他进程;每个进程有独立的内存空间。
同时跑多个 Python 脚本本身就是多进程,或者多种编程语言写的程序同时跑,也是多进程。
进程通信可以是 socket 网络套接字,多个进程在 TCP/IP 协议层进行通信。
C/C++、C#、Java、Python 等都可以利用 socket 编程实现 server-client 模型进行进程通信。这些进程既可以在同一台电脑上,通过本机的网络层进行通信,也可以在多台计算机上同时跑。比如利用 Python 进行人体图像动作捕捉,使用 socket 编程将捕捉到的坐标数据传输给 C#,在 Unity 的 3D 模型中做动作映射,Python 的人体动作坐标捕捉是二维的,映射到 3D 的模型中从正面看完全没问题,转到侧面就会偏离很大。



除了同时跑多个 Python 脚本这种简单粗暴的多进程方法,Python 也提供了多进程编程的 API —— multiprocessing 库。
Python 的多进程可以真正地跑 CPU 密集型任务。
基本用法:
import multiprocessingimport timedefwork(num):"""子进程执行函数""" print(f"子进程 {num} 开始,pid={multiprocessing.current_process().pid}") time.sleep(2) print(f"子进程 {num} 结束")if __name__ == "__main__":# Windows / macOS 必须写 if __name__ == "__main__" 防止进程递归创建 p1 = multiprocessing.Process(target=work, args=(1,)) p2 = multiprocessing.Process(target=work, args=(2,)) p1.start() # 启动进程 p2.start() p1.join() # 等待子进程执行完毕,主进程阻塞 p2.join() print("全部子进程完成")
.daemon=True:守护进程,主进程退出,子进程直接被杀掉
进程池
大量任务不要手动创建一堆 Process ,进程池复用进程,减少创建销毁开销。
from concurrent.futures import ProcessPoolExecutorimport timedefcalc(x):return x*xif __name__ == "__main__":with ProcessPoolExecutor(max_workers=4) as pool: res = pool.map(calc, [1, 2, 3, 4, 5, 6]) print(list(res))
max_workers:最大并发进程数,一般设置等于 CPU 物理核心数。
进程通信
进程内存隔离,普通全局变量在子进程是拷贝,修改不会影响主进程,需要特殊对象 Queue / Pipe / Manager 通信,不能用普通全局变量。
方式1:使用进程 Queue 队列通信:
import multiprocessingdefproducer(q): q.put("hello from child")defconsumer(q): msg = q.get() print(msg)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() p2.join()
方式2:使用 Pipe 管道,仅适合两个进程通信:
import multiprocessingdefchild(conn): conn.send({"name": "jackey", "num": 666}) msg = conn.recv() print(f"子进程收到: {msg}") conn.close()if __name__ == "__main__": parent_conn, child_conn = multiprocessing.Pipe() p = multiprocessing.Process(target=child, args=(child_conn,)) p.start()# 主进程接收 data = parent_conn.recv() print(f"主进程收到:{data}") parent_conn.send("收到你的消息啦") p.join() parent_conn.close()
multiprocessing.Pipe() 返回两个端点。
Pipe 是双向的。两端都可以 send / recv;recv() 会阻塞,直到收到消息。
方式3:Manager 共享字典 / 列表,Manager 启动一个额外管理子进程,所有共享的 dict/list 实际存在这个管理进程里,其他进程通过 IPC 访问,性能低,适合少量数据,不要高频循环读写。
import multiprocessingdefworker(shared_dict, shared_list):# 子进程修改共享对象 shared_dict["count"] = 100 shared_dict["name"] = "test" shared_list.append(999)if __name__ == "__main__": manager = multiprocessing.Manager() shared_dict = manager.dict() # 共享字典 shared_list = manager.list() # 共享列表 p = multiprocessing.Process(target=worker, args=(shared_dict, shared_list)) p.start() p.join()# 主进程可以读到子进程修改后的结果 print("shared_dict:", dict(shared_dict)) print("shared_list:", list(shared_list))
Python 异步编程
异步编程是一种在单线程内,通过 "事件循环 + 协程" 实现高并发 IO 的编程模型。
核心思想是当一个任务在等待 IO 时,主动让出 CPU,让其他任务执行。等到 IO 完成时,再回来继续执行。
Python 的异步与 node.js 的异步核心思想都是 事件循环 + 非阻塞 IO,语法层面也高度相似,但实现方式、细节差异比较大。
| | | |
|---|
| 多线程 | | | |
| 多进程 | | | |
| 异步IO | | 程序员主动让出 | 高并发IO(万级连接) |
异步编程核心机制 —— 事件循环:

协程 Coroutine
协程是异步编程的执行单元。它可以在执行过程中暂停(await),让出控制权交给事件循环。
import asyncio# 协程函数:使用 async def 定义asyncdefmy_coroutine(): print("开始执行")await asyncio.sleep(1) # 主动让出 CPU, 等待 1 秒 print("恢复执行")# 调用协程函数,不会立即执行,只返回协程对象coro = my_coroutine()print(type(coro)) # <class 'coroutine'>asyncio.run(coro) # 必须通过事件循环执行
同步对比异步:
import timeimport asyncio# ================== 同步方式defsync_task(name, delay): print(f"任务 {name} 开始") time.sleep(delay) # 阻塞,CPU 发呆 print(f"任务 {name} 完成")return namedefsync_main(): start = time.time()for i in range(3): sync_task(f"任务 {i}", 1) print(f"总耗时:{time.time() - start:.2f}s")# ================== 异步方式asyncdefasync_task(name, delay): print(f"任务 {name} 开始")await asyncio.sleep(delay) # 非阻塞,让出 CPU print(f"任务 {name} 完成")return nameasyncdefasync_main(): start = time.time()# 并发执行 3 个任务 tasks = [async_task(f"任务{i}", 1) for i in range(3)] results = await asyncio.gather(*tasks) # *tasks 列表解包操作符,asyncio.gather() 并发执行多个协程任务# 相当于 results = await asyncio.gather(task1, task2, task3)# 返回结果的顺序与传入顺序一致,这是 gather 的一个重要特性 print(f"总耗时: {time.time() - start:.2f}s")return results# 运行sync_main() # 3.00sasyncio.run(async_main()) # 1.00s
await等待可等待对象,只能用在 async 定义的函数中。可等待对象 Awaitable 有:
协程对象 coroutine:调用 async def fn() 得到。
Task 对象:asyncio.create_task() 创建,后台运行的任务。
import asyncioasyncdeffoo():await asyncio.sleep(1)return66asyncdefmain(): task = asyncio.create_task(foo()) # Task 对象 res = await task # await Task,不是 await foo() print(res)asyncio.run(main())
Future 对象:表示一个尚未完成的异步操作的结果占位符。它像一个"空盒子",承诺未来某个时刻里面会放一个值。我们可以通过 await future 等待结果,也可以检查 future.done() 判断是否完成。
import asyncioasyncdefmain(): loop = asyncio.get_running_loop() fut = loop.create_future() fut.set_result(888) print(await fut) # await Future 对象asyncio.run(main())
并发异步爬虫小案例:下面给出一个并发爬取我的网络博客上一个视频的代码。在网络上的很多视频都是分块传输的,一个完整的视频会被划分成很多块。如果我们要获取完整的视频,需要把所有的视频块全部请求。同步爬取需要一个个视频块按顺序请求,串行执行。而并发获取可以同时发起所有请求,但响应返回的顺序是乱的(谁先响应谁先回来),因此需要在全部下载完成后按索引排序,再顺序组装写入文件。
异步模式:
import aiohttp # pip install aiohttp 异步 http 模块import asyncioimport timefrom typing import List, Tuple, Optionalbase_url = 'https://v-blog.csdnimg.cn/asset/999efa6d97215aa8905a1a05f7398e9f/play_video/'video_chunk_list = ['32b018315e66b2f02a2c08433b42fcc0_0.ts','32b018315e66b2f02a2c08433b42fcc0_1.ts','32b018315e66b2f02a2c08433b42fcc0_2.ts','32b018315e66b2f02a2c08433b42fcc0_3.ts','32b018315e66b2f02a2c08433b42fcc0_4.ts']asyncdeffetch_chunk_with_retry( session: aiohttp.ClientSession, base_url: str, chunk_name: str, index: int, max_retries: int = 3) -> Tuple[int, Optional[bytes]]:""" 带重试机制的异步下载 """ url = base_url + chunk_namefor attempt in range(max_retries):try:asyncwith session.get(url, timeout=30) as response:if response.status == 200: content = await response.read() print(f"✅ 分片 {index} 下载完成 (尝试 {attempt + 1})")return index, contentelse: print(f"⚠️ 分片 {index} 状态码 {response.status}")except Exception as e: print(f"❌ 分片 {index} 下载失败 (尝试 {attempt + 1}): {e}")if attempt == max_retries - 1: print(f"💀 分片 {index} 所有重试失败")return index, Noneasyncdefasync_download_video_advanced( base_url: str, chunk_list: List[str], max_concurrent: int = 10) -> bool:""" 高级异步下载器:支持并发控制、重试、进度显示 """ total = len(chunk_list) chunks = [None] * total# 信号量:控制最大并发数 semaphore = asyncio.Semaphore(max_concurrent)asyncdeffetch_with_semaphore(session, chunk_name, idx):asyncwith semaphore:returnawait fetch_chunk_with_retry(session, base_url, chunk_name, idx)asyncwith aiohttp.ClientSession() as session:# 创建所有任务 tasks = [ fetch_with_semaphore(session, chunk_name, idx)for idx, chunk_name in enumerate(chunk_list) ]# 进度追踪 completed = 0for coro in asyncio.as_completed(tasks): idx, content = await coro chunks[idx] = content completed += 1 progress = (completed / total) * 100 bar_len = 30 filled = int(bar_len * completed // total) bar = '█' * filled + '░' * (bar_len - filled) print(f'\r📥 下载进度: [{bar}] {progress:.1f}%', end='', flush=True) print() # 换行# 检查是否有失败的分片 failed = [idx for idx, data in enumerate(chunks) if data isNone]if failed: print(f"❌ 以下分片下载失败: {failed}")returnFalse# 按顺序写入文件with open('async_get_video.ts', 'wb') as f:for chunk_data in chunks:if chunk_data isnotNone: f.write(chunk_data) print("✅ 所有分片已按顺序合并写入")returnTrueasyncdefmain(): print("🚀 异步下载器启动") print(f"📊 分片总数: {len(video_chunk_list)}") print(f"🔄 最大并发数: 10") print("-" * 60) start = time.time() success = await async_download_video_advanced(base_url, video_chunk_list) elapsed = time.time() - startif success: print(f"✅ 下载完成!总耗时: {elapsed:.2f}s")else: print("❌ 下载失败")if __name__ == '__main__': asyncio.run(main())
运行:
🚀 异步下载器启动📊 分片总数: 5🔄 最大并发数: 10------------------------------------------------------------✅ 分片 4 下载完成 (尝试 1)📥 下载进度: [██████░░░░░░░░░░░░░░░░░░░░░░░░] 20.0%✅ 分片 2 下载完成 (尝试 1)📥 下载进度: [████████████░░░░░░░░░░░░░░░░░░] 40.0%✅ 分片 3 下载完成 (尝试 1)📥 下载进度: [██████████████████░░░░░░░░░░░░] 60.0%✅ 分片 1 下载完成 (尝试 1)📥 下载进度: [████████████████████████░░░░░░] 80.0%✅ 分片 0 下载完成 (尝试 1)📥 下载进度: [██████████████████████████████] 100.0%✅ 所有分片已按顺序合并写入✅ 下载完成!总耗时: 1.53s
对比同步版本:
import requests # pip install requests 同步 http 请求模块 import timebase_url = 'https://v-blog.csdnimg.cn/asset/999efa6d97215aa8905a1a05f7398e9f/play_video/'video_chunk_list = ['32b018315e66b2f02a2c08433b42fcc0_0.ts','32b018315e66b2f02a2c08433b42fcc0_1.ts','32b018315e66b2f02a2c08433b42fcc0_2.ts','32b018315e66b2f02a2c08433b42fcc0_3.ts','32b018315e66b2f02a2c08433b42fcc0_4.ts']defsync_requests_video(base_url, video_chunk_list):for chunk in video_chunk_list: response = requests.get(base_url + chunk)with open('sync_get_video.ts', 'ab') as file_object: file_object.write(response.content) print("同步模式爬取完毕")start = time.time()sync_requests_video(base_url, video_chunk_list)print(f"同步爬取总耗时: {time.time() - start:.2f}s")
运行:
同步模式爬取完毕同步爬取总耗时: 3.86s
从生成器到协程 —— await/async 的进化之路
Python 的生成器有 send() 、 throw() 、 close() 方法,这意味着生成器本质上已经是一种协程,它可以暂停、恢复、接收外部数据。
生成器通过 yield 实现了控制流的主动让出,这正是协程的核心特征。
Python生成器作为原始协程:
defsimple_coroutine():"""这是一个基于生成器的协程""" print("协程启动") x = yield"第一次暂停"# 接收外部数据 print(f"收到外部数据: {x}") y = yield"第二次暂停" print(f"再次收到: {y}")return"协程结束"# 使用协程coro = simple_coroutine()# 1. 启动协程(必须执行到第一个yield)print(next(coro)) # 输出: 协程启动 \n 第一次暂停# 2. 发送数据给协程print(coro.send(100)) # 输出: 收到外部数据: 100 \n 第二次暂停# 3. 再次发送并捕获结束try: coro.send(200)except StopIteration as e: print(e.value) # 输出: 协程结束
历史上 Python3.4 使用装饰器 @asyncio.coroutine 和 yield from 来实现异步,但这种方式也有两个痛点:
- 语义混淆:
yield 既用于迭代又用于协程,让人困惑
Python 3.5 引入了 async 和 await 关键字,创建了原生协程(Native Coroutine):
# Python 3.4 风格(已过时)@asyncio.coroutinedefold_style_coro(): result = yieldfrom some_future()return result# Python 3.5+ 风格(现代标准)asyncdefmodern_coro(): result = await some_async_function()return result
异步上下文管理器(async with)
资源需要异步初始化和清理时使用(如异步数据库连接、异步文件)。
import aiofilesasyncdefread_file_async():"""异步读取文件"""# async with 管理异步上下文asyncwith aiofiles.open('data.txt', 'r') as f: content = await f.read()return content
异步迭代器
import asyncioclassAsyncCounter:"""异步计数器迭代器"""def__init__(self, max_count): self.max_count = max_count self.count = 0def__aiter__(self):return selfasyncdef__anext__(self):if self.count >= self.max_count:raise StopAsyncIterationawait asyncio.sleep(0.5) self.count += 1return self.countasyncdefmain():# async for 异步遍历asyncfor num in AsyncCounter(5): print(f"异步迭代: {num}")asyncio.run(main())
Python 的并发与异步是一个从入门到精通的分水岭。掌握多线程让你能高效处理 IO 任务;掌握多进程让你能真正利用多核 CPU;而掌握异步编程,则让你拥有处理万级并发的能力。
这三种工具没有绝对的好坏——多线程适合中低并发 IO,多进程适合 CPU 密集型计算,异步适合高并发网络服务。关键在于:在正确的场景选择正确的工具。
如果你认真读到了这里,我相信你已经迈过了这道分水岭。
如果这篇文章对你有帮助,欢迎点赞、收藏、评论!你的支持是我持续输出的动力。
END.