ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

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

第 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:int8,)-list[TrackingStatus]:withThreadPoolExecutor(max_workersmax_workers,thread_name_prefixtracking,)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:currentcustomer_totals.get(customer_id,0)customer_totals[customer_id](currentamount_cents)两个线程可能同时读到旧值 100线程 A读取 100 线程 B读取 100 线程 A写入 130 线程 B写入 150正确结果应是 180最终却只留下 150。GIL 只能约束某一时刻谁执行字节码不能保证“读取、计算、写回”这一组业务动作不可分割。可以用锁保护这段临界区fromthreadingimportLock customer_totals:dict[int,int]{}customer_totals_lockLock()defadd_order_amount(customer_id:int,amount_cents:int,)-None:withcustomer_totals_lock:currentcustomer_totals.get(customer_id,0)customer_totals[customer_id](currentamount_cents)锁内只保留共享状态更新不要顺手把网络调用也放进去withcustomer_totals_lock:responselogistics_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:contextmultiprocessing.get_context(spawn)withProcessPoolExecutor(max_workers8,mp_contextcontext,)asexecutor:resultslist(executor.map(calculate_risk_score,load_order_snapshots(),chunksize50,))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_workers8)asexecutor:resultslist(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_loopasyncio.get_running_loop()returnawaitevent_loop.run_in_executor(process_pool,calculate_risk_score,order,)边界仍然存在order要被序列化取消等待这个 Future 也不一定能立即终止已经在子进程里运行的函数。协程只是负责等待进程结果没有消除进程池的成本。限制并发不要一次创建所有任务异步代码很容易写出这种版本resultsawaitasyncio.gather(*(fetch_tracking_status(order_id)fororder_idinone_million_order_ids))事件循环不会同时执行一百万个 Python 指令但这里会一次创建大量协程与 Task保存参数、状态和结果。内存先涨起来下游也可能在很短时间内收到远超容量的请求。任务数量有限时可以用 Semaphore 约束在途请求asyncdeffetch_many_tracking_statuses(order_ids:list[str],*,concurrency:int20,)-list[TrackingStatus]:semaphoreasyncio.Semaphore(concurrency)asyncdeffetch_one(order_id:str)-TrackingStatus:asyncwithsemaphore:returnawaitfetch_tracking_status(order_id)returnawaitasyncio.gather(*(fetch_one(order_id)fororder_idinorder_ids))这限制了同时进入外部调用的数量但仍然一次创建了与输入等量的协程。面对持续流或海量输入更适合有界队列fromdataclassesimportdataclassfromtypingimportOptionaldataclass(frozenTrue)classTrackingOutcome:order_id:strstatus:Optional[TrackingStatus]error_message:Optional[str]asyncdeftracking_worker(queue:asyncio.Queue[Optional[str]],outcomes:list[TrackingOutcome],)-None:whileTrue:order_idawaitqueue.get()try:iforder_idisNone:returntry:statusawaitfetch_tracking_status(order_id)outcomeTrackingOutcome(order_idorder_id,statusstatus,error_messageNone,)exceptExceptionaserror:outcomeTrackingOutcome(order_idorder_id,statusNone,error_message(f{type(error).__name__}:{error}),)outcomes.append(outcome)finally:queue.task_done()生产者和入口负责有界写入fromcollections.abcimportAsyncIterableasyncdeffetch_from_stream(order_ids:AsyncIterable[str],*,worker_count:int20,queue_size:int100,)-list[TrackingOutcome]:queue:asyncio.Queue[Optional[str]]asyncio.Queue(maxsizequeue_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_exceptionsTrue,)queue_size控制已读取但还没处理的数据worker_count控制外部调用并发。队列满时生产者停下来这就是最直接的背压。这段示例把每个失败记录成字符串结果既避免 Worker 提前退出并把queue.join()永久卡住也不会长期保留异常 traceback 引用的局部对象。outcomes本身仍会随输入增长真正的海量任务应把结果分批写出而不是全部留到函数返回。生产系统还要定义哪些异常可重试、结果是否按输入顺序返回。那些细节不能靠gather()默认替我们决定。请求上下文不能放在线程级全局变量里异步服务中的多个请求通常运行在同一线程。若把请求 ID 放进普通全局变量或只按线程保存协程切换后会相互覆盖。请求级上下文应使用ContextVarfromcontextvarsimportContextVarfromtypingimportOptional request_id_context:ContextVar[Optional[str]]ContextVar(request_id,defaultNone,)asyncdefhandle_request(request:Request)-Response:tokenrequest_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_idCPU 任务只执行 Python 整数运算不调用可能释放 GIL 的第三方扩展defpure_python_cpu_task(seed:int,iterations:int,)-int:valueseedforindexinrange(iterations):value(value*1_664_5251_013_904_223index)0xFFFFFFFFreturnvalue测试环境是 Python 3.9.6、Apple M5 Pro、18 个逻辑 CPU。线程、进程和协程的并发上限都设为 8每个实现返回相同结果后才记录时间表格使用三次运行的中位数。第一组是 40 个等待任务执行方式总耗时串行1.3351 s8 线程0.1725 s8 进程每轮冷启动2.1486 s协程最多 8 个在途0.1551 s理论等待时间是40 × 0.03 1.2秒。线程和协程把等待叠成约五批接近5 × 0.03 0.15秒。进程当然也能同时睡眠但 macOSspawn启动解释器的成本远高于这点工作没有使用价值。第二组是 24 个纯 Python 计算任务每个执行 100 万轮执行方式总耗时串行1.2350 s8 线程1.2537 s8 进程每轮冷启动2.1653 s24 个 CPU 协程1.2572 s线程和协程都没有多核收益。协程版本内部没有await只是依次占用事件循环。冷启动进程池仍然太贵。接着复用已经预热的池把单任务计算提高到 200 万轮执行 16 个任务执行方式总耗时串行1.7822 s8 个复用线程1.6412 s8 个复用进程0.2192 s线程比串行少的那一点不应解读成稳定加速。CPU 调频、系统噪声和中位数样本都可能造成小幅波动两者仍在同一量级。复用进程池后启动成本不再进入每批任务纯 Python 计算才真正分布到多个核心。最后传递同一个 5 MiBbytes对象给 16 个任务。池先预热任务内部等待 30 毫秒确保所有 Worker 都参与执行方式总耗时参与执行单元Worker 峰值 RSS 之和8 线程0.0683 s1 个进程32.2 MiB8 进程0.0834 s8 个进程363.1 MiBbytes已经是对序列化相当友好的连续数据进程传输仍要复制。若换成几十万个嵌套 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、故意让出执行权和高重复次数扩大窗口而不是跑一遍没出错就算通过。一份并发模型检查清单准备让任务“同时跑”之前可以先回答当前瓶颈是 CPU、网络、磁盘、数据库还是锁等待优化目标是单次延迟、整体吞吐还是资源成本热点代码是纯 Python还是会释放 GIL 的扩展库同步依赖是否成熟改成异步要穿透多少层下游真正允许多少并发连接池能提供多少资源线程之间有哪些共享可变状态谁负责加锁锁内是否包含网络、磁盘或其他无界等待进程池是否长期复用任务粒度能否覆盖启动成本传给子进程的参数有多大是否能被pickle每个进程会复制哪些模型、缓存与连接池是否一次创建了远超处理能力的任务对象队列是否有上限满了以后压力传到哪里异步调用链中是否混入阻塞 I/O 或纯 CPU 长任务请求上下文是否使用ContextVar跨进程时是否显式传递部署环境使用fork、spawn还是其他启动方式Worker 异常退出时任务会丢失、重试还是重复关闭程序时正在运行的工作如何收尾基准是否包含真实参数大小、初始化和稳态运行如果这些问题只答得出“先设成 CPU 核心数乘二看看”那还不是并发设计只是一轮参数试探。结语回到开头的订单对账。最后并没有选出一个模型包办所有事情。渠道查询继续用有限线程池因为 SDK 是同步的改造成本低规则计算按批次交给长期复用的进程池汇总不让 Worker 共同修改一份字典而是在主进程合并局部结果。数据库写入仍然受独立连接池约束。这个方案不如“全异步”或“开满多核”听起来漂亮却更贴合那条任务真实的时间分布。线程、进程和协程其实都在做一件事当当前工作无法或不值得继续占住执行资源时让其他工作向前走。区别在于它们在哪里切换、共享什么、隔离什么又为此付出多少成本。所以选择时别先问“哪一个性能最好”。先看程序在忙还是在等再看数据要不要跨边界、状态能不能共享、下游能接住多少。答案常常没有那么戏剧化几十个阻塞请求用线程就够了一段纯 Python 计算交给进程已经是异步调用链的服务则继续用协程。真正麻烦的从来不是把任务创建出来而是任务超时、失败或被取消以后谁负责把剩下的资源收回来。下一期我们会沿着这句话深入asyncio。事件循环怎么调度 Task超时如何触发取消CancelledError为什么不能随便吞掉以及怎样用结构化并发确保一个请求结束时不会在后台留下仍然运行的协程。
返回列表