python 多进程、多线程,异步
介绍
1.1 QA
1.2 mian函数的用法
| 多线程 | 多进程 | |
|---|---|---|
是否需要if __name__ == '__main__' |
不需要 | ✅在 Windows/macOS 上必须加(Linux 可不加但强烈建议加) |
| 根本原因 | 共享同一进程内存,直接启动新线程 | 需要创建新进程,新进程会重新导入主模块,导致无限递归创建进程 |
1. 多进程
queue = multiprocessing.Queue()
无参数函数启动如何启动
Python 多进程是 利用操作系统的多进程机制,并行执行多个独立任务 的编程方式,核心解决 CPU 密集型任务 的性能瓶颈(如计算、数据分析),也能规避 Python 全局解释器锁(GIL)对 CPU 并行的限制。
Python 中的 GIL(Global Interpreter Lock)会导致 同一时刻只有一个线程执行 Python 字节码,因此:
- IO 密集型(等待时间 > 计算时间):优先用多线程(ThreadPoolExecutor),或更高效的异步(asyncio);
- CPU 密集型(计算时间 > 等待时间):优先用多进程(ProcessPoolExecutor/multiprocessing);
多进程的本质是:操作系统为每个进程分配独立的内存空间、Python 解释器和 GIL,进程间完全隔离,可真正并行执行。
1.1 启动方式
import multiprocessing
import time
def task(name, sleep_time):
"""子进程执行的任务"""
print(f"子进程 {name} 启动,休眠 {sleep_time} 秒")
time.sleep(sleep_time)
print(f"任务 {name} {sleep_time} 秒 完成.................")
if __name__ == '__main__':
start_time = time.time()
# 1. 创建进程实例(target指定函数,args指定参数元组)
p1 = multiprocessing.Process(target=task, args=("P1", 2))
p2 = multiprocessing.Process(target=task, args=("P2", 1))
# 2. 启动进程
p1.start()
p2.start()
# 3. 等待子进程结束(主进程阻塞)
p1.join()
p2.join()
cost_time = time.time() - start_time
print(f"所有子进程执行完毕,主进程退出{cost_time:.2f}S")
1.2 进程池
multiprocessing.Pool(以下简称 Pool)和 concurrent.futures.ProcessPoolExecutor(以下简称 ProcessPoolExecutor)都是 Python 中用于管理多进程池的工具,核心目的是简化批量并行任务的开发、自动调度进程资源,但二者在 设计定位、API 风格、功能特性、使用场景 上有显著区别。
核心结论先明确:
- Pool 是 multiprocessing 模块的核心组件,设计偏向 “底层控制”,提供了更多进程池管理的细节方法(如 close()、terminate()、join()),API 风格更贴近原生多进程操作,适合需要灵活调整进程池状态的场景。
- ProcessPoolExecutor 是 Python 3.2+ 引入的 concurrent.futures 模块的组件,设计目标是 “简化并行编程”—— 它统一了进程池和线程池的 API(与 ThreadPoolExecutor 接口完全一致),无需关心底层是进程还是线程,切换成本极低,适合快速开发。(向多线程对齐)
| 对比维度 | multiprocessing.Pool | concurrent.futures.ProcessPoolExecutor (推荐使用) |
|---|---|---|
| 设计定位 | 传统多进程池,功能全面、控制粒度细 | 高层封装,简化并行编程,与线程池 API 统一 |
| API 风格 | 类原生接口(如 apply/apply_async/map`) |
现代上下文管理器风格,接口简洁(map<br/>/submit) |
| 任务提交方式 | 支持同步(apply)、异步(apply_async)、批量(map`) |
支持批量(map)、单个异步(submit+Future`) |
| 结果获取 | map直接返回列表,apply_async需 get() |
map返回迭代器,submit返回 Future对象(result()/add_done_callback()) |
| 异常处理 | 需手动捕获,get()`会抛出原始异常 | Future统一封装异常,result()`抛出,支持回调处理 |
| 进程池复用 | 支持(创建后可多次提交任务,需手动 close()/join()) |
不支持(上下文管理器自动关闭,单次使用) |
| 超时控制 | apply_async/map_async支持 timeout参数 |
Future.result(timeout)支持超时,map`无直接超时(需手动处理) |
| 回调函数 | apply_async支持 callback`参数(仅成功回调) |
add_done_callback()支持回调(成功 / 失败均触发) |
| 取消任务 | 不支持(任务提交后无法取消) | 支持(Future.cancel()`,未执行的任务可取消) |
| 适用场景 | 复杂进程管理、多次复用进程池、细粒度控制 | 简单并行任务、与线程池切换、快速开发 |
1.2.1 进程池(multiprocessing.Pool 传统方式)
- 同步一般使用map, 异步一般使用apply_async
- map: 主进程完全阻塞,直到全部任务完成
- apply_async : 非阻塞提交 + 按需阻塞获取
- close()和 join()的关系
- close() 和 join() 是使用 multiprocessing.Pool 时必须成对出现的关键方法,它们协同工作来安全、干净地关闭进程池。理解它们的关系,是避免僵尸进程、资源泄漏和程序卡死的关键!
| 方法 | 说明 |
|---|---|
pool.map(func, iterable) |
同步:将func应用于iterable中每个元素,阻塞直到全部完成,返回结果列表(顺序与输入一致) |
pool.map_async(func, iterable, callback=None) |
异步版本:立即返回AsyncResult对象,可用.get()获取结果 |
pool.apply(func, args=(), kwds={}) |
同步:提交单个任务(类似func(*args, **kwds)),阻塞等待结果 |
pool.apply_async(func, args=(), kwds={}, callback=None, error_callback=None) |
异步:提交单个任务,不阻塞;返回AsyncResult |
pool.starmap(func, iterable) |
类似map,但iterable中每个元素是参数元组(解包传参),如func(*item) |
pool.close() |
关闭池:不再接受新任务 |
pool.join() |
等待所有 worker 进程结束(必须先** **close()** 或 ****terminate()**) |
pool.terminate() |
立即终止所有 worker 进程(不等任务完成) |
1.2.1.1 map (最常用)
会阻塞当前线程,需要自己手动关闭
- 不支持多参数函数, 只支持单参数
示例代码
import multiprocessing
import time
def task(sleep_time):
"""子进程任务:接收单个参数,返回结果"""
print(f"进程 {multiprocessing.current_process().name} 启动,休眠 {sleep_time} 秒")
time.sleep(sleep_time)
print(f"进程 {multiprocessing.current_process().name} 完成,休眠 {sleep_time} 秒")
return sleep_time
if __name__ == '__main__':
start_time = time.time()
# 1. 创建进程池(size=CPU核心数)
pool = multiprocessing.Pool(processes=4)
# 2. 批量提交任务(map自动分配任务,参数为可迭代对象)
tasks = [2, 1, 3] # 每个元素作为task的参数
# map是 “阻塞型调用” 本身会阻塞主进程,直到所有任务完成并返回结果
results = pool.map(task, tasks)
print(f"程序运行中 1..........................................")
# 3. 关闭进程池+等待任务完成
pool.close() # 禁止新任务提交, 不影响已提交任务 close() 是 join() 的 “前置操作”,只有先调用 close(),再调用 pool.join(),
pool.join() # 主进程才会阻塞等待所有任务完成(如果不调用 close() 直接 join(),会报错)。
print(f"程序运行中 2..........................................")
cost_time = time.time() - start_time
# 4. 打印结果
for res in results:
print(res)
print(f"所有任务结果:运行耗时{cost_time:.3f}s")
1.2.1.2 apply_async(异步单任务)
- 非阻塞提交
示例代码:
import multiprocessing
import time
def task(name, sleep_time):
print(f"进程 {name} 启动,休眠 {sleep_time} 秒")
time.sleep(sleep_time)
return f"进程 {name} 完成"
# multiprocessing.Pool 支持 apply_async, map 不支持 submit
# apply_async 中join() 之前一定要close, close之后线程池中不能再添加任务
if __name__ == '__main__':
start_time = time.time()
pool = multiprocessing.Pool(processes=4)
results = []
# 异步提交多个任务(多参数用args元组)
results.append(pool.apply_async(task, args=("P1", 2)))
results.append(pool.apply_async(task, args=("P2", 1)))
results.append(pool.apply_async(task, args=("P3", 3)))
print(f"程序运行中 1..........................................")
pool.close() # 禁止新任务提交, 不影响已提交任务 close() 是 join() 的 “前置操作”,只有先调用 close(),再调用 pool.join(),
pool.join() # 这个有区别 阻塞当前所在的主进程,
print(f"程序运行中 2..........................................")
# 获取结果(get()会阻塞直到任务完成)
for res in results:
print(res.get())
# pool.apply_async(task, args=("P4", 4))
cost_time = time.time() - start_time
print(f"程序运行结束.........{cost_time:.2f}s")
1.2.1.3 map VS apply_async
| 特性 | map | apply_async |
|---|---|---|
| 任务模式 | ✅批量、同构任务(对 iterable 中每个元素执行相同函数 | ✅单个、异构任务**(可提交不同函数/不同参数) |
| 同步/异步 | ❌ 同步(阻塞,直到所有任务完成) | ✅ 异步(立即返回AsyncResult对象) |
| 输入结构 | 一个函数 + 一个可迭代对象(如[1,2,3]) |
一个函数 + 参数元组(如(10,),kwds={'a':1}) |
| 结果顺序 | ✅ 严格保持输入顺序([f(x) for x in xs]) |
❌ 按完成顺序获取(需.get()显式取值) |
| **适用场景 | 把一批数据都做同样处理”(如:全部开方、全部下载) | 提交多个不同任务,谁先完成谁先处理”(如:多个独立 API 请求) |
示例代码
import multiprocessing
import time
def taskA(sleep_time):
"""子进程任务:接收单个参数,返回结果"""
print(f"进程 {multiprocessing.current_process().name} 启动,休眠 {sleep_time} 秒")
time.sleep(sleep_time)
print(f"进程 {multiprocessing.current_process().name} 完成,休眠 {sleep_time} 秒")
return sleep_time
def taskB(name, sleep_time):
print(f"[{time.strftime('%H:%M:%S')}] 进程 {name} 启动,休眠 {sleep_time} 秒")
time.sleep(sleep_time)
return f"[{time.strftime('%H:%M:%S')}] 进程 {name} 完成"
def test_mulitprocess_map(): # 线程池中的任务通常只能是一个
print(f"test_mulitprocess_map----------------------------------------")
# 1. 创建进程池(size=CPU核心数)
pool = multiprocessing.Pool(processes=4)
# 2. 批量提交任务(map自动分配任务,参数为可迭代对象)
params = [2, 1, 3] # 每个元素作为task的参数
# map是 “阻塞型调用” 本身会阻塞主进程,直到所有任务完成并返回结果
results = pool.map(taskA, params)
print(f"程序运行中 1..........................................")
# 3. 关闭进程池+等待任务完成
pool.close() # 禁止新任务提交, 不影响已提交任务 close() 是 join() 的 “前置操作”,只有先调用 close(),再调用 pool.join(),
pool.join() # 主进程才会阻塞等待所有任务完成(如果不调用 close() 直接 join(),会报错)。
print(f"程序运行中 2..........................................")
# 4. 打印结果
for res in results:
print(res)
def test_multiprocess_submit(): # 线程池中的任务可以不同,例如A, B, C, D
print(f"test_multiprocess_submit----------------------------------------")
pool = multiprocessing.Pool(processes=4)
results = []
# 异步提交多个任务(多参数用args元组)
results.append(pool.apply_async(taskA, args=(2, )))
results.append(pool.apply_async(taskB, args=("P2", 1)))
results.append(pool.apply_async(taskB, args=("P3", 3)))
print(f"程序运行中 1..........................................")
pool.close()
pool.join()
print(f"程序运行中 2..........................................")
# 获取结果(get()会阻塞直到任务完成)
for res in results:
print(res.get())
"""
多进程_map, 会阻塞当前线程,需要自己手动关闭
multiprocessing.Pool.map() 不支持多参数函数, 只支持单参数/无参数函数
"""
if __name__ == '__main__':
start_time = time.time()
# test_mulitprocess_map()
test_multiprocess_submit()
cost_time = time.time() - start_time
print(f"所有任务结果:运行耗时{cost_time:.3f}s")
1.2.2 进程池(concurrent.futures.ProcessPoolExecutor)
map: 相同任务多次调用
submit: 可以是相同任务,也可以是多个不同任务,一次只能提交一个任务
from concurrent.futures import ProcessPoolExecutor
# 推荐:用 with(自动 shutdown ✅)
with ProcessPoolExecutor(max_workers=4) as executor:
# 提交任务...
pass
# 或手动管理(不推荐,易漏资源)
executor = ProcessPoolExecutor(max_workers=4)
# ... 使用 ...
executor.shutdown(wait=True) # 等价于 pool.close() + pool.join()
| 方法 | 作用 | 返回类型 | 阻塞 |
|---|---|---|---|
submit(fn, *args, **kwargs) |
提交单个任务 | Future | ❌ 非阻塞 |
map(fn, *iterables, timeout=None, chunksize=1) |
批量提交,类似Pool.map | 迭代器(结果按输入顺序) | ⚠️调用时不阻塞,但迭代时阻塞 |
1.2.2.1 submit
- 一般是多任务,多参数
- as_completed 按照完成顺序返回
通常和concurrent.futures.as_completed(futures)搭配使用
import concurrent.futures
import time
def task(name, sleep_time):
print(f"[{time.strftime('%H:%M:%S')}] 进程 {name} 启动,休眠 {sleep_time} 秒")
time.sleep(sleep_time)
return f"[{time.strftime('%H:%M:%S')}] 进程 {name} 完成"
# concurrent.futures.ProcessPoolExecutor(推荐) 支持 submit, map 不支持apply_async
if __name__ == '__main__':
start_time = time.time()
# 用with语句创建进程池,自动关闭
with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor:
# 提交任务,返回Future对象列表
futures = [
executor.submit(task, "P1", 2),
executor.submit(task, "P2", 1),
executor.submit(task, "P3", 3)
]
print(f"程序运行中 1..........................................")
# 遍历Future对象,获取结果(as_completed()按完成顺序返回)
# for future in concurrent.futures.as_completed(futures): # as_completed按完成结果顺序结果返回,而不是futures的顺序
for future in futures: # as_completed按完成结果顺序结果返回,而不是futures的顺序
print(future.result())
print(f"程序运行中 2..........................................")
cost_time = time.time() - start_time
print(f"程序运行结束.........{cost_time:.2f}s")
1.2.1.2 map
- 一般是单任务
- 支持多参数
示例代码:
import concurrent.futures
import time
def task(name, sleep_time):
print(f"[{time.strftime('%H:%M:%S')}] 进程 {name} 启动,休眠 {sleep_time} 秒")
time.sleep(sleep_time)
return f"[{time.strftime('%H:%M:%S')}] 进程 {name} 完成"
if __name__ == '__main__':
start_time = time.time()
with concurrent.futures.ProcessPoolExecutor(max_workers=4) as executor:
names = ["P1", "P2", "P3"]
sleep_times = [2, 1, 3]
# map返回结果的顺序与任务提交顺序一致
results = executor.map(task, names, sleep_times) # 按照任务提交顺序返回 # 这个有区别 阻塞当前所在的主进程
# 2. 这行代码会正常执行!因为 map() 未阻塞,迭代器已返回
print(f"程序运行中 1..........................")
for res in results: # 关键:迭代 results 时,才会开始阻塞!
print(res)
print(f"程序运行中 2..........................................")
cost_time = time.time() - start_time
print(f"程序运行结束.........{cost_time:.2f}s")
1.3 总的测试代码
import multiprocessing
import time
# 定义单参数任务函数(接收一个包含多个参数的元组)
def task(args):
name, sleep_time = args
print(f"进程 {name} 启动,休眠 {sleep_time} 秒(进程ID:{multiprocessing.current_process().pid})")
time.sleep(sleep_time)
return f"进程 {name} 完成"
def test_processpool_map():
# 1. 创建进程池(指定4个进程,与CPU核心数无关,可自定义)
pool = multiprocessing.Pool(processes=4)
# 2. 批量提交任务:使用 zip 将多个参数组合成元组
names = ["P1", "P2", "P3"] # 第一个参数的所有值(3个元素)
sleep_times = [1, 2, 3] # 第二个参数的所有值(3个元素,与names个数一致)
# 使用 zip 组合参数,避免多参数 map 的问题
# tasks = list(zip(names, sleep_times))
# 关键:map 是阻塞型调用,会一直阻塞到所有任务完成才返回结果列表
results = pool.map(task, zip(names, sleep_times))
# 4. 关闭进程池(规范收尾)
pool.close() # 禁止新任务提交
pool.join() # 等待所有任务完成(此处 map 已阻塞过,join 更多是释放资源)
# 5. 打印结果
for res in results:
print(res)
if __name__ == '__main__':
test_processpool_map()
2. 多线程
Python 中的多线程(Multithreading)是一种并发编程技术,允许程序在同一进程中同时运行多个“线程”(轻量级的执行单元),从而提高 I/O 密集型任务的效率。但由于全局解释器锁(GIL)的存在,Python 的多线程不能真正并行执行 CPU 密集型任务(在 CPython 实现中)。
2.1 创建多线程方式
2.1.1 直接实例化Thread类
import threading
import time
# 定义线程要执行的函数
def task(name, delay):
print(f"线程 {name} 启动,延迟 {delay} 秒")
time.sleep(delay)
print(f"线程 {name} 结束")
if __name__ == '__main__':
# 1. 创建 Thread 对象(target 传入函数,args 传入函数参数,元组形式)
t1 = threading.Thread(target=task, args=("Thread-1", 2))
t2 = threading.Thread(target=task, args=("Thread-2", 1))
# 2. 启动线程(调用 start(),底层会调用 target 对应的函数)
t1.start()
t2.start()
# 可选:等待线程结束(主线程阻塞直到子线程完成)
t1.join()
t2.join()
print("所有线程执行完毕")
2.1.2 继承 Thread 类
import threading
import time
# 自定义线程类,继承 threading.Thread
class MyThread(threading.Thread):
# 重写 __init__ 方法(可选,用于传递参数)
def __init__(self, name, delay):
super().__init__() # 必须调用父类构造方法
self.name = name
self.delay = delay
# 重写 run() 方法:线程要执行的逻辑
def run(self):
print(f"线程 {self.name} 启动,延迟 {self.delay} 秒")
time.sleep(self.delay)
print(f"线程 {self.name} 结束")
if __name__ == '__main__':
# 创建自定义线程实例
t1 = MyThread("Thread-1", 2)
t2 = MyThread("Thread-2", 1)
# 启动线程(自动调用 run() 方法)
t1.start()
t2.start()
# 等待线程结束
t1.join()
t2.join()
print("所有线程执行完毕")
2.2 线程池
submit(fn, *args, **kwargs) → Future
- 提交单个任务,立即返回 Future对象
- 非阻塞!
map(func, *iterables, timeout=None, chunksize=1) → 迭代器
- 类似内置
map(),但并行执行,按输入顺序返回结果 - 适合处理列表/生成器
shutdown(wait=True, *, cancel_futures=False)
- 关闭线程池:
wait=True(默认):阻塞直到所有已提交任务完成wait=False:立即返回,后台继续执行cancel_futures=True(Python 3.9+):取消所有未开始的任务
import concurrent.futures
import time
from loguru import logger
def worker(number):
logger.info(f"Worker {number} 开始")
time.sleep(2) # 模拟耗时操作
logger.info(f"Worker {number} 结束")
return number
def main():
start_time = time.time()
with concurrent.futures.ThreadPoolExecutor(max_workers=5) as executor:
futures = [executor.submit(worker, i) for i in range(5)]
multi_result = []
for future in concurrent.futures.as_completed(futures):
multi_result.append(future.result())
cost_time = time.time() - start_time
logger.info(f"结果: {multi_result}")
logger.info(f"方法耗时 {cost_time:.2f}s")
if __name__ == "__main__":
main()
2.2.1 submit(非阻塞)
1. as_completed
as_completed 按照任务完成顺序返回结果(而不是提交顺序)
2. shutdown
在 Python 的多线程编程中,shutdown 是线程池(特别是 concurrent.futures.ThreadPoolExecutor)生命周期管理的关键操作,用于优雅地关闭线程池,释放资源。理解它对避免资源泄漏、确保程序健壮性至关重要。
● wait=True(默认)—— 优雅关闭
● wait=False —— 快速关闭(后台继续)
3. callback
add_done_callback() 必须在 submit() 之后、result() 之前
(或不调用 result())调用,否则可能失效或失去意义。
2.2.1.1 submit
from concurrent.futures import ThreadPoolExecutor
import time
from loguru import logger
def task(name, delay):
logger.info(f"线程 {name} 启动,延迟 {delay} 秒")
time.sleep(delay)
logger.info(f"任务 {name} {delay} 秒 完成")
return delay # 任务返回值
if __name__ == '__main__':
start_time = time.time()
# 创建线程池(max_workers 指定最大线程数)
with ThreadPoolExecutor(max_workers=2) as executor:
# 提交任务到线程池(返回 Future 对象,用于获取结果) 假设线程池中有空闲线程,会立即启动改线程
future1 = executor.submit(task, "Thread-1", 2)
future2 = executor.submit(task, "Thread-2", 1)
logger.info(f"mian.........................................1")
# submit 非阻塞, task不需要执行完成, future1.result(timeout=3) 可指定超时时间,避免无限等待,
result1 = future1.result()
result2 = future2.result()
logger.info(f"mian.........................................2")
logger.info([result1, result2])
cost_time = time.time() - start_time
logger.info(f"所有线程执行完 耗时: {cost_time:.2f}S")
2.2.1.2 as_completed
from concurrent.futures import ThreadPoolExecutor, as_completed
import time
from loguru import logger
def task(name, delay):
logger.info(f"线程 {name} 启动,延迟 {delay} 秒")
time.sleep(delay)
logger.info(f"任务 {name} {delay} 秒 完成")
return delay # 任务返回值
if __name__ == '__main__':
start_time = time.time()
# 创建线程池(max_workers 指定最大线程数)
with ThreadPoolExecutor(max_workers=2) as executor:
# 提交任务到线程池(返回 Future 对象,用于获取结果) 假设线程池中有空闲线程,会立即启动改线程
future1 = executor.submit(task, "Thread-1", 2)
future2 = executor.submit(task, "Thread-2", 1)
future_list = [future1, future2]
result_list = []
logger.info(f"mian.........................................1")
# as_completed 按照任务完成顺序返回结果, 先返回Thread-2, 再返回Thread-1
for future in as_completed(future_list):
cur_result = future.result()
result_list.append(cur_result)
logger.info(f"mian.........................................2")
logger.info(result_list)
cost_time = time.time() - start_time
logger.info(f"所有线程执行完 耗时: {cost_time:.2f}S")
2.2.1.2 shutdown
在 Python 的多线程编程中,**shutdown**<是线程池(特别是 concurrent.futures.ThreadPoolExecutor)生命周期管理的关键操作,用于优雅地关闭线程池,释放资源。理解它对避免资源泄漏、确保程序健壮性至关重要。
wait=True(默认)—— 优雅关闭wait=False<—— 快速关闭(后台继续)
from concurrent.futures import ThreadPoolExecutor
import time
from loguru import logger
def task(name, delay):
logger.info(f"任务 {name} 启动,延迟 {delay} 秒执行")
time.sleep(delay)
logger.info(f"任务 {name} {delay} 秒 完成")
return delay
if __name__ == '__main__':
flag = False
# 手动创建线程池(不用 with)
executor = ThreadPoolExecutor(max_workers=1)
# 提交任务
future1 = executor.submit(task, "Task-1", 1)
future2 = executor.submit(task, "Task-2", 2)
logger.info(f"main shutdown之前............................. 1")
# 关键:所有任务提交完成后,调用 shutdown, 线程池中线程继续执行, 后续不能在提交线程
# wait=True:阻塞主线程,直到所有已提交任务完成后再关闭 wait=False:不阻塞主线程,立即关闭,任务在后台继续执行
executor.shutdown(wait=flag)
logger.info(f"main shutdown 之后............................. 2")
try:
executor.submit(task, "Task-3", 3) # 报错:cannot schedule new futures after shutdown
except Exception as e:
logger.info(f"main shutdown之后不再接受新任务................ ")
logger.error(e)
logger.info("线程池已手动关闭")
2.2.1.2 callback
add_done_callback()必须在 **submit()** 之后、**result()**** 之前(或不调用 **result()**)调用**,否则可能失效或失去意义。下面我们深入解释 为什么 以及 最佳调用时机。
from concurrent.futures import ThreadPoolExecutor
import time
from loguru import logger
def task(name, delay):
logger.info(f"线程 {name} 启动,延迟 {delay} 秒")
time.sleep(delay)
logger.info(f"任务 {name} {delay} 秒 完成")
return delay # 任务返回值
# 定义回调函数(参数必须是 Future 对象)
def callback_task(future):
# 通过 future.result() 获取任务结果
logger.info(f"回调函数触发:{future.result()}")
if __name__ == '__main__':
start_time = time.time()
# 创建线程池(max_workers 指定最大线程数)
with ThreadPoolExecutor(max_workers=2) as executor:
# 提交任务到线程池(返回 Future 对象,用于获取结果) 假设线程池中有空闲线程,会立即启动改线程
future1 = executor.submit(task, "Thread-1", 2)
future2 = executor.submit(task, "Thread-2", 1)
future1.add_done_callback(callback_task)
future2.add_done_callback(callback_task)
logger.info(f"mian.........................................1")
# submit 非阻塞, task不需要执行完成, future1.result(timeout=3) 可指定超时时间,避免无限等待,
result1 = future1.result()
result2 = future2.result()
logger.info(f"mian.........................................2")
logger.info([result1, result2])
# 为每个任务绑定回调函数
cost_time = time.time() - start_time
logger.info(f"所有线程执行完 耗时: {cost_time:.2f}S")
2.2.2 map
1. results = executor.map(task, names, delays) # 非阻塞
2. results = list(executor.map(task, names, delays)) # 阻塞,相当于已经调用迭代器
results = executor.map(task, names, delays)→ **非阻塞 **✅
- ✅ 仅提交任务到线程池队列;
- ✅ 立即返回一个 惰性迭代器(lazy iterator);
- ✅ 工作线程开始并行执行任务,但主线程不等待;
- ✅ 后续代码(如下一行
logger.info)立即执行; - 🔍 类型:
<class 'map'>(底层是concurrent.futures._base.ResultIterator)
📌 类比:就像按下“开始下载”按钮,任务已派发,但你还没点“查看下载完成的文件”。
results = list(executor.map(task, names, delays))→ **阻塞 **✅
- ✅
executor.map(...)先返回迭代器(非阻塞); - ❗
list(...)立即遍历整个迭代器; - ⏳ 遍历过程会:
- 等待第 1 个任务完成 → 取结果;
- 等待第 2 个任务完成 → 取结果;
- ……
- 直到最后一个任务完成;
- ✅ 最终构建完整列表
[res0, res1, ..., resN]并赋值给results; - ⏱️ 耗时 ≈ 最长任务的执行时间(因为是并行执行 + 按序取结果);
📌 类比:点“开始下载”后,立刻要求“等所有文件下完再打包给我”——你只能干等着。
2.2.2.1 非阻塞式
from concurrent.futures import ThreadPoolExecutor
import time
from loguru import logger
def task(name, delay):
logger.info(f"线程 {name} 启动,延迟 {delay} 秒")
time.sleep(delay)
logger.info(f"任务 {name} {delay} 秒 完成")
return delay # 任务返回值
"""
executor.map(task, names, delays) 会一次性提交所有任务到线程池,线程池会用空闲线程并行执行(最多 3 个,对应 max_workers=3);
当执行 for res in results 时,迭代器会先阻塞等待第一个任务完成,拿到结果后立即打印;
打印完第一个结果后,迭代器继续阻塞等待第二个任务完成(无论其他任务是否先执行完,都按提交顺序等),拿到结果后打印;
以此类推,直到所有任务的结果按提交顺序打印完毕 —— 全程任务是并行执行的,只是结果输出严格遵循提交顺序,而非执行完成顺序。
"""
if __name__ == '__main__':
start_time = time.time()
# 批量任务的参数(迭代器形式)
names = ["Thread-1", "Thread-1", "Thread-1"]
delays = [2, 1, 3]
with ThreadPoolExecutor(max_workers=4) as executor:
results = executor.map(task, names, delays) # map 自动将 names 和 delays 中的元素作为参数传入 task,返回结果迭代器
logger.info(f"mian.........................................1")
logger.info(f"reuslts: {list(results)} type_list: {type(results)}")
for res in results: # 迭代获取结果(按提交顺序返回,即使 Task-2 先完成,也会在 Task-1 之后输出)
logger.info(res)
logger.info(f"mian.........................................2")
cost_time = time.time() - start_time
logger.info(f"所有批量任务执行完毕 耗时{cost_time:.3f}S")
2.2.2.2 阻塞式
from concurrent.futures import ThreadPoolExecutor
import time
from loguru import logger
def task(name, delay):
logger.info(f"线程 {name} 启动,延迟 {delay} 秒")
time.sleep(delay)
logger.info(f"任务 {name} {delay} 秒 完成")
return delay # 任务返回值
"""
executor.map(task, names, delays) 会一次性提交所有任务到线程池,线程池会用空闲线程并行执行(最多 3 个,对应 max_workers=3);
当执行 for res in results 时,迭代器会先阻塞等待第一个任务完成,拿到结果后立即打印;
打印完第一个结果后,迭代器继续阻塞等待第二个任务完成(无论其他任务是否先执行完,都按提交顺序等),拿到结果后打印;
以此类推,直到所有任务的结果按提交顺序打印完毕 —— 全程任务是并行执行的,只是结果输出严格遵循提交顺序,而非执行完成顺序。
"""
if __name__ == '__main__':
start_time = time.time()
# 批量任务的参数(迭代器形式)
names = ["Thread-1", "Thread-1", "Thread-1"]
delays = [2, 1, 3]
with ThreadPoolExecutor(max_workers=4) as executor:
# executor.map() 本身是非阻塞的(立即返回),但它返回的是一个惰性迭代器(lazy iterator),真正阻塞发生在你迭代这个迭代器(如 for res in results 或 list(results))时。
# results = executor.map(task, names, delays) # map 自动将 names 和 delays 中的元素作为参数传入 task,返回结果迭代器
results = list(executor.map(task, names, delays)) # map 自动将 names 和 delays 中的元素作为参数传入 task,返回结果迭代器
logger.info(f"mian.........................................1")
logger.info(f"reuslts: {list(results)} type_list: {type(results)}")
for res in results: # 迭代获取结果(按提交顺序返回,即使 Task-2 先完成,也会在 Task-1 之后输出)
logger.info(res)
logger.info(f"mian.........................................2")
cost_time = time.time() - start_time
logger.info(f"所有批量任务执行完毕 耗时{cost_time:.3f}S")
2.3 锁
参考连接: https://www.runoob.com/python/python-multithreading.html
多线程中的锁(Lock)是一种同步机制,用于控制多个线程对共享资源的并发访问,避免因竞争条件(Race Condition)导致的数据不一致、程序崩溃或逻辑错误等问题。
threading.Lock —— 基础互斥锁
- 最简单的锁,用于确保某段代码(临界区)同一时刻只被一个线程执行。
- 非重入:同一个线程重复
acquire()会导致死锁。
threading.RLock —— 可重入锁(Reentrant Lock)
- 同一线程可多次
acquire(),不会死锁,需对应次数release()。 - 适用于递归调用或函数嵌套中需**多次加锁**的场景。
2.3.1. 死锁
如何避免?
- 方法 1:固定加锁顺序(推荐!)
- 使用
acquire(timeout=...)避免无限等待
import threading
import time
# 创建两把锁
lock_A = threading.Lock()
lock_B = threading.Lock()
def thread_1(): # Lock-A Lock-B
print("Thread-1: 尝试获取 lock_A...")
lock_A.acquire()
print("Thread-1: ✅ 拿到 lock_A")
time.sleep(0.5) # 模拟处理时间 —— 关键!制造交叉时机
print("Thread-1: 尝试获取 lock_B...")
lock_B.acquire() # ❗ 此时 Thread-2 已持有 lock_B,等待中...
print("Thread-1: ✅ 拿到 lock_B")
# 临界区
print("Thread-1: 执行任务...")
# 释放锁(但由于死锁,永远到不了这里)
lock_B.release()
lock_A.release()
def thread_2(): # Lock-B Lock-A
print("Thread-2: 尝试获取 lock_B...")
lock_B.acquire()
print("Thread-2: ✅ 拿到 lock_B")
time.sleep(0.5) # 同样延迟 —— 此时 Thread-1 已拿到 lock_A
print("Thread-2: 尝试获取 lock_A...")
lock_A.acquire() # ❗ 此时 Thread-1 已持有 lock_A,等待中...
print("Thread-2: ✅ 拿到 lock_A")
print("Thread-2: 执行任务...")
lock_A.release()
lock_B.release()
# 触发死锁
if __name__ == '__main__':
# 启动两个线程
t1 = threading.Thread(target=thread_1, name="Thread-1")
t2 = threading.Thread(target=thread_2, name="Thread-2")
t1.start()
t2.start()
t1.join()
t2.join()
print("✅ 所有线程结束") # ❌ 永远不会打印!
2.3.2 吃火锅(互斥锁)
# coding=utf-8
import threading
import time
def chiHuoGuo(people, operator):
print("%s 吃火锅的小伙伴:%s" % (time.strftime("H:%M:%S"), people))
time.sleep(1)
for i in range(3):
time.sleep(1)
print("%s %s正在 %s 鱼丸"% (time.strftime("H:%M:%S"), people, operator))
print("%s 吃火锅的小伙伴:%s" % (time.strftime("H:%M:%S"), people))
lock = threading.Lock() # 1. 定义全局锁(所有线程共享同一个锁) 如果这个锁加载方法里面 那么就相当于是没有加锁了
def hot_pot(people:str, operator:str, lock_flag:bool):
print(f"开始线程: {threading.current_thread().name}")
if lock_flag: # 2. 判断是否加锁
lock.acquire()
chiHuoGuo(people, operator) # 加锁 只能所有等小明把所有鱼丸添加之后才能吃
lock.release()
else:
chiHuoGuo(people, operator)
print(f"结束线程: {threading.current_thread().name}")
if __name__ == '__main__':
print("海底捞火锅开始......................................")
# 设置线程组
threads = []
# 创建新线程
thread1 = threading.Thread(target=hot_pot, name="Thread-1", args=("xiaoming", "添加", True))
thread2 = threading.Thread(target=hot_pot, name="Thread-2", args=("xiaowang", "吃掉", True))
# 添加到线程组
threads.append(thread1)
threads.append(thread2)
# 开启线程
for thread in threads:
thread.start()
# 阻塞主线程,等子线程结束
for thread in threads:
thread.join()
time.sleep(0.1)
print("退出主线程:吃火锅结束,结账走人")
2.3.3 银行取钱(互斥锁)
前提: 假定这是你的银行存款为0, 调用一次change_it方法期间主要是进行了存款和取款的操作, 正常情况下
-
- 存款后一定是5或者8
-
- 取款后一定是0: 如果不是则说明有问题
# multithread
import time, threading
# 前提: 假定这是你的银行存款为0, 调用一次change_it方法期间主要是进行了存款和取款的操作, 正常情况下
# 1. 存款后一定是5或者8
# 2. 取款后一定是0: 如果不是则说明有问题
balance = 0
lock = threading.Lock()
def change_it(n, i):
# 先存后取,结果应该为0:
global balance
balance = balance + n
time.sleep(0.1)
print(f"{threading.current_thread().name} balance 轮次{i}-存后-:............................{balance}")
balance = balance - n
print(f"{threading.current_thread().name} balance 轮次{i}-取后-:............................{balance}")
def run_thread(money, is_lock):
for round in range(2):
if is_lock:
lock.acquire() # 先要获取锁:
change_it(money, round) # 存款取款操作
lock.release() # 释放锁
else:
change_it(money, round)
# 不能保证change_it操作是线程安全的, 可能导致存后是13,取后是5、8 存在潜在问题
if __name__ == '__main__':
lock_flag = True
t1 = threading.Thread(target=run_thread, name="Thread5 ", args=(5, lock_flag))
t2 = threading.Thread(target=run_thread, name="Thread8 ", args=(8, lock_flag))
t1.start()
t2.start()
t1.join()
t2.join()
print(balance)
2.3.4 银行取钱(可重入锁)
# multithread
import time, threading
# 前提: 假定这是你的银行存款为0, 调用一次change_it方法期间主要是进行了存款和取款的操作, 正常情况下
# 1. 存款后一定是5或者8
# 2. 取款后一定是0: 如果不是则说明有问题
balance = 0
# lock = threading.Lock() # 互斥锁 会卡在第二次锁中
lock = threading.RLock() # 互斥锁 会卡在第二次锁中
def change_it(money:str, round:str, is_lock:bool):
# 先存后取,结果应该为0:
global balance
print(f"change_it...........................................")
with lock:
balance = balance + money
time.sleep(0.1)
print(f"{threading.current_thread().name} balance 轮次{round}-存后-:............................{balance}")
balance = balance - money
print(f"{threading.current_thread().name} balance 轮次{round}-取后-:............................{balance}")
def run_thread(money, is_lock):
for round in range(2):
if is_lock:
lock.acquire() # 先要获取锁:
change_it(money, round, is_lock) # 存款取款操作
lock.release() # 释放锁
else:
change_it(money, round, is_lock)
# 不能保证change_it操作是线程安全的, 可能导致存后是13,取后是5、8 存在潜在问题
if __name__ == '__main__':
lock_flag = True
t1 = threading.Thread(target=run_thread, name="Thread5 ", args=(5, lock_flag))
t2 = threading.Thread(target=run_thread, name="Thread8 ", args=(8, lock_flag))
t1.start()
t2.start()
t1.join()
t2.join()
print(balance)
2.4 threading.Event
threading.Event() 是 Python 多线程中用于 线程间通信与同步 的核心工具,本质是一个「线程间的信号开关」—— 通过简单的布尔标志(True/False),让一个或多个线程等待其他线程的「通知」,实现线程间的协作(比如 “等待条件满足”“暂停 / 恢复线程”“批量唤醒”)。
简单类比:它就像一个「红绿灯」——
wait()的线程 = 等待放行的车辆;set()= 绿灯(放行所有等待的线程);clear()= 红灯(让后续线程继续等待);- 核心价值是「传递信号」,而非保护共享资源(区别于锁)。
核心原理Event 内部维护一个布尔值标志(默认 False):
- 标志为
False时:调用wait()的线程会 阻塞暂停,直到标志变为True; - 标志为
True时:调用wait()的线程会 直接通过,不阻塞; set():将标志设为True,唤醒所有正在wait()的线程;clear():将标志设为False,后续调用wait()的线程会再次阻塞。
核心方法(4 个必掌握)
| 方法 | 作用 |
|---|---|
event.set() |
标志设为 True,唤醒所有等待的线程(线程被唤醒后继续执行)。 |
event.clear() |
标志设为 False,后续调用 wait()的线程会阻塞。 |
event.wait(timeout=None) |
阻塞线程,直到标志为 True或超时:+ timeout=None:无限等待;+ timeout=数字:最多等 数字秒(支持小数);+ 返回值: True= 被唤醒,False= 超时。 |
event.is_set() |
判断标志是否为 True(返回布尔值),用于查询信号状态。 |
2.4.1 wait用法
import threading
import time
init_event = threading.Event() # 创建Event对象
def sub_task():
print("sub_task :初始化中(如连接数据库)...")
time.sleep(1) # 模拟耗时操作
print("sub_task :初始化完成!")
init_event.set() # 发送“完成”信号 放行 这之后一般不执行代码
print("sub_task :完成")
def test_wait():
t1 = threading.Thread(target=sub_task)
t1.start()
print("主线程:等待初始化.............................")
success = init_event.wait(timeout=3) # 最多等3秒
if success:
print("主线程:开始执行业务!...........................") # 收到信号后执行
t1.join()
else:
print("主线程:初始化超时!.............................")
# wait 等待 set 信号(标志位为True,立马放行)
if __name__ == "__main__":
start_time = time.time()
test_wait()
cost_time = time.time() - start_time
print(f"总耗时:{cost_time:.2f}")
2.4.2 标志位set()
- 设置一个函数, pause_event.is_set() 标志位为False的时候退出这个子线程
import threading
import time
init_event = threading.Event() # 创建Event对象
def sub_task():
print("sub_task :初始化中(如连接数据库)...")
time.sleep(1) # 模拟耗时操作
print("sub_task :初始化完成!")
init_event.set() # 发送“完成”信号 放行
print("sub_task :开始执行业务...")
def test_wait():
t1 = threading.Thread(target=sub_task)
t1.start()
print("主线程:等待初始化.............................")
success = init_event.wait(timeout=3) # 最多等3秒
if success:
print("主线程:开始执行业务!...........................") # 收到信号后执行
t1.join()
else:
print("主线程:初始化超时!.............................")
pause_event = threading.Event()
pause_event.set() # 初始为True(线程可执行)
number = 0
def worker():
global number
while pause_event.is_set():
number = number + 1
print(f"worker:执行业务.........................{number}")
time.sleep(1)
print("worker:结束.........................")
def test_set_sleep():
"""
clear 和 set之间 设置sleep 可以保证True -> False, False -> True 中 panuse_event.set() 可以检测到False,
否则可能 clear() set() 之后 worker正在sleep, 一觉醒来还是True, 导致worker线程无法结束
"""
t1 = threading.Thread(target=worker, daemon=True)
t1.start()
time.sleep(2)
pause_event.clear()
print("主线程:暂停子线程...运行主线程")
time.sleep(1.5) # 只要睡眠时间大于子线程的睡眠时间(此时while循环必有一次条件为False),则子线程必然会结束
print("主线程:恢复子线程...")
pause_event.set()
time.sleep(2)
print("主线程:结束...")
def test_set_no_sleep():
"""
没有sleep, worker不会正确关闭
"""
t1 = threading.Thread(target=worker, daemon=True)
t1.start()
time.sleep(2)
pause_event.clear()
print("主线程:暂停子线程...运行主线程")
# time.sleep(1.5) # 没有sleep, 标志位初始化为True, -> False, -> True, 由于过程很快所以下一次判断标志还是True, 所以子线程不会结束
print("主线程:恢复子线程...")
pause_event.set()
time.sleep(2)
print("主线程:结束...")
#
# wait 等待 set 信号(标志位为True,立马放行)
# 初始化pause_event.set() # 初始为True(线程可执行)
if __name__ == "__main__":
start_time = time.time()
# test_wait()
# test_set_sleep()
test_set_no_sleep()
cost_time = time.time() - start_time
print(f"总耗时:{cost_time:.2f}")
2.4.3 常见用法
import threading
import time
def worker():
global number
number = 0
while True:
# 关键:如果信号被清除(暂停),则阻塞等待;信号被设置(恢复)则继续
pause_event.wait() # 等价于:while not pause_event.is_set(): time.sleep(0.1)
# 业务逻辑
number += 1
print(f"worker:执行业务.........................{number}")
time.sleep(1)
# 检查是否需要终止(单独用一个 Event 控制终止,更清晰)
if stop_event.is_set():
break
print("worker:结束.........................")
if __name__ == "__main__":
# 两个 Event 分工:一个控制“暂停/恢复”,一个控制“终止”(职责分离,更易维护)
pause_event = threading.Event() # 控制暂停/恢复
stop_event = threading.Event() # 控制终止
pause_event.set() # 初始为 True:允许子线程执行(不暂停)
number = 0
t1 = threading.Thread(target=worker, daemon=True)
t1.start()
# 阶段1:子线程运行2秒
time.sleep(2)
print("主线程:暂停子线程...")
pause_event.clear() # 清除信号 → 子线程阻塞(暂停)
# 阶段2:主线程运行3秒(子线程此时暂停)
time.sleep(3)
print("主线程:恢复子线程...")
pause_event.set() # 设置信号 → 子线程继续执行
# 阶段3:恢复后运行2秒,再终止子线程
time.sleep(2)
print("主线程:终止子线程...")
stop_event.set() # 发送终止信号
pause_event.set() # 确保子线程从阻塞中唤醒,才能检查到终止信号
t1.join() # 等待子线程优雅退出
print("主线程:结束...")
3. 异步
参考博客: https://blog.csdn.net/qq_45056135/article/details/150404125?spm=1001.2014.3001.5502
更多推荐



所有评论(0)