Skip to content

Python 队列与流水线:让多阶段任务稳定推进

对应视频:EP10《生产者-消费者与任务队列》
视频入口:B站 EP10 定时稿(计划公开:2026-08-14 09:47)。 对应合集:Python 并发实战:从基础模型到语音工程。 本文为视频的工程展开版,补充代码模板、排查链路、状态字段和检查清单。 最后核验:2026/08/09 08:26:33

这篇文章是视频的工程展开版。视频先说明队列是什么、为什么需要它,再进入批量语音处理示例。文章补充完整模板、停止信号、失败隔离和检查清单。

有界队列流水线结构图

这张图把有界队列放在阶段之间:队列既传递任务,也传递压力,下游慢时上游必须感知并停下来。

这篇解决什么问题

队列在并发工程里不是简单的“列表”。它更常见的作用,是把一个多阶段任务拆成几个稳定边界:上游生产任务,下游消费任务,中间用队列传递数据和压力。

批量语音处理通常不是一个动作,而是一条链路:

text
扫描文件 -> 转码/预处理 -> ASR 调用 -> 结果写入

这些阶段的瓶颈不一样。扫描可能卡磁盘,转码可能卡 CPU 和外部命令,ASR 可能卡网络和配额,写入关心顺序、幂等和 checkpoint。

如果把所有阶段丢进同一个线程池,程序看起来并发了,但压力没有边界。最快的阶段会不断制造任务,最慢的阶段会让内存、重试和失败记录一起堆起来。

视频核心判断

多阶段任务不能只靠一个并发池。要用有界队列把阶段拆开,让每个阶段有自己的 worker 数、失败处理和观测指标。

机制展开

元素作用工程意义
Producer扫描文件、读取任务表、产生任务不直接做重活,只负责稳定产出
Queue连接两个阶段阶段边界,也是压力信号
Consumer从队列取任务并处理可以按阶段配置 worker 数
Bounded queue限制队列容量下游变慢时让上游等待
Sentinel停止信号保证 worker 能退出
Writer单独写结果集中处理幂等、checkpoint 和失败记录

最小工程代码

下方保留可直接阅读的示例代码。

关键结构是三段队列:

python
scan_to_transcode = queue.Queue(maxsize=3)
transcode_to_asr = queue.Queue(maxsize=2)
asr_to_write = queue.Queue(maxsize=8)

maxsize 不是随便填的,它代表阶段之间允许积压多少任务。ASR 如果有限流,transcode_to_asr 就不应该无限增长。

停止信号

队列流水线最容易卡死在停止阶段。

错误写法现象根因修正
不发送停止信号worker 一直阻塞消费者不知道上游结束生产者结束后发送 sentinel
多个 worker 只发一个 sentinel只有一个 worker 退出停止信号数量不足按下游 worker 数传播
异常后不 task_donejoin 永远不返回队列计数没有归还finally 里调用
writer 多处写文件结果缺行或乱序写入状态分散单写入者集中落盘

工程补充:怎么把这篇用到项目里

读图方式

  • 扫描、转码、ASR、写入最好拆成不同阶段,因为每个阶段的瓶颈和并发上限不同。
  • Queue(maxsize) 是容量预算,不是随手填的数字;它决定系统能承受多少排队。
  • sentinel、状态表和 checkpoint 决定流水线能不能优雅停止和恢复。

排查路径

  • 持续记录每个队列的 qsize、put 等待时间、get 等待时间和每阶段吞吐。
  • 如果某个队列持续增长,下游就是瓶颈;如果上游频繁 put 等待,背压已经生效。
  • 中断恢复时核对状态表:DONE 跳过、FAILED 重试或隔离、PROCESSING 超时回收。

问题、证据和解决动作

问题现场先查什么解决动作
扫描太快撑爆内存待处理任务数、RSS有界队列 + 分批扫描
转码慢拖住 ASR转码队列长度、CPU/磁盘 util拆独立 worker 池,限制转码并发
写入重复或漏写状态表、输出文件校验单写入者 + 幂等写入
停止时丢任务sentinel 数量、未完成计数task_done / join / checkpoint

落地边界

  • 队列不是越大越好;队列过大只会把问题延迟暴露并放大恢复成本。
  • 每个阶段都要能独立降并发、暂停和重试。

常见失败模式

错误写法现象根因修正
无界队列内存持续上涨下游压力被隐藏使用有界队列
一个池跑所有阶段某阶段拖垮全局不同资源混在一起抢 worker按阶段拆 worker
失败直接抛出一条坏文件中断全局错误没有任务级隔离转成失败结果进入状态表
不记录队列长度只知道慢,不知道哪里慢缺少阶段指标定时记录 qsize
重跑时从头开始成本和时间浪费没有 checkpointwriter 维护状态表

实战检查清单

  • 每个阶段的主要瓶颈是否写清楚。
  • 每个队列是否有 maxsize
  • 每个 worker 是否能收到停止信号。
  • 异常是否被转成任务状态,而不是直接杀掉进程。
  • 是否有单独 writer 负责结果、失败和 checkpoint。
  • 是否记录队列长度、成功数、失败数和阶段耗时。
  • 是否能从 failed 状态恢复,而不是从头重跑。

读完之后能完成什么

  • 能把一个大循环拆成生产者-消费者流水线。
  • 能写出有界队列、停止信号和单写入者模板。
  • 能通过队列长度判断哪个阶段正在成为瓶颈。
  • 能为后续限流、背压和重试打好结构基础。

结尾总结

这一集真正要留下的不是一个 Queue API,而是一种组织并发系统的方式。

先看问题:批量音频处理里,扫描、转码、ASR 和写入的瓶颈不同。它们如果挤进同一个并发池,资源竞争、错误传播和重试压力会混在一起。

再看结构:用 Queue 把阶段拆开,队列既是边界,也是压力信号。下游变慢时,队列长度通常比总耗时更早暴露问题;有界队列还能把压力传回上游。

最后看工程落地:worker 固定做“取任务、处理、交给下一段”,单独 writer 收敛结果、失败和 checkpoint。停止阶段要处理好 sentinel、task_donejoin;失败阶段要把错误变成任务状态;运行阶段要持续观察 qsize、成功数、失败数、平均耗时和 P95。

当这些结构稳定以后,下一步才是限流、背压、重试和更细的调参。

配套示例代码

源码文件:examples/ep10/bounded_pipeline.py

python
from __future__ import annotations

import queue
import threading
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Callable, Iterable


SENTINEL = object()


@dataclass(frozen=True)
class AudioJob:
    path: Path
    attempt: int = 1


@dataclass(frozen=True)
class StageResult:
    path: Path
    status: str
    detail: str


def scan_files(paths: Iterable[Path], output: queue.Queue[AudioJob | object]) -> None:
    for path in paths:
        output.put(AudioJob(path))
    output.put(SENTINEL)


def worker(
    name: str,
    input_queue: queue.Queue[AudioJob | object],
    output_queue: queue.Queue[AudioJob | StageResult | object],
    handler: Callable[[AudioJob], AudioJob | StageResult],
    worker_count_next: int = 1,
) -> None:
    while True:
        item = input_queue.get()
        try:
            if item is SENTINEL:
                for _ in range(worker_count_next):
                    output_queue.put(SENTINEL)
                return
            result = handler(item)
            output_queue.put(result)
        except Exception as exc:
            output_queue.put(StageResult(item.path, "failed", f"{name}: {exc!r}"))
        finally:
            input_queue.task_done()


def transcode(job: AudioJob) -> AudioJob:
    time.sleep(0.05)
    return AudioJob(job.path.with_suffix(".wav"), job.attempt)


def call_asr(job: AudioJob) -> StageResult:
    time.sleep(0.08)
    if "bad" in job.path.name:
        return StageResult(job.path, "failed", "asr rejected bad audio")
    return StageResult(job.path, "success", "transcript-id")


def writer(input_queue: queue.Queue[StageResult | object], stop_count: int) -> None:
    stopped = 0
    while True:
        item = input_queue.get()
        try:
            if item is SENTINEL:
                stopped += 1
                if stopped >= stop_count:
                    return
                continue
            print(f"{item.status:7s} {item.path.name:14s} {item.detail}")
        finally:
            input_queue.task_done()


def main() -> None:
    source_files = [Path(f"audio_{index}.mp3") for index in range(6)] + [Path("bad_audio.mp3")]

    scan_to_transcode: queue.Queue[AudioJob | object] = queue.Queue(maxsize=3)
    transcode_to_asr: queue.Queue[AudioJob | StageResult | object] = queue.Queue(maxsize=2)
    asr_to_write: queue.Queue[StageResult | object] = queue.Queue(maxsize=8)

    threads = [
        threading.Thread(target=scan_files, args=(source_files, scan_to_transcode), name="scanner"),
        threading.Thread(
            target=worker,
            args=("transcode", scan_to_transcode, transcode_to_asr, transcode, 2),
            name="transcode-0",
        ),
        threading.Thread(target=worker, args=("asr", transcode_to_asr, asr_to_write, call_asr), name="asr-0"),
        threading.Thread(target=worker, args=("asr", transcode_to_asr, asr_to_write, call_asr), name="asr-1"),
        threading.Thread(target=writer, args=(asr_to_write, 2), name="writer"),
    ]

    for thread in threads:
        thread.start()
    for thread in threads:
        thread.join()


if __name__ == "__main__":
    main()
别急,先让缓存热一下。