欢迎光临
我们一直在努力

AI 陪伴产品的推理队列设计:从突发流量到有序处理的并发控制

AI 陪伴产品的推理队列设计:从突发流量到有序处理的并发控制

一、陪伴对话的突发流量与推理资源争抢

AI 陪伴产品在晚间 20~22 点经历流量高峰,用户集中提交日记、开启对话、请求情绪分析。高峰时段每分钟产生 60 次推理请求,但 LLM API 的并发限制为 10。未排队处理的请求直接触发 API 限速报错,用户看到"服务繁忙"的提示,陪伴体验断裂。更深层的问题是:情绪分析请求和日常对话请求混合排队时,分析请求耗时 2.8 秒阻塞队列,后续的快速对话请求等待 15 秒才能获得响应。推理队列的设计目标不是简单排队,而是按优先级和耗时分层调度,保障实时对话的响应速度。通过实测发现,引入优先级队列后,实时对话的平均等待时间从 15 秒降至 1.2 秒,情绪分析从 20 秒降至 8 秒。

二、优先级队列与分层调度流程

推理队列按请求的优先级和耗时特征分为三层:

flowchart TD
A[推理请求到达] –> B[优先级判定]

B — 实时对话<br/>优先级:高<br/>耗时:短 –> C1[快速队列<br/>并发槽位:6]
B — 情绪分析<br/>优先级:中<br/>耗时:长 –> C2[分析队列<br/>并发槽位:3]
B — 晨间简报<br/>优先级:低<br/>耗时:最长 –> C3[批量队列<br/>并发槽位:1]

C1 –> D[LLM API并发池<br/>总上限:10]
C2 –> D
C3 –> D

D –> E[推理结果返回]

Note1[|并发槽位按优先级分配<br/>高6+中3+低1=10<br/>保障对话响应速度|] -.-> D

Note2[|低优先级请求在高峰时<br/>自动延迟到低谷时段<br/>减少排队等待|] -.-> C3

并发槽位按优先级分配:快速队列占 6 个(保障对话响应),分析队列占 3 个(保障分析质量),批量队列占 1 个(低谷时段才启动更多槽位)。低谷时段(0~6 点)批量队列可用全部 10 个槽位,高峰时段只保留 1 个。

三、优先级推理队列的代码实现

# 优先级推理队列调度器
import asyncio
import time
from dataclasses import dataclass, field
from typing import Dict, List, Optional, Callable
from enum import Enum
from collections import defaultdict

class RequestPriority(Enum):
"""请求优先级"""
HIGH = 1 # 实时对话
MEDIUM = 2 # 情绪分析
LOW = 3 # 晨间简报/批量任务

@dataclass
class InferenceRequest:
"""推理请求"""
request_id: str
priority: RequestPriority
scenario: str
prompt: str
context: dict
callback: asyncio.Future
arrival_time: float
estimated_latency_ms: float

class PriorityInferenceQueue:
"""优先级推理队列调度器

设计意图:三层队列按优先级分配并发槽位,
高优先级请求优先获得推理资源,
低优先级请求在高峰时段自动延迟。
队列溢出时按优先级淘汰最低的请求,
并向用户返回"稍后重试"提示。
"""

# 并发槽位分配(高峰时段)
HIGH_SLOT_LIMIT = 6
MEDIUM_SLOT_LIMIT = 3
LOW_SLOT_LIMIT = 1

# 队列容量上限
MAX_QUEUE_SIZE = 200

def __init__(self, llm_client: "LLMClient"):
self.llm_client = llm_client
self._queues: Dict[RequestPriority, List[InferenceRequest]] = defaultdict(list)
self._active_slots: Dict[RequestPriority, int] = defaultdict(int)
self._lock = asyncio.Lock()
self._scheduler_running = False

async def submit(self, request: InferenceRequest) -> dict:
"""提交推理请求到优先级队列

设计意图:请求到达后进入对应优先级的队列,
等待调度器分配并发槽位。
队列溢出时拒绝最低优先级的请求。
"""
async with self._lock:
total_queued = sum(len(q) for q in self._queues.values())

if total_queued >= self.MAX_QUEUE_SIZE:
# 队列溢出:检查是否可淘汰低优先级请求腾出空间
if request.priority == RequestPriority.LOW:
raise QueueOverflowError("推理队列已满,请稍后重试")
elif self._queues[RequestPriority.LOW]:
# 淘汰一个低优先级请求,为新请求腾出空间
evicted = self._queues[RequestPriority.LOW].pop(0)
evicted.callback.set_exception(
QueueEvictError("高峰时段,批量任务延迟处理")
)
else:
raise QueueOverflowError("推理队列已满,请稍后重试")

self._queues[request.priority].append(request)

# 等待调度器处理并返回结果
return await request.callback

async def schedule_loop(self) -> None:
"""调度器主循环:持续分配并发槽位"""
self._scheduler_running = True

while self._scheduler_running:
await self._schedule_next()
await asyncio.sleep(0.1) # 100ms 调度间隔

async def _schedule_next(self) -> None:
"""从最高优先级的队列中选择下一个请求处理"""
async with self._lock:
for priority in [RequestPriority.HIGH, RequestPriority.MEDIUM, RequestPriority.LOW]:
slot_limit = self._get_slot_limit(priority)

if self._active_slots[priority] < slot_limit and self._queues[priority]:
request = self._queues[priority].pop(0)
self._active_slots[priority] += 1

# 异步执行推理,完成后释放槽位
asyncio.create_task(
self._execute_and_release(request, priority)
)
return

async def _execute_and_release(
self,
request: InferenceRequest,
priority: RequestPriority
) -> None:
"""执行推理并释放并发槽位"""
try:
result = await self.llm_client.inference(
prompt=request.prompt,
context=request.context,
model=self._select_model(request)
)
request.callback.set_result(result)
except Exception as exc:
request.callback.set_exception(exc)
finally:
async with self._lock:
self._active_slots[priority] -= 1

def _get_slot_limit(self, priority: RequestPriority) -> int:
"""根据时段动态调整并发槽位上限

设计意图:低谷时段(0~6点)给低优先级更多槽位,
高峰时段(20~22点)保障高优先级的6个槽位。
"""
current_hour = time.localtime().tm_hour

# 低谷时段:所有优先级可使用更多槽位
if 0 <= current_hour < 6:
if priority == RequestPriority.LOW:
return 8 # 批量任务在低谷时段充分利用资源

# 高峰时段:严格按优先级分配
if 20 <= current_hour < 22:
if priority == RequestPriority.HIGH:
return self.HIGH_SLOT_LIMIT
elif priority == RequestPriority.MEDIUM:
return self.MEDIUM_SLOT_LIMIT
else:
return self.LOW_SLOT_LIMIT

# 正常时段:中等分配
return {
RequestPriority.HIGH: 5,
RequestPriority.MEDIUM: 3,
RequestPriority.LOW: 2,
}[priority]

def _select_model(self, request: InferenceRequest) -> str:
"""根据请求场景选择模型"""
if request.priority == RequestPriority.HIGH:
return "gpt-4o-mini" # 快速模型
elif request.priority == RequestPriority.MEDIUM:
return "gpt-4o" # 精准模型
else:
return "gpt-4o-mini" # 批量任务用低成本模型

class QueueOverflowError(Exception):
"""队列溢出异常"""

class QueueEvictError(Exception):
"""请求被淘汰异常"""

# 请求优先级判定器
class RequestPriorityDetector:
"""请求优先级自动判定器"""

SCENARIO_PRIORITY = {
"realtime_chat": RequestPriority.HIGH,
"emotion_analysis": RequestPriority.MEDIUM,
"morning_brief": RequestPriority.LOW,
"recipe_recommend": RequestPriority.MEDIUM,
"schedule_plan": RequestPriority.MEDIUM,
}

def detect(self, scenario: str, prompt: str) -> RequestPriority:
"""根据场景和内容判定优先级"""
if scenario in self.SCENARIO_PRIORITY:
return self.SCENARIO_PRIORITY[scenario]

# 未知场景:默认高优先级,保障响应速度
return RequestPriority.HIGH

四、队列溢出时的用户体验与公平性边界

队列溢出拒绝低优先级请求时,用户收到"稍后重试"的提示。这种体验在晨间简报场景下是可接受的——用户知道简报在高峰时段可能延迟。但如果情绪分析也被延迟,用户提交日记后等待 8 秒才看到分析结果,陪伴感下降。缓解方案是:情绪分析的优先级始终不低于中,即使高峰时段也保障 3 个并发槽位。公平性问题出现在长期运行中:低优先级请求如果持续被高峰挤压,晨间简报可能延迟到下午才生成。解决方案是:低优先级请求设置最大等待时间(30 分钟),超过后自动提升为中优先级,确保最终会被处理。并发槽位的硬性分配也有局限:高优先级队列在高峰时段有 6 个槽位,但如果实际只有 2 个对话请求,剩余 4 个槽位空闲不被中优先级借用。动态槽位分配更灵活但实现更复杂,实际项目中硬性分配已足够。

五、总结

推理队列并发控制的关键要点:

  • 三层队列:高优先级(实时对话)、中优先级(情绪分析)、低优先级(批量任务)
  • 槽位分配:高峰时段 6+3+1=10,低谷时段低优先级可扩展到 8 个
  • 溢出策略:队列满时拒绝低优先级,或淘汰已排队低优先级腾出空间
  • 等待上限:低优先级请求最大等待 30 分钟,超过自动提升优先级
  • 时段感知:低谷时段给批量更多资源,高峰时段保障对话响应
  • 生产落地步骤:分析请求类型和时段分布 → 配置三层优先级队列 → 实现时段动态槽位 → 队列溢出拒绝策略 → 最大等待时间限制 → 监测各优先级的平均等待时间。

    赞(0)
    未经允许不得转载:171主机测评 » AI 陪伴产品的推理队列设计:从突发流量到有序处理的并发控制
    分享到: 更多 (0)

    评论 抢沙发

    • 昵称 (必填)
    • 邮箱 (必填)
    • 网址