news 2026/10/3 19:35:43

第 12 期:线程、进程和协程,到底应该选哪一个

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
第 12 期:线程、进程和协程,到底应该选哪一个

并发工具解决不了“程序慢”这四个字。它只能处理某一种具体的慢,而且每多推进一份工作,也会多带来一份调度、状态和资源成本。

写在前面

有一次,订单对账任务从每天处理两万笔涨到了二十万笔。

原来的程序很老实:查一批订单,逐个调用渠道接口补齐流水,再计算差异,最后写报表。串行跑完要四十多分钟。有人看见循环里都是独立订单,很自然地把它改成了 64 个线程。

第一版上线后,时间确实降了,但没有想象中那么多。渠道查询快了,后半段的规则计算几乎没变;数据库连接池却被占满,偶尔还出现汇总金额对不上的情况。继续把线程调到 128,任务反而更慢,内存和上下文切换都涨了。

后来把整条链路拆开看,事情就清楚了:

读取订单:等待数据库 查询渠道:等待网络 解析响应:少量 CPU 规则匹配:大量纯 Python 计算 写入报表:等待存储

它不是一种慢,而是几种完全不同的工作被串在了一起。

线程在等待网络时能去推进别的订单,所以渠道查询明显加快。到了纯 Python 规则计算阶段,64 个线程仍然要争用同一个解释器锁,并没有让 64 个 CPU 核心同时执行 Python 字节码。与此同时,这些线程还共享汇总字典和数据库客户端;没有并发上限、没有锁,也没有按阶段隔离资源,问题自然一起冒了出来。

我见过不少类似改造。真正的误区通常不是不会用ThreadPoolExecutor,而是一看到“任务彼此独立”,就跳过了更前面的问题:它们大部分时间究竟在等什么?

这一期仍然沿着订单任务往下走。我们会用同一组等待型任务和纯计算任务,分别交给串行、线程池、进程池和协程,看看吞吐、启动成本与内存发生了什么。比 API 更重要的是选择依据:I/O 等待、CPU 计算、共享状态、任务粒度,还有系统真正能承受的并发量。

并发和并行,不是同一件事

这两个词经常被混着用。

并发描述的是一段时间内,多个任务都在向前推进。它们不一定在同一时刻执行。

单个执行核心: 任务 A:执行 ─── 等网络 ───────── 执行 任务 B: 执行 ─ 等磁盘 ───── 执行 时间 :──────────────────────────────>

任务 A 等网络时,执行权交给任务 B。一个核心也能并发。

并行则是多个任务真的在同一时刻执行:

核心 1:任务 A ───────────────────> 核心 2:任务 B ───────────────────> 核心 3:任务 C ───────────────────>

并发擅长填补等待空档,并行才能缩短可拆分计算的墙上时间。

线程和协程都能实现并发。多个进程可以并发,也可以在多个 CPU 核心上并行。线程能不能并行执行 Python 代码,还要看解释器、GIL 以及底层扩展是否释放了它。

先把目标说得更具体也很重要。

如果用户一次请求要查 20 个下游,我们关心的是这次请求的端到端延迟;如果后台每晚处理 500 万条记录,我们可能更关心总吞吐;如果服务长期运行,还要关心 P99、内存、连接数和故障恢复。只说“快了两倍”,常常掩盖了到底哪项指标变好、哪项资源正在变坏。

先看任务在忙,还是在等

判断并发模型前,我通常先把任务时间粗略分成两类:

CPU 时间:执行 Python、压缩、解析、计算、加密 等待时间:网络、磁盘、数据库、锁、限流与休眠

如果进程 CPU 长期只有 20%,大部分调用栈停在 socket 读取,继续优化 Python 循环很可能收效不大。反过来,单核已经接近 100%,剖析结果又集中在业务计算函数,增加线程通常不是答案。

证据可以从这些地方来:

  • 业务阶段耗时,区分排队、连接、首字节、读取和计算。
  • 进程与单核 CPU 利用率。
  • py-spy top或火焰图中的热点函数。
  • 数据库、HTTP 客户端和连接池等待时间。
  • 任务运行期间的上下文切换、内存与运行队列。
  • 增加并发后,吞吐是否继续增长,P99 是否开始恶化。

别只看整个函数用了 500 毫秒。一次 HTTP 调用可能有 480 毫秒在等对方,真正占 CPU 的只有 20 毫秒;另一段 JSON 规则处理虽然也用 500 毫秒,却一直在执行 Python。两者的优化方向完全不同。

GIL 限制的是什么

在传统 CPython 构建中,GIL,也就是全局解释器锁,保证同一时刻通常只有一个线程执行 Python 字节码。

这句话有两个容易走偏的读法。

第一种是“Python 线程没有用”。不对。线程在阻塞 I/O 时会释放 GIL,另一个线程可以继续运行。下载文件、访问数据库、等待外部接口,本来就有大量空档,线程可以把这些空档叠起来。

第二种是“线程永远不能使用多个核心”。也不严谨。NumPy、压缩库、图像处理和部分加密实现会在 C 扩展中释放 GIL,底层代码可能并行。是否释放、在哪个数据规模释放,要看具体库,不能只看函数名猜。

我们真正能稳妥地说的是:

  • 大量纯 Python 计算通常不能靠多个线程获得多核并行。
  • 阻塞 I/O 任务通常能从线程并发中受益。
  • 调用释放 GIL 的扩展库时,线程可能并行,需要实测。
  • GIL 不是业务锁,不能替你保护一组共享状态的不变量。

近年的 CPython 已经出现 free-threaded 构建,但它和传统构建在运行方式、扩展兼容及性能特征上仍有差异。项目如果明确部署这种运行时,应按那个环境重新测试;不能把新运行时的能力倒推到现有 Python 3.9 服务上。

线程池适合把等待叠起来

假设我们要查询一批订单的物流状态,客户端是同步阻塞接口:

fromconcurrent.futuresimportThreadPoolExecutorfromcollections.abcimportIterabledeffetch_tracking_status(order_id:str,)->TrackingStatus:returnlogistics_client.fetch(order_id)deffetch_tracking_in_threads(order_ids:Iterable[str],*,max_workers:int=8,)->list[TrackingStatus]:withThreadPoolExecutor(max_workers=max_workers,thread_name_prefix="tracking",)asexecutor:returnlist(executor.map(fetch_tracking_status,order_ids,))

一个线程发出请求后等待响应,其他线程仍可继续发请求或解析结果。对于已有同步 SDK、任务数量中等的系统,线程池往往是改动最小的选择。

它也有很实际的优点:

  • 线程共享进程内对象,不需要把参数序列化后传给另一个解释器。
  • 创建成本通常低于进程。
  • 原有同步代码不必整体改成async。
  • 调试栈与异常传播相对直接。

可“能开线程”不代表应该把max_workers写成 100。

并发上限至少要同时看:

下游允许的 QPS 与并发数 HTTP / 数据库连接池大小 本机文件描述符与内存 单任务平均等待时间 超时、重试和限流策略 同一进程里其他请求的资源需求

数据库连接池只有 20 个连接,开 100 个查询线程,只会让 80 个线程换个地方排队。下游限制每秒 50 次请求,200 个线程可能带来更多429和重试,吞吐反而下降。

还有一个不太显眼的问题:把百万条输入一次性提交给线程池,任务对象本身也会占内存。有限线程数只限制“正在执行多少”,不一定限制“已经排队多少”。对于持续输入,应使用有界queue.Queue、分批提交,或者让生产者在队列满时等待。

并发上限是容量决定,不是语法参数。

线程共享内存,也共享麻烦

线程访问同一份字典、列表和客户端对象很方便,这也是竞态条件的来源。

下面这段汇总代码看起来很普通:

customer_totals:dict[int,int]={}defadd_order_amount(customer_id:int,amount_cents:int,)->None:current=customer_totals.get(customer_id,0)customer_totals[customer_id]=(current+amount_cents)

两个线程可能同时读到旧值 100:

线程 A:读取 100 线程 B:读取 100 线程 A:写入 130 线程 B:写入 150

正确结果应是 180,最终却只留下 150。GIL 只能约束某一时刻谁执行字节码,不能保证“读取、计算、写回”这一组业务动作不可分割。

可以用锁保护这段临界区:

fromthreadingimportLock customer_totals:dict[int,int]={}customer_totals_lock=Lock()defadd_order_amount(customer_id:int,amount_cents:int,)->None:withcustomer_totals_lock:current=customer_totals.get(customer_id,0)customer_totals[customer_id]=(current+amount_cents)

锁内只保留共享状态更新,不要顺手把网络调用也放进去:

withcustomer_totals_lock:response=logistics_client.fetch(order_id)update_total(response)

如果请求等两秒,其他线程也会跟着等两秒,并发几乎退化回串行。

更好的办法常常是减少共享写入。每个线程返回局部结果,由单一汇总阶段合并;或者按客户 ID 分片,让同一键只由固定 Worker 处理。锁不是坏东西,但共享状态越少,证明正确性越容易。

线程安全还包括客户端本身。某个 HTTP Client 能在线程间共享,不代表所有 SDK 都可以;Session、游标和事务对象尤其需要查看文档。不要拿“压测时没出错”代替线程安全契约。

进程池绕开 GIL,也隔开了内存

纯 Python 规则计算长期占满一个核心时,进程池更有机会提高吞吐。每个子进程有自己的解释器和 GIL,可以分布到多个 CPU 核心。

importmultiprocessingfromconcurrent.futuresimportProcessPoolExecutordefcalculate_risk_score(order:OrderSnapshot,)->RiskResult:returnrisk_engine.calculate(order)defmain()->None:context=multiprocessing.get_context("spawn")withProcessPoolExecutor(max_workers=8,mp_context=context,)asexecutor:results=list(executor.map(calculate_risk_score,load_order_snapshots(),chunksize=50,))write_results(results)if__name__=="__main__":multiprocessing.freeze_support()main()

入口保护不是装饰。使用spawn时,子进程会重新导入主模块;如果创建进程池的代码位于模块顶层,导入时会再次创建子进程,最后不是报错就是无限递归启动。

传给进程池的函数与参数通常还要能被pickle。局部函数、lambda、打开的数据库连接、锁和许多 C 扩展对象无法直接传递。工作函数尽量放在模块顶层,输入使用明确、紧凑的数据结构。

chunksize用来把多个小任务打包发送,减少进程间通信次数。值太小,调度与序列化成本可能压过计算;值太大,某个进程拿到慢任务后又会导致负载不均。没有一个适合所有数据的数字。

示例还假设load_order_snapshots()返回的是一批规模受控的数据。对持续流或海量输入,不要让executor.map()提前积累大量待执行任务;应按批次读取和提交,让父进程的待处理队列也有明确上限。

进程隔离也改变了状态语义:

  • 子进程修改普通全局变量,父进程看不到。
  • 每个进程会建立自己的数据库和网络连接。
  • 日志句柄、指标客户端与初始化代码要确认是否支持多进程。
  • 结果和异常需要通过 IPC 返回。
  • Worker 崩溃时,正在执行的任务可能需要重试。

少了线程级共享竞态,却多了数据交接与生命周期问题。

进程数量也不该直接等于宿主机逻辑核心数。容器可能只分到 2 核,os.cpu_count()却仍能看到更大的宿主机;每个进程还会复制解释器、模型和缓存。数值计算库若在每个进程内部再启动一组原生线程,8 个进程乘 8 个底层线程会造成过度并行。应以容器 CPU 配额、单 Worker 内存和实测吞吐为准,并为同机其他服务留下余量。

为什么短 CPU 任务放进进程池反而更慢

把一个函数交给进程池,至少会经历:

创建或唤醒子进程 ↓ 序列化函数参数 ↓ 通过进程间通道发送 ↓ 子进程反序列化 ↓ 执行计算 ↓ 序列化结果并传回

如果计算本身只花 2 毫秒,这一圈交接可能比工作还贵。

进程池更适合粒度足够大的任务,或者被长期复用。批处理每来一条记录就创建一个新池,是很昂贵的写法;Web 请求里临时启动八个子进程,也会把延迟和内存抬得很难看。

这次实测里,单个纯 Python 任务大约 50 毫秒。每轮都新建进程池时,24 个任务用时约 2.17 秒,串行只有约 1.24 秒。把进程预热并将单任务计算提高到约 100 毫秒后,进程池才明显拉开差距。具体数据后面会完整列出。

所以“CPU 密集就用进程”还不够。后面应该跟一句:计算要足以覆盖进程启动、序列化和调度成本。

大对象会在进程边界上交过路费

线程传递一个 100 MiB 对象,通常只是把同一对象引用交给另一个线程。进程不能直接读取另一个进程的 Python 堆,常规进程池会序列化并复制数据。

假设每个任务都收到一份巨大的订单列表:

withProcessPoolExecutor(max_workers=8)asexecutor:results=list(executor.map(calculate_batch,repeated_large_batches,))

内存可能同时存在:

  • 父进程里的原始订单对象。
  • 序列化缓冲区。
  • 进程间通道中的数据。
  • 多个子进程各自反序列化出的对象。
  • 返回结果的序列化副本。

Linux 的fork有写时复制,看起来可以先共享父进程内存页;一旦对象被修改,相应页面仍会复制。macOS 和 Windows 常用spawn,子进程从新的解释器开始,启动与初始化成本更明显。部署环境不同,不能直接照搬本地结论。

常见的改进方向有:

  • 只传文件路径、ID、偏移量等小参数,让 Worker 自己读取。
  • 使用进程池initializer,每个 Worker 只加载一次只读模型。
  • 把小任务合成批次,降低每条消息的序列化次数。
  • 对连续数值数据使用共享内存、NumPy 共享缓冲或内存映射。
  • 避免把 ORM 对象、客户端和层层嵌套的对象图跨进程传递。

共享内存减少复制,也把同步和生命周期责任还给了我们。谁创建、谁释放、并发写如何协调,都需要明确。它不是免费加速开关。

协程靠自愿让出执行权

协程通常运行在单个事件循环线程中。一个任务执行到await,发现 I/O 尚未完成,才把控制权交回事件循环,让其他任务继续。

importasyncioasyncdeffetch_tracking_status(order_id:str,)->TrackingStatus:returnawaitasync_logistics_client.fetch(order_id)

这里真正重要的不是函数前面的async,而是调用链中的 I/O 客户端也支持异步,并且在等待时正确await。

下面这段代码虽然写在协程里,仍会堵住整个事件循环:

importtimeasyncdeffetch_badly(order_id:str)->bytes:time.sleep(1)returnblocking_client.fetch(order_id)

time.sleep()不会把控制权交给事件循环,阻塞客户端也一样。其他协程只能干等。

如果暂时必须调用同步阻塞函数,Python 3.9 可以用asyncio.to_thread()把它放到线程池:

asyncdeffetch_legacy_client(order_id:str,)->TrackingStatus:returnawaitasyncio.to_thread(blocking_client.fetch,order_id,)

这是一座兼容桥,不会把同步库变成真正的非阻塞 I/O。底层仍占用线程,线程池大小、超时与共享状态问题仍然存在。

协程的优势在于单个任务对象通常比线程轻量,尤其适合大量同时等待的连接。它的代价是调用链需要配合,取消、超时、资源释放和异常传播也更容易被写漏。下一期会专门展开这些生命周期问题。

async不能让纯计算自动并行

把 CPU 函数改成async def,并不会凭空出现让出点:

asyncdefcalculate_score(order:OrderSnapshot,)->RiskResult:returnrisk_engine.calculate(order)

如果risk_engine.calculate()连续执行 300 毫秒纯 Python 代码,这 300 毫秒里事件循环无法处理其他 socket、超时和取消。把一百个这样的协程交给gather(),它们仍然会依次占住同一个线程。

CPU 计算可以显式交给进程池:

importasynciofromconcurrent.futuresimportProcessPoolExecutorasyncdefcalculate_in_process(process_pool:ProcessPoolExecutor,order:OrderSnapshot,)->RiskResult:event_loop=asyncio.get_running_loop()returnawaitevent_loop.run_in_executor(process_pool,calculate_risk_score,order,)

边界仍然存在:order要被序列化,取消等待这个 Future 也不一定能立即终止已经在子进程里运行的函数。协程只是负责等待进程结果,没有消除进程池的成本。

限制并发,不要一次创建所有任务

异步代码很容易写出这种版本:

results=awaitasyncio.gather(*(fetch_tracking_status(order_id)fororder_idinone_million_order_ids))

事件循环不会同时执行一百万个 Python 指令,但这里会一次创建大量协程与 Task,保存参数、状态和结果。内存先涨起来,下游也可能在很短时间内收到远超容量的请求。

任务数量有限时,可以用 Semaphore 约束在途请求:

asyncdeffetch_many_tracking_statuses(order_ids:list[str],*,concurrency:int=20,)->list[TrackingStatus]:semaphore=asyncio.Semaphore(concurrency)asyncdeffetch_one(order_id:str)->TrackingStatus:asyncwithsemaphore:returnawaitfetch_tracking_status(order_id)returnawaitasyncio.gather(*(fetch_one(order_id)fororder_idinorder_ids))

这限制了同时进入外部调用的数量,但仍然一次创建了与输入等量的协程。面对持续流或海量输入,更适合有界队列:

fromdataclassesimportdataclassfromtypingimportOptional@dataclass(frozen=True)classTrackingOutcome:order_id:strstatus:Optional[TrackingStatus]error_message:Optional[str]asyncdeftracking_worker(queue:asyncio.Queue[Optional[str]],outcomes:list[TrackingOutcome],)->None:whileTrue:order_id=awaitqueue.get()try:iforder_idisNone:returntry:status=awaitfetch_tracking_status(order_id)outcome=TrackingOutcome(order_id=order_id,status=status,error_message=None,)exceptExceptionaserror:outcome=TrackingOutcome(order_id=order_id,status=None,error_message=(f"{type(error).__name__}:{error}"),)outcomes.append(outcome)finally:queue.task_done()

生产者和入口负责有界写入:

fromcollections.abcimportAsyncIterableasyncdeffetch_from_stream(order_ids:AsyncIterable[str],*,worker_count:int=20,queue_size:int=100,)->list[TrackingOutcome]:queue:asyncio.Queue[Optional[str]]=asyncio.Queue(maxsize=queue_size,)outcomes:list[TrackingOutcome]=[]workers=[asyncio.create_task(tracking_worker(queue,outcomes))for_inrange(worker_count)]try:asyncfororder_idinorder_ids:# 队列满时在这里等待,把压力传回数据源。awaitqueue.put(order_id)for_inworkers:awaitqueue.put(None)awaitqueue.join()awaitasyncio.gather(*workers)returnoutcomesfinally:# 数据源异常或外层取消时,不能把 Worker 留在后台等待。forworkerinworkers:ifnotworker.done():worker.cancel()awaitasyncio.gather(*workers,return_exceptions=True,)

queue_size控制已读取但还没处理的数据,worker_count控制外部调用并发。队列满时,生产者停下来,这就是最直接的背压。

这段示例把每个失败记录成字符串结果,既避免 Worker 提前退出并把queue.join()永久卡住,也不会长期保留异常 traceback 引用的局部对象。outcomes本身仍会随输入增长;真正的海量任务应把结果分批写出,而不是全部留到函数返回。生产系统还要定义哪些异常可重试、结果是否按输入顺序返回。那些细节不能靠gather()默认替我们决定。

请求上下文不能放在线程级全局变量里

异步服务中的多个请求通常运行在同一线程。若把请求 ID 放进普通全局变量或只按线程保存,协程切换后会相互覆盖。

请求级上下文应使用ContextVar:

fromcontextvarsimportContextVarfromtypingimportOptional request_id_context:ContextVar[Optional[str]]=ContextVar("request_id",default=None,)asyncdefhandle_request(request:Request)->Response:token=request_id_context.set(request.request_id)try:returnawaitprocess_request(request)finally:request_id_context.reset(token)

ContextVar会随异步任务上下文传播,finally中恢复旧值,避免上下文泄漏到后续请求。

进程边界不会自动继承这份业务上下文。线程执行器的传播行为也取决于入口:asyncio.to_thread()会复制当前 Context,而普通run_in_executor()不应被默认当作上下文传播机制。跨边界需要的追踪 ID,最好作为显式参数传递。

混合任务,不必强迫一种模型包办

订单对账既有网络等待,也有纯 Python 计算。一个更合理的结构是分阶段:

异步或线程并发下载 ↓ 有界队列 进程池执行 CPU 规则 ↓ 有界结果队列 单独批量写入数据库

下载阶段的并发由渠道连接上限控制,计算阶段的进程数由 CPU 和内存控制,写入阶段则按数据库容量批量提交。每个阶段可以独立观察队列长度和处理速率。

用一种模型包办所有阶段,配置会互相牵制:

  • 为网络等待开很多进程,浪费内存。
  • 为纯 Python 计算开很多线程,无法获得多核并行。
  • 为了异步而把成熟同步 SDK 全部重写,风险可能高于收益。
  • 在事件循环里直接计算,又会拖住所有 I/O。

混合模型并不天然高级。阶段越多,交接、取消和排障越复杂。只有剖析证明瓶颈确实分布在不同类型工作上,这种拆分才值得。

用同一组任务做一次实测

为了把差异落到数字上,我写了两类可控任务。

I/O 任务只等待 30 毫秒,用来模拟网络响应空档:

importasyncioimporttimedefblocking_wait_task(task_id:int,delay_seconds:float,)->int:time.sleep(delay_seconds)returntask_idasyncdefasync_wait_task(task_id:int,delay_seconds:float,)->int:awaitasyncio.sleep(delay_seconds)returntask_id

CPU 任务只执行 Python 整数运算,不调用可能释放 GIL 的第三方扩展:

defpure_python_cpu_task(seed:int,iterations:int,)->int:value=seedforindexinrange(iterations):value=(value*1_664_525+1_013_904_223+index)&0xFFFFFFFFreturnvalue

测试环境是 Python 3.9.6、Apple M5 Pro、18 个逻辑 CPU。线程、进程和协程的并发上限都设为 8,每个实现返回相同结果后才记录时间,表格使用三次运行的中位数。

第一组是 40 个等待任务:

执行方式总耗时
串行1.3351 s
8 线程0.1725 s
8 进程,每轮冷启动2.1486 s
协程,最多 8 个在途0.1551 s

理论等待时间是40 × 0.03 = 1.2秒。线程和协程把等待叠成约五批,接近5 × 0.03 = 0.15秒。进程当然也能同时睡眠,但 macOSspawn启动解释器的成本远高于这点工作,没有使用价值。

第二组是 24 个纯 Python 计算任务,每个执行 100 万轮:

执行方式总耗时
串行1.2350 s
8 线程1.2537 s
8 进程,每轮冷启动2.1653 s
24 个 CPU 协程1.2572 s

线程和协程都没有多核收益。协程版本内部没有await,只是依次占用事件循环。冷启动进程池仍然太贵。

接着复用已经预热的池,把单任务计算提高到 200 万轮,执行 16 个任务:

执行方式总耗时
串行1.7822 s
8 个复用线程1.6412 s
8 个复用进程0.2192 s

线程比串行少的那一点不应解读成稳定加速。CPU 调频、系统噪声和中位数样本都可能造成小幅波动;两者仍在同一量级。复用进程池后,启动成本不再进入每批任务,纯 Python 计算才真正分布到多个核心。

最后传递同一个 5 MiBbytes对象给 16 个任务。池先预热,任务内部等待 30 毫秒,确保所有 Worker 都参与:

执行方式总耗时参与执行单元Worker 峰值 RSS 之和
8 线程0.0683 s1 个进程32.2 MiB
8 进程0.0834 s8 个进程363.1 MiB

bytes已经是对序列化相当友好的连续数据,进程传输仍要复制。若换成几十万个嵌套 Python 对象,编码、对象重建和内存开销通常会更明显。

这里的 RSS 是参与执行的 Worker 各自ru_maxrss峰值之和,只用于观察量级。进程池一行没有计入仍持有原始对象的父进程,同时又可能重复计算共享库页面;它也不是某一时刻的精确物理内存。要做容量规划,应在实际容器中观察完整进程组的 RSS、PSS 和 cgroup 内存。

这些数字只属于这台机器、这个 Python 版本和这组任务。真正可迁移的结论是增长形状:

  • 等待型任务能从线程或协程并发中受益。
  • 纯 Python CPU 任务在线程和协程中没有多核加速。
  • 进程池要有足够任务粒度,并尽量复用。
  • 进程并行的收益要扣除启动、序列化和内存。

如何做选择,我更看重这几条

如果现有代码使用同步 SDK,任务主要等网络,并发量几十到几百,线程池通常最务实。改动小,库兼容性也好。

如果服务从入口到数据库、HTTP 客户端都是异步,存在大量同时等待的连接,而且团队能正确处理超时、取消和资源关闭,协程更合适。

如果热点是纯 Python 计算,任务可以拆分且粒度足够大,进程池值得尝试。先控制进程数,检查参数大小,再看稳态吞吐。

如果热点位于 NumPy、压缩或图像库,不要只凭“CPU 密集”决定。确认库是否释放 GIL,线程有时能避免进程复制并获得并行。

如果任务只执行几十次、每次几毫秒,串行可能已经是最好的方案。并发框架本身也需要时间。

可以用下面这张表做起点:

任务特征优先考虑先检查的代价
同步阻塞 I/O,中等并发线程池线程安全、连接池、并发上限
异步 I/O,大量连接协程调用链兼容、阻塞代码、取消与超时
纯 Python CPU 计算进程池启动、序列化、内存、任务粒度
释放 GIL 的 C 扩展计算线程或库自身并行底层线程数、过度并行
大对象共享与少量计算线程或串行竞态、锁竞争
极小且有限的任务串行或批处理并发开销可能更高

这不是决策树的终点。压测结果如果与预期不同,应该回到调用栈和资源指标,而不是继续机械增加 Worker。

并发以后,应该测什么

总耗时只是第一项。

一套有意义的对比还应观察:

  • 每秒完成任务数,以及输入增长后的吞吐曲线。
  • 单任务 P50、P95、P99 与排队时间。
  • CPU 总利用率和每个核心的使用情况。
  • 进程组内存、线程数、上下文切换。
  • 下游连接数、429、超时和错误率。
  • 队列长度、最老任务等待时间。
  • 失败时未完成任务能否重试或恢复。
  • 服务关闭后是否仍有线程、子进程或异步任务残留。

并发从 8 增到 16,吞吐提高 30%,也许值得;增到 64 后吞吐不变,P99 翻倍,说明瓶颈已经转移。那个拐点比某篇文章推荐的 Worker 数更有用。

测试还要包含共享状态。把结果数量对上不够,金额、去重与顺序都要验证。竞态问题通常不稳定,可以用Barrier、故意让出执行权和高重复次数扩大窗口,而不是跑一遍没出错就算通过。

一份并发模型检查清单

准备让任务“同时跑”之前,可以先回答:

  1. 当前瓶颈是 CPU、网络、磁盘、数据库,还是锁等待?
  2. 优化目标是单次延迟、整体吞吐,还是资源成本?
  3. 热点代码是纯 Python,还是会释放 GIL 的扩展库?
  4. 同步依赖是否成熟,改成异步要穿透多少层?
  5. 下游真正允许多少并发,连接池能提供多少资源?
  6. 线程之间有哪些共享可变状态,谁负责加锁?
  7. 锁内是否包含网络、磁盘或其他无界等待?
  8. 进程池是否长期复用,任务粒度能否覆盖启动成本?
  9. 传给子进程的参数有多大,是否能被pickle?
  10. 每个进程会复制哪些模型、缓存与连接池?
  11. 是否一次创建了远超处理能力的任务对象?
  12. 队列是否有上限,满了以后压力传到哪里?
  13. 异步调用链中是否混入阻塞 I/O 或纯 CPU 长任务?
  14. 请求上下文是否使用ContextVar,跨进程时是否显式传递?
  15. 部署环境使用fork、spawn还是其他启动方式?
  16. Worker 异常退出时,任务会丢失、重试还是重复?
  17. 关闭程序时,正在运行的工作如何收尾?
  18. 基准是否包含真实参数大小、初始化和稳态运行?

如果这些问题只答得出“先设成 CPU 核心数乘二看看”,那还不是并发设计,只是一轮参数试探。

结语

回到开头的订单对账。

最后并没有选出一个模型包办所有事情。渠道查询继续用有限线程池,因为 SDK 是同步的,改造成本低;规则计算按批次交给长期复用的进程池;汇总不让 Worker 共同修改一份字典,而是在主进程合并局部结果。数据库写入仍然受独立连接池约束。

这个方案不如“全异步”或“开满多核”听起来漂亮,却更贴合那条任务真实的时间分布。

线程、进程和协程其实都在做一件事:当当前工作无法或不值得继续占住执行资源时,让其他工作向前走。区别在于它们在哪里切换、共享什么、隔离什么,又为此付出多少成本。

所以选择时别先问“哪一个性能最好”。先看程序在忙还是在等,再看数据要不要跨边界、状态能不能共享、下游能接住多少。答案常常没有那么戏剧化:几十个阻塞请求用线程就够了,一段纯 Python 计算交给进程,已经是异步调用链的服务则继续用协程。

真正麻烦的,从来不是把任务创建出来,而是任务超时、失败或被取消以后,谁负责把剩下的资源收回来。

下一期,我们会沿着这句话深入asyncio。事件循环怎么调度 Task,超时如何触发取消,CancelledError为什么不能随便吞掉,以及怎样用结构化并发确保一个请求结束时,不会在后台留下仍然运行的协程。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/10/3 19:32:48

视频动态目标三维重构在危化品事故平战切换指挥中的应用技术解析

技术权属说明:危化品事故平战一体化态势感知、常态/应急无缝切换指挥机制、动态目标三维态势联动调度、极端工况指挥闭环技术体系由华东师范大学浙江普陀时空大数据研究院耿文海团队原创研发,镜像视界(浙江)科技有限公司为唯一产业…

作者头像 李华
网站建设 2026/10/3 19:32:45

MQTT从入门到实战:Windows搭建与485设备接入指南

1. 先弄明白MQTT的底层逻辑:它凭什么成为物联网事实标准做物联网项目做得久了,你会发现一个很有意思的现象:不管是做智能家居、工业数据采集,还是做智慧农业、车联网,大家最后都会不约而同地选MQTT来跑业务消息。我在几…

作者头像 李华
网站建设 2026/10/3 19:29:18

单视频三维重构赋能化工装置泄漏扩散三维态势推演技术解析

技术权属说明:化工泄漏气云三维重构、扩散态势时空推演、单视频抗扰感知推演体系由华东师范大学浙江普陀时空大数据研究院耿文海团队原创研发,镜像视界(浙江)科技有限公司为唯一产业化落地主体,具备完整自主知识产权。…

作者头像 李华