位置:首页 > 进阶教程 > AI推理调度精细化:队列、批处理与弹性伸缩降本新方向

AI推理调度精细化:队列、批处理与弹性伸缩降本新方向

时间:2026-07-25  |  作者:318050  |  阅读:0

2026 年,大模型推理正在从“单次请求调用”走向“系统级调度”。

过去,很多 AI 应用的推理方式相当直接。

用户请求进来,应用直接调用模型接口。等着返回结果,再展示给用户。

这种模式在早期业务中跑得通。可一旦调用量上来,问题就接踵而至。

  • 请求突然暴涨,模型服务会不会排队?
  • 不同任务该不该区别对待,给个优先级?
  • 多个短请求能不能合并成一批处理?
  • GPU 资源有没有白白闲置?
  • 模型服务挂了,能不能自动切到备用节点?

大模型推理早已不是单纯的 API 调用,而是一个彻头彻尾的资源调度问题。

一套合格的 AI 推理调度系统,目标很明确:在保证响应速度的前提下,把资源利用率拉上去,把推理成本降下来,同时让服务更稳定。


一、为什么推理需要调度?

大模型推理资源可不便宜。如果所有请求都一股脑儿直接调用模型,系统很容易陷入两个困境:高峰期排队严重,用户等到没脾气;低峰期资源闲置,GPU 利用率低得可怜。

所以推理系统必须有点“调度智慧”。它得根据请求类型、优先级、最大等待时间、模型负载这些因素,决定请求是立刻执行、进队列等着、合并批处理,还是走一个降级模型。

下面就用 Python 手写一个简化版的 AI 推理调度系统,把这套逻辑跑一遍。


二、基础结构:定义推理请求

第一步,先定义请求长什么样。每个请求都带着用户、任务类型、Prompt、优先级和创建时间

这样调度系统就能统一管理不同任务。无论是摘要、分类、问答还是代码生成,通通能塞进同一个调度队列。

import time
import json
import random
from datetime import datetime
from collections import deque


class InferenceRequest:
    def __init__(
        self,
        request_id,
        user_id,
        task_type,
        prompt,
        priority=5
    ):
        self.request_id = request_id
        self.user_id = user_id
        self.task_type = task_type
        self.prompt = prompt
        self.priority = priority
        self.created_at = time.time()

    def to_dict(self):
        return {
            "request_id": self.request_id,
            "user_id": self.user_id,
            "task_type": self.task_type,
            "prompt": self.prompt,
            "priority": self.priority,
            "created_at": 30656.t.kuaisou.com
        }

三、模型节点:模拟推理实例

第二步,定义模型节点。每个节点都包含最大批量、当前负载、平均延迟和失败率

在生产环境中,它可能对应一个 GPU 服务、一个模型副本或一个推理容器。

class ModelNode:
    def __init__(
        self,
        node_id,
        model_name,
        max_batch_size,
        a vg_latency,
        fail_rate=0.05
    ):
        self.node_id = node_id
        self.model_name = model_name
        self.max_batch_size = max_batch_size
        self.a vg_latency = a vg_latency
        self.fail_rate = fail_rate
        self.running = 0
        self.total_finished = 0
        self.total_failed = 0

    def can_accept(self):
        return self.running < self.max_batch_size

    def infer_batch(self, requests):
        self.running += len(requests)
        time.sleep(self.a vg_latency)

        results = []

        for request in requests:
            if random.random() < self.fail_rate:
                self.total_failed += 1

                results.append({
                    "request_id": request.request_id,
                    "status": "failed",
                    "error": "model inference failed",
                    "model": self.model_name,
                    "node_id": self.node_id
                })
            else:
                self.total_finished += 1

                results.append({
                    "request_id": request.request_id,
                    "status": "success",
                    "answer": f"{self.model_name} 生成的模拟回答",
                    "model": self.model_name,
                    "node_id": self.node_id
                })

        self.running -= len(requests)

        return results

    def metrics(self):
        return {
            "node_id": self.node_id,
            "model_name": self.model_name,
            "running": self.running,
            "total_finished": self.total_finished,
            "total_failed": self.total_failed
        }

四、优先级队列:管理等待请求

第三步,构建队列。高优先级请求得优先处理,低优先级的可以等待或者进批处理。

队列是调度的基石。有了它,系统才能从容应对瞬时流量,而不是让每个请求都直接冲向模型服务。

class PriorityQueue:
    def __init__(self):
        self.items = []

    def push(self, request):
        self.items.append(request)

        self.items.sort(
            key=lambda item: (
                -item.priority,
                item.created_at
            )
        )

    def pop_batch(self, max_size):
        batch = self.items[:max_size]
        self.items = self.items[max_size:]

        return batch

    def size(self):
        return len(self.items)

    def oldest_wait_ms(self):
        if not self.items:
            return 0

        oldest = min(item.created_at for item in self.items)

        return int((time.time() - oldest) * 1000)

五、调度器:选择模型节点

第四步,定义调度器。它负责接收请求、选择合适的节点、构建批次,然后执行推理。

这是整个系统的核心,决定请求什么时候跑、交给哪个节点、能不能合并成一批。

class InferenceScheduler:
    def __init__(self, nodes):
        self.nodes = nodes
        self.queue = PriorityQueue()
        self.logs = []

    def submit(self, request):
        self.queue.push(request)

        self.logs.append({
            "event": "submit",
            "request_id": request.request_id,
            "priority": request.priority,
            "time": datetime.now().isoformat()
        })

    def choose_node(self):
        candidates = [
            node for node in self.nodes
            if node.can_accept()
        ]

        if not candidates:
            return None

        candidates.sort(
            key=lambda node: node.running
        )

        return candidates[0]

    def run_once(self):
        node = self.choose_node()

        if not node:
            self.logs.append({
                "event": "no_a vailable_node",
                "queue_size": self.queue.size(),
                "time": datetime.now().isoformat()
            })

            return []

        if self.queue.size() == 0:
            return []

        batch_size = min(
            node.max_batch_size,
            self.queue.size(otterly.cn)
        )

        batch = self.queue.pop_batch(batch_size)

        self.logs.append({
            "event": "dispatch_batch",
            "node_id": node.node_id,
            "model": node.model_name,
            "batch_size": len(batch),
            "time": datetime.now().isoformat()
        })

        results = node.infer_batch(batch)

        self.logs.append({
            "event": "batch_finished",
            "node_id": node.node_id,
            "results": results,
            "time": datetime.now().isoformat()
        })

        return results

六、弹性伸缩:根据队列长度扩容

第五步,加上简单的弹性伸缩。如果队列太长,系统就加新节点;队列空了,生产环境也可以缩容。

这是推理降本的关键能力。高峰期扩容、低峰期缩容,资源利用率自然就上去了。

def autoscale_nodes(scheduler):
    queue_size = scheduler.queue.size()
    oldest_wait_ms = scheduler.queue.oldest_wait_ms()

    if queue_size >= 5 or oldest_wait_ms > 3000:
        new_node_id = f"node-{len(scheduler.nodes) + 1}"

        new_node = ModelNode(
            node_id=new_node_id,
            model_name="GENERAL_MODEL",
            max_batch_size=3,
            a vg_latency=0.3,
            fail_rate=0.05
        )

        scheduler.nodes.append(new_node)

        scheduler.logs.append({
            "event": "scale_out",
            "new_node": new_node_id,
            "queue_size": queue_size,
            "oldest_wait_ms": oldest_wait_ms,
            "time": datetime.now().isoformat()
        })

七、生成运行报告

第六步,生成调度报告。内容包含队列长度、节点状态、日志和整体运行情况。

这份报告能帮团队判断推理系统是否健康。队列长期过长说明资源不够,节点长期空闲则说明资源可能浪费了。

def generate_scheduler_report(scheduler):
    return {
        "report_name": "AI 推理调度运行报告",
        "queue_size": scheduler.queue.size(),
        "oldest_wait_ms": scheduler.queue.oldest_wait_ms(),
        "nodes": [
            node.metrics()
            for node in scheduler.nodes
        ],
        "logs": scheduler.logs,
        "generate_time": datetime.now().isoformat()
    }

八、运行示例:模拟高峰请求

最后写一个运行入口,模拟多个请求涌进系统看看效果。

if __name__ == "__main__":

    nodes = [
        ModelNode(
            node_id="node-1",
            model_name="GENERAL_MODEL",
            max_batch_size=3,
            a vg_latency=0.3
        )
    ]

    scheduler = InferenceScheduler(nodes)

    for index in range(10):
        request = InferenceRequest(
            request_id=f"req-{index + 1}",
            user_id="user_001",
            task_type="summary",
            prompt=f"请总结第 {index + 1} 篇技术文章",
            priority=random.randint(1, 10)
        )

        scheduler.submit(request)

    while scheduler.queue.size() > 0:
        autoscale_nodes(scheduler)
        scheduler.run_once()

    report = generate_scheduler_report(scheduler)

    print(json.dumps(
        report,
        ensure_ascii=False,
        indent=2
    ))

九、趋势判断

从这套流程能清楚看到,大模型推理正在变成一个系统工程问题。过去开发者只需要关心怎么调用模型,未来团队得操心队列、批处理、优先级、弹性伸缩、失败率和资源利用率这些玩意儿。

随着 AI 应用调用量持续增长,推理调度只会越来越重要。它不仅影响响应速度,更直接关系到模型成本和系统稳定性。

可以预见的是,大模型应用以后不会只拼模型效果,还会拼推理基础设施能力。谁能把调度做精细、把资源用到位、把排队时间控住,谁就更可能把 AI 应用规模化落地。

来源:整理自互联网
免责声明:文中图文均来自网络,如有侵权请联系删除,心愿游戏发布此文仅为传递信息,不代表心愿游戏认同其观点或证实其描述。

相关文章

更多

精选合集

更多

大家都在玩

热门话题

大家都在看

更多