跳转到正文
作者

1. 引言:并发下的隐形杀手

上一篇文章中,通过 multiprocessing 实现了多模型实例的并行推理。但在实际压测和流式业务上线后,我发现了新的性能瓶颈:

  1. IPC 通信开销:当并发量上来后,主进程与 Worker 进程间通过 Queue 传递大量音频数据,序列化(Pickle)占用了大量 CPU,导致 RTF(实时率)波动。
  2. IO 阻塞雪崩:高并发上传文件或下载音频时,任何同步 IO 操作都会阻塞 EventLoop,导致 WebSocket 心跳丢失。
  3. 资源争抢:没有精细的限流,瞬间涌入的请求会打爆显存或导致 OOM。

本文将分享我是如何组合使用 SharedMemoryAnyIOAsyncio 全家桶,处理上述问题。文末记录的是整体方案对比,包含设备数量变化,不能据此分离各项软件优化的贡献。


2. 核心手术:SharedMemory 替换 Queue

这一调整旨在减少进程间序列化和大块音频传输开销,将音频传输方式从队列负载改为共享内存引用。代码仍包含数据写入和复制,不是整个链路完全零拷贝;当前实验也未单独测量这一项的收益。

2.1 主进程:写入与指针传递

不再直接把音频数据(Bytes)塞入队列,而是申请一块共享内存,将数据写入 buffer,通过队列仅传递内存块的名称(name)。

import multiprocessing.shared_memory
import numpy as np

async def process_chunk(self, audio_array: np.ndarray, is_final: bool) -> str:
    audio_size = audio_array.nbytes

    # 1. 创建共享内存 (Zero-Copy 基础)
    shm = multiprocessing.shared_memory.SharedMemory(
        create=True,
        size=audio_size,
        name=f"asr_{task_id}"
    )

    # 2. 写入数据
    audio_bytes = audio_array.tobytes()
    shm.buf[:audio_size] = audio_bytes

    # 3. 引用管理(防止被过早 GC)
    async with self.factory._shared_memory_lock:
        self.factory._shared_memory_pool[shm.name] = shm

    # 4. 仅发送“指针”
    task = {
        "task_id": task_id,
        "shm_name": shm.name,
        "audio_size": audio_size,
        # ... 其他元数据
    }
    self._task_queue.put(task)

2.2 Worker 进程:映射读取与本地复制

Worker 收到任务后,根据 shm_name 映射内存区域。下面的实现会通过 .copy() 创建本地数组,然后关闭共享内存引用,交由主进程进行 unlink 释放;因此不能将这段实现称为零拷贝。

def worker_process():
    # ... 获取任务 ...
    if shm_name and audio_size > 0:
        # 直接连接已存在的共享内存
        shm = multiprocessing.shared_memory.SharedMemory(name=shm_name)
        # 映射 buffer 后复制为本地数组,再关闭共享内存引用
        audio_array = np.frombuffer(shm.buf[:audio_size], dtype=np.float32).copy()
        shm.close()
    # ... 模型推理 ...

3. 防线构筑:多级并发控制

有了高效的数据传输通道,还需要红绿灯来控制流量。我在不同层级使用了不同的并发原语。

3.1 入口层:Asyncio Semaphore 限流

为了保护服务不被瞬间流量冲垮,在 HTTP 和 WebSocket 入口处均设置了信号量。

WebSocket 连接控制:

# 限制最大并发连接数,超过直接拒绝或等待
_streaming_semaphore = asyncio.Semaphore(MAX_CONCURRENT_REQUESTS)

@router.websocket("/streaming")
async def transcription(websocket: WebSocket):
    try:
        # 快速失败机制:等待超时则拒绝服务
        await asyncio.wait_for(semaphore.acquire(), timeout=REQUEST_TIMEOUT)
    except asyncio.TimeoutError:
        await websocket.close(code=1008, reason="服务繁忙")
        return
    # ... 业务逻辑 ...

3.2 逻辑层:AnyIO 管道与任务组

在处理流式音频时,我使用了 anyio 库。相比原生的 asyncio.Queue,anyio 的 Memory Object Stream 提供了更现代的生产者-消费者模型,且与 Structured Concurrency(结构化并发)结合得更好。

import anyio

# 创建有界内存流,防止内存无限增长
send_stream, receive_stream = anyio.create_memory_object_stream(max_buffer_size=100)

async with anyio.create_task_group() as tg:
    # 消费者:后台模型处理
    tg.start_soon(audio_processor_task, receive_stream)

    # 生产者:WebSocket 接收
    async with send_stream:
        while True:
            msg = await websocket.receive()
            await send_stream.send(audio_message)

3.3 数据层:线程安全与连接池

对于 HTTP 下载和文件写入等操作,必须严格隔离 async 和 sync 操作,防止阻塞事件循环。

  • HTTP 并发下载:使用 httpx.AsyncClient 复用连接池。
  • 文件 IO:使用 run_in_executor 将磁盘写入扔到线程池。
# 典型的异步避坑:不要在 async def 里直接写 open()
def _write_file():
    with tempfile.NamedTemporaryFile(delete=False) as tmp:
        tmp.write(content)
        return tmp.name

tmp_path = await loop.run_in_executor(None, _write_file)

4. 调度策略:负载均衡

当后端有多个 GPU Worker 时,如何分配任务?我实现了一个简单的负载均衡器,始终将任务派发给"当前待处理任务最少"的 Factory 实例。

@staticmethod
def select_from_factories(factories: List["ASRFactory"]) -> "ASRFactory":
    # 贪心算法:选择积压任务最少的 Worker 组
    return min(factories, key=lambda f: len(f._pending_tasks))

5. 效果实战:性能对比

为了验证上述优化的效果,我针对流式 ASR 场景进行了梯度压测。

测试环境:

  • 硬件:双卡 NVIDIA L20
  • 部署:双进程模型实例(每卡1实例)
  • 测试工具:自研 Python 脚本模拟 WebSocket 并发流

5.1 优化前 vs 优化后 (50并发)

50 路并发 下,以下为作者历史记录中的两个整体方案。基准为单卡单实例,修改方案为双卡双实例并叠加软件调整;该比较存在硬件和软件混杂因素。

指标单卡单实例 (基准)双卡双实例 + 优化 (Modified)提升幅度
RTF (实时率)0.610.23性能提升 2.6 倍
总耗时184.4s75.3s速度提升 2.4 倍
平均队列延迟1273ms187ms延迟降低 85%

注:RTF (Real Time Factor) 越低越好,0.23 意味着处理 1 秒的音频只需要 0.23 秒。

5.2 高并发档位的观测结果

在优化后的架构下,我们测试了从 30 到 50 并发的表现(见下图数据趋势):

  • 30 并发:RTF 1.28 (异常) -> 优化为 RTF 0.24
  • 40 并发:RTF 0.48 (单卡) -> 优化为 RTF 0.23
  • 50 并发:RTF 0.61 (单卡) -> 优化为 RTF 0.23

局限:这些记录说明修改后的整体方案在所测档位获得了更低的 RTF,但未提供同硬件消融、重复测量或误差区间。因此不能据此证明双卡线性扩展、持续稳定性,或不存在通信和锁竞争开销。下一步需固定硬件、音频集与模型版本,分别测量各项软件调整。


6. 总结

本次调整同时包含硬件扩容和软件改动。以下是设计意图,不是已通过独立消融确认的因果收益:

  • SharedMemory 用于减少跨进程音频传输负载,仍需测试与其他调整分离后的收益。
  • Semaphore & Lock 解决了秩序问题(资源竞争)。
  • AnyIO & ThreadPool 解决了阻塞问题(CPU 调度)。