介绍

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 传统方式)

  1. 同步一般使用map, 异步一般使用apply_async
  • map: 主进程完全阻塞,直到全部任务完成
  • apply_async : 非阻塞提交 + 按需阻塞获取
  1. 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))   		# 阻塞,相当于已经调用迭代器
  1. results = executor.map(task, names, delays) → **非阻塞 **
  • ✅ 仅提交任务到线程池队列;
  • ✅ 立即返回一个 惰性迭代器(lazy iterator)
  • ✅ 工作线程开始并行执行任务,但主线程不等待;
  • ✅ 后续代码(如下一行 logger.info立即执行
  • 🔍 类型:<class 'map'>(底层是 concurrent.futures._base.ResultIterator

📌 类比:就像按下“开始下载”按钮,任务已派发,但你还没点“查看下载完成的文件”。


  1. 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方法期间主要是进行了存款和取款的操作, 正常情况下

    1. 存款后一定是5或者8
    1. 取款后一定是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()

  1. 设置一个函数, 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

Logo

Agent 垂直技术社区,欢迎活跃、内容共建。

更多推荐