PERSONAL LAB / ai/agent开发
AI Agent 开发

14 AI 应用后端

Agent 在命令行里跑通,离给用户用还差一个后端:用户点一下按钮,后端要启动一个可能跑几分钟的任务,把每一步的进度实时推到网页上,同时不能让两个任务互相干扰,用户关掉网页也不能把服务搞乱。AI 应用岗的面试里,这部分和普通后端岗的问题很像,但有几个 AI 应用特有的点:流式输出、长任务、慢且贵的外部调用。

这篇要讲清楚:FastAPI 的基本用法;async defdef 到底有什么区别,什么时候会把整个服务卡住;SSE 流式输出的协议格式、浏览器怎么接、和 WebSocket 怎么选;长任务怎么放到后台跑、客户端断开时怎么办、并发怎么控制;生产环境要加的任务队列、Redis、Nginx 配置、心跳。用我项目 109 行的 server/app.py 逐段讲,包括一个”断开连接会提前放锁”的真实 bug。

怎么读:

前置:02 模型 API 工程(流式、超时、重试)、05 Agent 循环(生成器 yield 事件)。

第一部分 速记页

问题一句话答案
FastAPI 是什么Python 的 Web 框架,用类型注解声明参数,自动校验、自动生成接口文档
async def vs def 接口async def 在事件循环里跑,里面不能有阻塞调用;def 被放到线程池里跑,阻塞不影响其他请求
什么会卡住整个服务async def 里调用 time.sleep、同步的 requests、同步数据库驱动、跑 CPU 密集计算
流式输出为什么重要模型生成慢,逐步显示能让用户马上看到进展;Agent 任务几分钟,没进度用户会以为卡死
SSE 是什么服务器单向推送的 HTTP 长连接,响应类型 text/event-stream,每条消息 data: ... 加空行
SSE 字段data 内容、event 事件类型、id 事件编号、retry 重连间隔;冒号开头是注释,可当心跳
浏览器怎么接EventSource,断线会自动重连并带上 Last-Event-ID
SSE vs WebSocketSSE 单向、基于普通 HTTP、自动重连、实现简单;WebSocket 双向、要单独的协议升级
EventSource 的限制只能 GET、不能加自定义请求头;HTTP/1 下每个浏览器对同一域名最多 6 个连接
长任务怎么跑不在请求处理函数里直接跑完;放到后台线程或任务队列,通过队列把事件传给响应
客户端断开怎么办任务要不要继续由业务决定;并发锁不能跟着连接释放,要跟着任务结束释放
并发控制单进程用锁,拿不到返回 409;多进程、多机器要用 Redis 或数据库做锁
Nginx 反代要注意默认开启响应缓冲,SSE 要关(响应头 X-Accel-Buffering: no);读超时默认 60 秒,长时间没数据会断,要发心跳
任务队列Celery、RQ 等:请求只负责提交任务,后台工作进程执行,状态存 Redis 或数据库
Redis 常见用途缓存、限流计数、分布式锁、任务队列、会话、发布订阅
我项目的做法同步接口 + 后台线程 + queue.Queue + StreamingResponsethreading.Lock 同时只允许一个构建,否则 409
我项目修过的 bug锁原来在响应生成器结束时释放,客户端断开就提前放锁,后台构建还在跑时能再开一个

第二部分 易混对照

容易混的两个区别一句话记法
并发 vs 并行并发是多个任务交替推进(一个厨师照看几口锅);并行是真的同时执行(几个厨师)轮流 vs 同时
线程 vs 协程线程由操作系统调度,可以被随时切换;协程在 await 处主动让出,由事件循环调度被动切换 vs 主动让出
阻塞 vs 非阻塞阻塞调用等待期间占着当前线程;非阻塞调用等待期间把控制权交出去站着排队 vs 取号后去干别的
StreamingResponse vs EventSourceResponse前者是通用的流式响应,SSE 格式要自己拼;后者是 FastAPI 0.135 起专门的 SSE 响应,自动加心跳和响应头手写格式 vs 现成
SSE vs 轮询SSE 服务器有新数据就推;轮询是客户端定时问快递上门 vs 自己去问
SSE vs WebSocket单向 HTTP vs 双向独立协议广播 vs 电话
流式输出 vs 长任务流式是结果分段返回;长任务是执行时间长。Agent 两者都有说话慢 vs 活多
请求内执行 vs 任务队列请求内执行,进程重启任务就没了;任务队列把任务存下来,由独立的工作进程执行当场做 vs 下单排队
进程内锁 vs 分布式锁threading.Lock 只管同一个进程;多进程、多机器要用 Redis 等外部锁屋里的门闩 vs 大楼门禁
409 vs 429409 Conflict 表示和当前状态冲突(已有构建在跑);429 Too Many Requests 表示请求太频繁撞车 vs 超速

第三部分 面试口述稿

3.1 “你项目的后端是怎么设计的?”

后端是 FastAPI,一共一百来行,三个接口:一个启动构建并流式推送进度,两个查历史运行记录。

启动构建的接口是核心。Agent 跑一次要几十秒到几分钟,不能在请求里跑完再返回。我的做法是:请求进来先检查参数,然后尝试拿一把全局锁,拿不到说明已经有构建在跑,直接返回 409;拿到了就启动一个后台线程跑 Agent,Agent 每产生一个事件就放进一个队列,接口返回一个流式响应,响应的生成器从队列里取事件,按 SSE 格式写给浏览器。前端用 EventSource 接收,实时显示每一步调了什么工具、结果怎样。

为什么同时只允许一个构建:Agent 会设置进程级的环境变量,比如目标平台、toolchain 路径,两个构建同时跑会互相覆盖;工作目录也按项目名放,同名项目会冲突。这是单机项目的简化,要支持多用户就得每个任务独立的工作目录和环境,放到任务队列里跑。

3.2 “SSE 和 WebSocket 怎么选?”

看是不是需要双向通信。

Agent 进度推送、模型逐字输出,基本是服务器单向往浏览器推,SSE 就够了。它就是一个普通的 HTTP 响应,内容类型是 text/event-stream,每条消息是 data: 开头加一个空行,浏览器用 EventSource 接,断线会自动重连,还会带上最后收到的事件编号,服务端可以据此续传。不需要额外的协议,经过普通的反向代理和负载均衡也比较容易。

WebSocket 是双向的,适合用户在任务进行中也要频繁发消息的场景,比如语音对话、协同编辑。代价是要做协议升级,要自己处理心跳和重连。

SSE 也有限制:浏览器的 EventSource 只能发 GET、不能带自定义请求头,所以要传认证信息或者较大的输入时,常见做法是先用 POST 创建任务拿到任务 ID,再用 GET 订阅这个任务的事件流。另外 HTTP/1 下同一个域名浏览器最多 6 个连接,开很多标签页会不够用,HTTP/2 就没这个问题。

3.3 “FastAPI 里 async defdef 有什么区别?”

async def 的接口直接在事件循环里运行,事件循环是单线程的,靠 await 的时候让出控制权来同时处理很多请求。所以 async def 里面绝对不能有阻塞调用,比如 time.sleep、同步的 requests、同步的数据库驱动,否则整个事件循环都被卡住,所有请求排队。

def 的接口 FastAPI 会放到线程池里跑,里面阻塞不会卡住别的请求。

我做过一个实验,同样睡 0.5 秒,5 个并发请求:async def 里用 time.sleep,总共花了 2.5 秒,完全串行了;普通 def 里用 time.sleep,0.53 秒;async def 里用 await asyncio.sleep,0.51 秒。

所以原则是:用异步库(比如 httpx 的异步客户端、异步数据库驱动)就写 async def;用同步库就写 def;CPU 密集的计算放到进程池或者单独的工作进程。我项目的 Agent 用的是同步的 OpenAI 客户端和 subprocess,所以接口写的是 def,Agent 本身放在后台线程里跑。

3.4 “用户关掉网页,后台任务怎么办?”

先要决定业务上任务该不该继续。Agent 构建这种任务,已经花了很多 token,中途停掉很浪费,而且用户可能只是刷新了一下页面,我选择让它继续跑完,结果照样存进数据库,用户回来能在历史记录里看到。如果是聊天回答这种,用户走了就可以取消,省钱。

这个决定会影响并发控制,我项目在这里修过一个 bug。最早的写法是在流式响应的生成器里放锁,生成器结束就释放。结果客户端一断开,生成器就被关闭,锁跟着释放了,但后台线程里的构建还在跑,这时候另一个请求就能拿到锁再开一个构建,两个构建同时改环境变量,就乱了。

修复是让锁的生命周期跟着任务走,而不是跟着连接走:锁在后台线程的 finally 里释放。我写了回归测试,模拟一个慢构建,请求返回之后检查锁仍然被持有,构建结束后锁才释放。

如果要做得更完整,还应该支持重连后续看:事件带编号存起来,浏览器重连时带上最后的编号,服务端从那之后接着推。我项目没做这一步,断开后只能去历史记录里看。

3.5 “Agent 服务要上线,还需要补什么?”

我项目是单机单用户的,上线要补几块。

第一是任务执行。不能在 Web 进程的线程里跑,要用任务队列,比如 Celery 或者 RQ,Web 进程只负责提交任务、查状态、推事件;工作进程可以单独扩容,Web 进程重启也不影响正在跑的任务。任务状态和事件存 Redis 或数据库。

第二是隔离。每个任务独立的工作目录和环境变量,最好每个任务一个容器,11 篇讲的安全问题也一起解决。

第三是并发和配额。锁要换成 Redis 这种多进程共享的;每个用户限制同时运行的任务数和每天的 token 用量,超了返回 429。

第四是部署细节。前面有 Nginx 的话,SSE 要关掉响应缓冲,响应头加 X-Accel-Buffering: no,我项目加了;Nginx 默认 60 秒读不到数据就断开,编译命令可能跑两分钟没有输出,所以要每隔十几秒发一个心跳注释,这个我项目没做。

第五是可观测,10 篇讲过。

第四部分 逐个详解

4.1 一个 AI 应用后端长什么样

flowchart LR
    B[浏览器<br/>EventSource] -->|GET /api/build| N[Nginx<br/>反向代理]
    N --> API[FastAPI<br/>校验 鉴权 限流]
    API -->|提交任务| Q[(任务队列<br/>Redis)]
    Q --> W1[工作进程 1<br/>跑 Agent]
    Q --> W2[工作进程 2]
    W1 --> LLM[模型 API]
    W1 --> SB[沙箱<br/>执行工具]
    W1 -->|事件| PS[(Redis<br/>发布订阅 / 事件流)]
    PS --> API
    API -->|SSE 推送| N
    W1 --> DB[(数据库<br/>运行记录)]
    API -->|查历史| DB

这是一个比较完整的形态。我项目是它的极简版:

flowchart LR
    B[浏览器<br/>EventSource] -->|Vite 开发代理 /api| API[FastAPI 单进程]
    API -->|启动| T[后台线程<br/>build_events]
    T -->|事件| QQ[queue.Queue]
    QQ --> API
    API -->|SSE| B
    T --> DB[(runs.db)]
    API -->|/api/runs| DB
完整形态我项目差距带来的限制
任务队列 + 工作进程后台线程Web 进程重启,正在跑的构建就没了
Redis 事件流进程内 queue.Queue只有发起请求的那个连接能收到事件,断开后无法续看
分布式锁 + 每用户配额threading.Lock,全局只允许一个构建只能单用户
沙箱工具层检查见 11 篇
Nginx开发时 Vite 代理没有部署过

4.2 FastAPI 基础:逐段读 server/app.py

先说为什么会有 FastAPI。 在它之前,Python 写后端最常用的是 Django 和 Flask(Flask 是 Armin Ronacher 2010 年发布的)。它们都建立在 WSGI 上。WSGI 是 Phillip J. Eby 2003 年在 PEP 333 里定下的约定,规定 Web 服务器怎么把一个请求交给 Python 程序、程序再怎么把响应交回来。这个约定是同步的:调用一次,处理一个请求,处理完才返回。做普通网页够用,放到 AI 应用里有两处不合适:

  • 一个请求要等模型几秒到几十秒。等的这段时间线程什么也没干,却一直被占着。想同时服务几百个用户,就得开几百个线程
  • SSE、WebSocket 这种在一个连接上要来回发很多条消息的场景,“进来一个请求、出去一个响应”的 WSGI 描述不了。ASGI 文档开头讲它为什么存在,举的就是 WebSocket 和长轮询这两个例子

后来 Python 语言本身补上了两块。2015 年发布的 Python 3.5 里,PEP 492(Yury Selivanov 提出)加了 async def / await 语法,PEP 484(Guido van Rossum 等人提出)加了类型注解。Django Channels 的作者 Andrew Godwin 又牵头做出了 ASGI,把 WSGI 的思路改成异步,一个连接上可以收发多条消息。

FastAPI 就是把这几样拼到一起的框架。Sebastián Ramírez 2018 年 12 月发布了它:底层用 Starlette 这个 ASGI 工具包处理 HTTP,用 Pydantic 校验数据,接口参数直接写成类型注解,框架根据注解自动校验参数、自动生成 OpenAPI 格式的接口文档。他在官方文档的历史页里说,自己好几年都在回避写新框架,先是拿各种框架和插件去拼这些功能,一直拼不齐,最后才动手写。

底层约定等模型时参数校验、接口文档
Flask 这类 WSGI 框架WSGI,同步一个请求占一个线程自己写,或装插件
FastAPIASGI,异步async defawait 时让出线程,一个线程能同时挂着很多请求写好类型注解就自动有

FastAPI:Python 的 Web 框架。你用普通函数加类型注解写接口,它负责解析请求、校验参数、把返回值转成 JSON,还自动生成交互式接口文档。底层基于 Starlette(处理 HTTP)和 Pydantic(数据校验)。

运行方式:FastAPI 程序本身不监听端口,要由一个 ASGI 服务器来跑,最常用的是 uvicorn。ASGI(异步服务器网关接口)是 Python Web 服务器和 Web 框架之间的约定。

uvicorn server.app:app --port 8765

server.app:app 的意思是:server/app.py 这个模块里名叫 app 的对象。

创建应用和全局锁:

app = FastAPI(title="CrossBuild Agent")
running = threading.Lock()
  • threading.Lock(),同一时间只能被一个线程持有。running 是模块级变量,所有请求共享同一把锁

查询接口:

@app.get("/api/runs")
def api_runs(limit: int = 50):
    return rows(
        "SELECT id, task, model, started_at, finished_at, status, steps,"
        " prompt_tokens, completion_tokens FROM runs ORDER BY id DESC LIMIT ?",
        (limit,),
    )
  • @app.get("/api/runs"):装饰器,把函数注册为 GET /api/runs 的处理函数,这叫路由
  • limit: int = 50:函数参数不在路径里,FastAPI 就把它当作查询参数(URL 里 ?limit=10 那部分)。类型注解 int 让 FastAPI 自动把字符串 "10" 转成整数,传 ?limit=abc 会自动返回 422 错误
  • 返回一个列表(里面是字典),FastAPI 自动转成 JSON
  • SQL 里的 ?占位符,参数通过第二个参数传入,不拼接字符串,防止 SQL 注入(11 篇)
@app.get("/api/runs/{run_id}")
def api_run(run_id: int):
    run = rows("SELECT * FROM runs WHERE id = ?", (run_id,))
    if not run:
        raise HTTPException(404, "没有这条运行记录")
  • {run_id}路径参数,函数里同名参数接收,int 注解同样会自动转换和校验
  • raise HTTPException(404, ...):抛出这个异常,FastAPI 就返回对应的状态码和 {"detail": "..."}
def rows(sql, params=()):
    conn = sqlite3.connect(DB)
    conn.row_factory = sqlite3.Row
    try:
        return [dict(r) for r in conn.execute(sql, params)]
    finally:
        conn.close()
  • 每次查询新建一个连接、用完关闭。SQLite 连接默认不能跨线程用,而 def 接口在线程池里跑,每次新建最简单
  • conn.row_factory = sqlite3.Row:让查询结果的每一行可以按列名访问,dict(r) 再转成普通字典

请求体和 Pydantic:我项目的接口参数都在 URL 里,没有请求体。要接收 JSON 请求体时,定义一个 Pydantic 模型:

class BuildRequest(BaseModel):
    source: str
    target: str = "aarch64-linux-musl"

@app.post("/api/build")
def create_build(req: BuildRequest):
    ...

参数类型是 Pydantic 模型时,FastAPI 就从请求体读 JSON 并校验。04 篇讲过 Pydantic 做工具参数校验,这里是同一个东西。这段只是写法示意,没有在项目里。

静态文件

if DIST.exists():
    app.mount("/assets", StaticFiles(directory=DIST / "assets"), name="assets")

    @app.get("/")
    def index():
        return FileResponse(DIST / "index.html")

前端构建产物存在时,FastAPI 顺便把网页也发出去,这样生产环境只需要起一个进程。开发时前端用 Vite 自己的服务器,通过 vite.config.js 里的代理把 /api 开头的请求转给后端的 8765 端口。

4.3 同步和异步:什么会卡住整个服务

先理解几个词。

  • 事件循环(event loop):一个单线程的调度器。它手里有很多任务,轮流推进每一个,某个任务在等网络、等定时器时,就先去推进别的
  • 协程(coroutine)async def 定义的函数。调用它不会立刻执行,而是得到一个协程对象,交给事件循环去跑
  • await:写在协程里,意思是”这里要等一会儿,先把控制权交还给事件循环”。只能 await 可等待的对象,比如 asyncio.sleep()、异步 HTTP 客户端的请求
  • 阻塞调用:调用期间一直占着当前线程不放,比如 time.sleep()、同步的 requests.get()subprocess.run()

关键点:事件循环只有一个线程。 如果某个协程里调用了阻塞函数,这个线程就被占住,事件循环上的所有其他任务都停了。

FastAPI 对两种接口的处理方式(官方文档 Concurrency and async / await):

接口写法FastAPI 怎么跑里面能不能阻塞
async def直接在事件循环里不能
def放到线程池里可以

动手实验 1:三种写法,各发 5 个并发请求。保存为 async_block.py,用项目虚拟环境运行。

import asyncio
import time

import httpx
from fastapi import FastAPI

app = FastAPI()


@app.get("/async-blocking")
async def async_blocking():
    time.sleep(0.5)
    return {"ok": True}


@app.get("/sync")
def sync_endpoint():
    time.sleep(0.5)
    return {"ok": True}


@app.get("/async-await")
async def async_await():
    await asyncio.sleep(0.5)
    return {"ok": True}


async def measure(path, n=5):
    transport = httpx.ASGITransport(app=app)
    async with httpx.AsyncClient(transport=transport, base_url="http://test") as client:
        start = time.perf_counter()
        await asyncio.gather(*(client.get(path) for _ in range(n)))
        return time.perf_counter() - start


async def main():
    for path in ("/async-blocking", "/sync", "/async-await"):
        print(f"{path:<16} 5 个并发请求共用 {await measure(path):.2f} 秒")


asyncio.run(main())

实际输出(Python 3.13,fastapi 0.141.1,httpx 0.28.1):

/async-blocking  5 个并发请求共用 2.53 秒
/sync            5 个并发请求共用 0.53 秒
/async-await     5 个并发请求共用 0.51 秒

逐段讲。

  • 三个接口都”等 0.5 秒”,区别只在怎么等
  • httpx.ASGITransport(app=app):让 httpx 客户端直接调用 FastAPI 应用对象,不经过真实网络和端口,适合测试
  • httpx.AsyncClient:httpx 的异步客户端,发请求要 await
  • asyncio.gather(*协程们)同时启动多个协程,等全部完成。* 把生成器展开成多个参数
  • (client.get(path) for _ in range(n)):生成器表达式,产生 5 个请求协程。_ 是约定俗成的”用不到的循环变量”名

看结果。

  1. async def + time.sleep:2.53 秒,约等于 5 × 0.5,完全串行了。每个请求都把事件循环卡住 0.5 秒,下一个只能等
  2. def + time.sleep:0.53 秒。5 个请求被分到线程池的 5 个线程里,同时睡
  3. async def + await asyncio.sleep:0.51 秒。await 时让出控制权,事件循环在一个线程里同时推进 5 个请求

实际项目里最常见的坑不是 time.sleep,而是在 async def 里调用了同步库:同步的 openai.OpenAI 客户端、requestssqlite3subprocess.run。它们和实验里的 time.sleep 效果一样。

怎么选:

情况写法
里面全是异步库(openai.AsyncOpenAIhttpx.AsyncClient、异步数据库驱动)async def
里面有同步库def
async def 里不得不调一个同步函数await asyncio.to_thread(函数, 参数) 把它放到线程里
CPU 密集计算(比如本地跑 embedding 模型)放到进程池或单独的工作进程

我项目:Agent 用同步的 OpenAI 客户端和 subprocess,所以 api_build 写的是 def;Agent 本身又放在单独的后台线程里(4.5 节),接口函数只负责启动和返回流式响应。

自己改一改:

  1. n 改成 50,看 /sync 的时间。提示:线程池有大小上限,超过之后也要排队
  2. async_blocking 里把 time.sleep(0.5) 改成 await asyncio.to_thread(time.sleep, 0.5),看时间变化

4.4 SSE:流式推送的协议

先说为什么会有 SSE。 HTTP 本来是”浏览器问一句、服务器答一句”,服务器没法主动开口。可网页上有很多事要服务器有了新消息就告诉浏览器:聊天消息、股价、现在的模型逐字输出。早年只有两种绕路的办法:

  • 轮询:浏览器每隔几秒发一次请求问”有新的吗”。间隔短了,大部分请求问回来都是”没有”,白白浪费;间隔长了,消息要晚几秒才到
  • 长轮询:浏览器发一个请求,服务器没消息就先不回,等有了再回;浏览器收到后马上再发下一个。实时性好了,但每条消息都要重新发一次完整的 HTTP 请求。2006 年 Alex Russell 给这一类”服务器推送”的做法起了个总称叫 Comet

HTML5 标准的主要编辑 Ian Hickson 从 2004 年起,在 WHATWG 的 Web Applications 1.0 草案里写进了一个更直接的办法:浏览器发一个普通 GET 请求,服务器不结束响应,有消息就往里写一段,浏览器端用 EventSource 接收,断了自动重连。Opera 浏览器 2006 年最先做了实验性实现,这个名字 Server-Sent Events 也是那时用起来的;W3C 在 2015 年 2 月把它发布为正式推荐标准(W3C Server-Sent Events)。同一时期还有 WebSocket,2011 年 12 月作为 RFC 6455 发布,解决的是双向通信;RFC 的背景部分列的正是轮询的毛病:每条消息都要带一遍 HTTP 头,服务器要为一个客户端维护好几条连接。

做法服务器能主动推吗每条消息的开销断线后
轮询不能,只能等浏览器来问一次完整请求,而且大多是空问下次照常问
长轮询勉强能,有消息才回每条消息一次完整请求自己写重连
SSE能,单向连接一直开着,只多一行 data:浏览器自动重连,还带上最后一个编号
WebSocket能,双向协议升级后的小数据帧自己写重连

SSE(Server-Sent Events,服务器发送事件):浏览器发一个普通的 HTTP 请求,服务器不马上结束响应,而是保持连接,有新数据就往里写一段。响应的内容类型是 text/event-stream

格式MDN:Using server-sent events):

data: {"type": "run_started", "project": "zlib"}

event: progress
id: 7
data: 第一行
data: 第二行

: 这是注释,浏览器会忽略,可以当心跳

retry: 5000
字段作用
data消息内容;连续多行 data: 会用换行拼接成一条
event事件类型名;有它时触发同名的监听器,没有时触发 onmessage
id事件编号;浏览器记住最后一个,断线重连时放在 Last-Event-ID 请求头里发回来
retry断线后多少毫秒再重连
冒号开头的行注释,被忽略,常用来当心跳防止连接被中间设备断开

一条消息以一个空行结束,也就是 \n\n。我项目只用了最简单的 data: 一种字段,每个事件是一行 JSON。

浏览器端(我项目 web/src/components/BuildRunner.vue):

es = new EventSource(`/api/build?${q}`)

es.onmessage = (e) => {
  const ev = JSON.parse(e.data)
  if (ev.type === 'run_started') {
    meta.value = ev
  } else if (ev.type === 'finished') {
    result.value = ev
    stop()
  }
}

es.onerror = () => {
  if (running.value && !result.value) {
    error.value = '连接断了。要么后端没起来,要么已经有一个构建在跑。'
  }
  stop()
}
  • new EventSource(url):建立 SSE 连接,只能是 GET 请求
  • onmessage:每收到一条没有 event 字段的消息就调用,e.datadata: 后面的内容
  • stop() 里调用 es.close() 关闭连接

onerror 里为什么要主动关闭? EventSource 默认断线会自动重连。对”订阅一个已有任务的事件”来说这是好事,但我项目的 /api/build 是”启动一个构建”,自动重连就等于再发一次启动请求:要么拿到 409(前一个还在跑),要么在前一个结束后又启动一个新的。所以出错时直接关掉,不让它重连。

这其实暴露了一个设计问题:用 GET 请求去做”启动任务”这种有副作用的事,本来就不太对。更好的设计是 4.6 节的”POST 创建任务、GET 订阅事件”。

SSE、WebSocket、轮询怎么选:

SSEWebSocket轮询
方向服务器 → 浏览器双向浏览器主动问
协议普通 HTTPHTTP 升级成独立协议普通 HTTP
断线重连浏览器自动,带 Last-Event-ID自己写天然没有长连接
实时性取决于间隔
浏览器 API 限制只能 GET,不能加自定义请求头;HTTP/1 下每个域名最多 6 个连接较少
适合模型逐字输出、Agent 进度语音对话、协同编辑、游戏很少更新的状态

FastAPI 现在有专门的 SSE 响应。 官方文档说从 0.135.0 起提供 fastapi.sse.EventSourceResponse:接口函数里直接 yield 数据,自动编成 JSON 放进 data:;用 ServerSentEvent 可以设 eventidretry;空闲时每 15 秒自动发一个心跳注释;自动加上 Cache-Control: no-cacheX-Accel-Buffering: no 响应头。我项目写的时候用的是通用的 StreamingResponse,格式和响应头自己拼,心跳没有做

4.5 长任务、断开连接和并发锁

我项目的实现server/app.py):

def start_build(source, target, arch):
    q = queue.Queue()

    def worker():
        try:
            for ev in build_agent.build_events(source, target, arch):
                q.put(ev)
        except Exception as e:
            q.put({"type": "error", "message": f"{type(e).__name__}: {e}"})
        finally:
            running.release()
            q.put(None)

    threading.Thread(target=worker, daemon=True).start()
    return q


def stream(q):
    while True:
        ev = q.get()
        if ev is None:
            break
        yield f"data: {json.dumps(ev, ensure_ascii=False)}\n\n"


@app.get("/api/build")
def api_build(source: str, target: str = "aarch64-linux-musl", arch: str = ""):
    check_source(source)
    if not running.acquire(blocking=False):
        raise HTTPException(409, "已经有一个构建在跑,等它结束再来")
    try:
        q = start_build(source, target, arch or target.split("-")[0])
    except Exception:
        running.release()
        raise
    return StreamingResponse(
        stream(q),
        media_type="text/event-stream",
        headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"},
    )
sequenceDiagram
    participant B as 浏览器
    participant API as api_build
    participant L as running 锁
    participant T as 后台线程 worker
    participant Q as queue.Queue
    B->>API: GET /api/build
    API->>L: acquire(blocking=False)
    L-->>API: 拿到
    API->>T: 启动线程
    API-->>B: StreamingResponse(stream(q))
    loop Agent 每个事件
        T->>Q: put(事件)
        Q-->>API: stream 生成器 get()
        API-->>B: data: {...}
    end
    T->>L: finally: release()
    T->>Q: put(None)
    Q-->>API: get() 得到 None,生成器结束
    API-->>B: 响应结束

逐段讲。

  • queue.Queue()线程安全的队列。一个线程往里 put,另一个线程从里 get,不用自己加锁。get() 在队列为空时会阻塞等待,直到有东西放进来
  • def worker(): 定义在 start_build 里面,是闭包,能直接用外层的 sourceq 等变量
  • threading.Thread(target=worker, daemon=True).start():创建并启动一个线程执行 workerdaemon=True 表示守护线程,主进程退出时不等它
  • q.put(None):约定用 None 表示”没有更多事件了”,这种做法叫哨兵值(sentinel)
  • running.acquire(blocking=False):尝试拿锁,拿不到立刻返回 False,不等待。于是第二个请求直接得到 409
  • try: ... except Exception: running.release(); raise:启动线程这一步如果出错,线程里的 finally 不会执行,锁就永远不会释放了,所以这里要补一次释放
  • stream(q) 是一个生成器,StreamingResponse 会不断从它取值写给客户端。因为 get() 会阻塞,而 stream 是普通函数(不是 async),Starlette 会把它放到线程池里迭代,不会卡住事件循环
  • json.dumps(ev, ensure_ascii=False):转 JSON 时保留中文原样,不转成 \uXXXX
  • X-Accel-Buffering: no:告诉 Nginx 不要缓冲这个响应,见 4.7 节

一个修过的 bug:锁跟着连接走,而不是跟着任务走。

修复前(提交 1941ce9 之前)的写法是:

def stream(source, target, arch):
    q = queue.Queue()

    def worker():
        ...
        finally:
            q.put(None)

    threading.Thread(target=worker, daemon=True).start()
    try:
        while True:
            ev = q.get()
            if ev is None:
                break
            yield f"data: {json.dumps(ev, ensure_ascii=False)}\n\n"
    finally:
        running.release()

锁在响应生成器finally 里释放。问题出在客户端断开时:服务器停止迭代生成器并关闭它,生成器的 finally 执行,锁被释放;但后台线程里的 Agent 还在跑。这时第二个请求就能拿到锁,启动第二个构建,两个构建同时修改 os.environ 里的目标平台和 toolchain 路径。

flowchart TD
    subgraph 修复前
        A1[客户端断开] --> A2[生成器被关闭]
        A2 --> A3[finally 释放锁]
        A3 --> A4[后台线程仍在构建]
        A4 --> A5[新请求拿到锁<br/>两个构建同时跑]
    end
    subgraph 修复后
        B1[客户端断开] --> B2[生成器被关闭]
        B2 --> B3[锁仍被持有]
        B3 --> B4[后台线程构建结束<br/>worker 的 finally 释放锁]
        B4 --> B5[新请求才能拿到锁]
    end

修复:锁的释放挪到 workerfinally 里,谁执行任务,谁在任务结束时放锁tests/test_regressions.pyServerLock 测试用一个会卡住的假构建,确认请求函数返回后锁仍被持有,放行假构建之后锁才释放。

动手实验 2:按项目的写法搭一个最小版本,用真实的 uvicorn 服务器验证流式推送、409 和断开连接后的锁。保存为 sse_lock.py,用项目虚拟环境运行。

import json
import queue
import threading
import time

from fastapi import FastAPI, HTTPException
from fastapi.responses import StreamingResponse
import httpx
import uvicorn

app = FastAPI()
running = threading.Lock()


def fake_build(project):
    yield {"type": "run_started", "project": project}
    for step in (1, 2, 3):
        time.sleep(0.2)
        yield {"type": "tool_result", "step": step, "ok": step != 2}
    yield {"type": "finished", "passed": True}


def start_build(project):
    q = queue.Queue()

    def worker():
        try:
            for ev in fake_build(project):
                q.put(ev)
        finally:
            running.release()
            q.put(None)

    threading.Thread(target=worker, daemon=True).start()
    return q


def stream(q):
    while True:
        ev = q.get()
        if ev is None:
            break
        yield f"data: {json.dumps(ev, ensure_ascii=False)}\n\n"


@app.get("/api/build")
def api_build(project: str):
    if not running.acquire(blocking=False):
        raise HTTPException(409, "已经有一个构建在跑")
    q = start_build(project)
    return StreamingResponse(stream(q), media_type="text/event-stream",
                             headers={"Cache-Control": "no-cache", "X-Accel-Buffering": "no"})


server = uvicorn.Server(uvicorn.Config(app, port=8799, log_level="warning"))
threading.Thread(target=server.run, daemon=True).start()
while not server.started:
    time.sleep(0.05)

base = "http://127.0.0.1:8799/api/build"
t0 = time.perf_counter()
with httpx.stream("GET", base, params={"project": "zlib"}) as resp:
    print("状态", resp.status_code, resp.headers["content-type"])
    for line in resp.iter_lines():
        if not line:
            continue
        print(f"  [{time.perf_counter() - t0:.1f}s] {line}")
        if '"step": 1' in line:
            second = httpx.get(base, params={"project": "libpng"})
            print(f"  [{time.perf_counter() - t0:.1f}s] 第二个请求 ->", second.status_code, second.json())

time.sleep(0.05)
with httpx.stream("GET", base, params={"project": "re2"}) as resp:
    first = next(resp.iter_lines())
    print("re2 收到第一条后断开:", first)
print("断开后立刻请求 ->", httpx.get(base, params={"project": "curl"}).status_code)
time.sleep(0.8)
with httpx.stream("GET", base, params={"project": "curl"}) as resp:
    print("等 re2 在后台跑完再请求 ->", resp.status_code)
server.should_exit = True

实际输出(Python 3.13,fastapi 0.141.1,uvicorn 0.38.0,httpx 0.28.1):

状态 200 text/event-stream; charset=utf-8
  [0.1s] data: {"type": "run_started", "project": "zlib"}
  [0.3s] data: {"type": "tool_result", "step": 1, "ok": true}
  [0.3s] 第二个请求 -> 409 {'detail': '已经有一个构建在跑'}
  [0.5s] data: {"type": "tool_result", "step": 2, "ok": false}
  [0.7s] data: {"type": "tool_result", "step": 3, "ok": true}
  [0.7s] data: {"type": "finished", "passed": true}
re2 收到第一条后断开: data: {"type": "run_started", "project": "re2"}
断开后立刻请求 -> 409
等 re2 在后台跑完再请求 -> 200

逐段讲新出现的部分。

  • fake_build:假的构建,是一个生成器,每 0.2 秒产出一个事件,模拟 build_agent.build_events
  • uvicorn.Server(uvicorn.Config(app, port=8799, ...)):在代码里创建 uvicorn 服务器,而不是用命令行。server.run 放到一个守护线程里跑,主线程接着当客户端
  • while not server.started:等服务器真正开始监听再发请求
  • httpx.stream("GET", ...):以流式方式发请求,with 块里可以一边接收一边处理。resp.iter_lines() 按行迭代响应内容
  • if not line: continue:SSE 每条消息后面有一个空行,跳过它
  • next(resp.iter_lines())next 从迭代器里只取一个元素。取到第一条后离开 with 块,连接就被关闭,模拟用户关掉网页
  • server.should_exit = True:通知 uvicorn 退出

为什么不用 FastAPI 自带的 TestClient?我先试过,它的流式响应不是实时到达的,看不出”一边推一边收”的效果,第二个请求的时机也不对,所以换成了真实的服务器。

看结果。

  1. 事件每隔约 0.2 秒到达一条,是真的流式推送,不是最后一次性返回
  2. 构建进行中再请求,立刻得到 409
  3. re2 构建开始后客户端马上断开,再请求仍然是 409,说明锁没有随连接释放,后台构建还在跑。这就是修复后的行为
  4. 等 0.8 秒让 re2 的假构建在后台跑完,再请求就是 200 了

自己改一改:

  1. 按”修复前”的写法改:把 running.release()workerfinally 挪到 stream 生成器的 finally 里,重跑,看”断开后立刻请求”的结果
  2. stream 里每次 q.get(timeout=0.1) 取不到事件时(会抛出 queue.Empty),yield ": ping\n\n" 发一个心跳注释,观察输出
  3. 给事件加上 id: 字段(用一个递增的计数器)

4.6 生产环境:任务队列、事件存储、续传

我项目的”线程 + 进程内队列”有三个根本限制:Web 进程一重启,任务就没了;事件只存在内存里,只有发起请求的那个连接能收到;只能一台机器一个进程。

先说为什么会有任务队列。 这三个限制不是 AI 应用才有的。十几年前做网站的人就碰到过:用户注册后要发一封邮件、上传视频后要转码,这些事几秒到几分钟,放在请求里做,用户就一直转圈;在 Web 进程里开个线程偷偷做,进程一重启任务就没了,而且没法把活分给别的机器。办法是在中间加一个”待办清单”:Web 进程只往清单里写一条任务就立刻返回,另外有一批专门的工作进程从清单里取任务来做。清单放在 Redis、RabbitMQ 这种独立的服务里,哪个进程重启都不会丢。Python 里最常用的 Celery 就是这么来的,Ask Solem 在 2009 年 4 月发布了第一个版本(Celery 早期更新记录),当时主要给 Django 网站跑后台任务。Agent 任务一跑几分钟,正好是这类场景。

更完整的做法:

sequenceDiagram
    participant B as 浏览器
    participant API as Web 进程
    participant R as Redis
    participant W as 工作进程
    B->>API: POST /api/tasks {source, target}
    API->>API: 校验、鉴权、检查用户配额
    API->>R: 提交任务
    API-->>B: 202 {task_id}
    W->>R: 取到任务
    loop 执行 Agent
        W->>R: 追加事件(带递增 id)
    end
    B->>API: GET /api/tasks/{id}/events(SSE)
    API->>R: 读已有事件,并订阅新事件
    API-->>B: id: 1 data: ...
    Note over B: 网络断开,自动重连
    B->>API: GET .../events,请求头 Last-Event-ID: 7
    API->>R: 从 id 8 开始读
    API-->>B: id: 8 data: ...
要点做法
创建和订阅分开POST 创建任务、返回任务 ID(有副作用的操作用 POST);GET 订阅事件(可以重复、可以多人看)
任务执行Celery、RQ 等任务队列,由独立的工作进程执行,可以单独扩容
事件存储事件带递增编号存进 Redis 或数据库,而不是只放内存
断线续传浏览器重连时带 Last-Event-ID,服务端从下一条开始推
状态查询另提供 GET /api/tasks/{id} 返回当前状态,页面刷新后先查状态再订阅
取消POST /api/tasks/{id}/cancel,工作进程在每一步之间检查取消标记
并发控制按用户限制同时运行的任务数,用 Redis 计数或数据库
进程崩溃任务队列的重试 + 07 篇讲的检查点和幂等

202 Accepted:HTTP 状态码,表示”请求已接受,但处理还没完成”,适合异步任务。

4.7 部署:Nginx、心跳、Docker

先说为什么会有 Nginx。 2000 年前后最常用的 Web 服务器是 Apache,它的经典做法是每个连接交给一个进程或线程。连接一多,光是这些进程线程占的内存和来回切换就把机器拖垮了。1999 年 Dan Kegel 写了一篇文章,把”一台服务器怎么同时撑住一万个连接”叫作 C10K 问题。俄罗斯程序员 Igor Sysoev 当时在门户网站 Rambler 工作,为了解决这个问题、让网站扛住大流量,2002 年开始写 Nginx,2004 年 10 月 4 日发布了第一个公开版本 0.1.0。它的思路是少量工作进程、每个进程用事件循环同时照看成千上万个连接,哪个连接有数据就处理哪个,和 4.3 节讲的事件循环是一个道理。所以它特别适合挡在最前面:收下所有连接、处理 HTTPS、发静态文件,再把动态请求转给后面的应用进程,这就是反向代理

连接怎么处理一万个慢连接时
Apache 经典模式一个连接一个进程或线程一万个进程或线程,内存和切换开销很大
Nginx几个工作进程,每个用事件循环管很多连接大部分连接只是在等,只占一点内存

Nginx 反向代理要注意的三个配置Nginx proxy 模块文档):

配置默认值对 SSE 的影响怎么处理
proxy_bufferingonNginx 先把后端的响应读进缓冲区再转给浏览器。我在 nginx 1.29 上实测,默认配置下事件仍是实时到达的,但这个行为取决于版本和缓冲区设置,不能依赖仍然显式关掉:后端响应头加 X-Accel-Buffering: no(文档说这个头可以按响应开关缓冲),或者在对应的 location 里写 proxy_buffering off
gzipoff,但很多站点全局开启text/event-stream 开 gzip 后,实测三条间隔一秒的事件在第 3 秒一起到达,流式效果没了gzip_types 里不要包含 text/event-stream14b 篇 4.10 节
proxy_read_timeout60s两次读取之间超过 60 秒没收到数据,连接就被关闭后端定期发心跳注释;或者调大这个值

注意 proxy_read_timeout 的原文说法:超时计算的是两次连续读取之间的间隔,不是整个响应的总时长。所以只要持续有数据,SSE 连接可以保持很久;问题出在长时间没有事件的时候。

这对我项目是个实际问题。 10 篇统计过,run_command 的 P95 耗时约 15 秒,最长的一次到了 120 秒超时。一条编译命令执行期间,Agent 不会产生任何事件。如果部署在默认配置的 Nginx 后面,一条跑了 60 秒以上的命令就会让连接断开,而前端的 onerror 会直接关闭连接、显示”连接断了”。项目加了 X-Accel-Buffering: no,但没有心跳,也没有真的在 Nginx 后面部署过。

心跳:每隔十几秒发一行 : ping\n\n。冒号开头是注释,浏览器会忽略,但中间的代理看到有数据流过,就不会判定超时。FastAPI 的 EventSourceResponse 默认每 15 秒自动发。

uvicorn 多进程uvicorn --workers 4 会启动 4 个进程。这时 threading.Lock 只在各自进程内有效,4 个进程能同时各跑一个构建,而且查询 /api/build 被分到哪个进程是不确定的。所以一旦多进程,锁和事件都必须放到进程外(Redis、数据库)。

先说为什么会有 Docker。 部署最常听到的一句话是”在我电脑上是好的”:开发机上 Python 3.13、某个库 2.1 版,服务器上是 3.10、库 1.8 版,一上线就报错。以前的办法要么是写一长串安装文档让运维照着装,要么是给每个应用开一台虚拟机,前者容易漏步骤,后者每台虚拟机都要带一整个操作系统,又大又慢。Linux 内核其实早就有了隔离进程的能力(命名空间、cgroups,来历见 11 篇 4.9 节),只是用起来很麻烦。2013 年 3 月,做云平台的 dotCloud 公司的 Solomon Hykes 在 PyCon 上做了一个 5 分钟的闪电演讲,展示了他们内部用的工具 Docker,同月开源。它做的事是把”应用加它需要的所有文件”打成一个镜像,任何装了 Docker 的 Linux 机器上都能一条命令跑起来,而且跑起来的环境和打包时一样。公司后来也改名叫 Docker。

Docker:把应用和依赖打包成镜像,部署时环境一致。AI 应用里要特别注意:本地模型文件(比如 bge-m3)体积大,不要打进镜像,挂载进去或启动时下载;API key 通过环境变量或密钥管理注入,不要写进镜像;执行不可信代码的沙箱容器和 Web 服务容器要分开(11 篇)。

4.8 Redis 在 AI 应用里常见的用途

Redis 是怎么来的、它和进程内字典、MySQL、Memcached 比各差在哪,在 14b 篇 4.2 节讲。

Redis:一个把数据放在内存里的键值数据库,读写非常快,支持字符串、列表、哈希、有序集合、流等数据结构,也支持给键设置过期时间。

用途怎么用AI 应用里的例子
缓存键是请求的哈希,值是结果,设过期时间相同问题的回答、embedding 结果
限流按用户和时间窗口计数,超了返回 429每分钟请求数、每天 token 用量
分布式锁设置一个带过期时间的键,只有设置成功的才算拿到锁同一用户同时只能跑一个 Agent 任务
任务队列列表或流Celery、RQ 可以用 Redis 做消息中间件
事件流流结构,每条带递增 ID4.6 节的 SSE 断线续传
会话对话历史、短期记忆06 篇的短期记忆
发布订阅一个进程发布,多个进程订阅工作进程产生事件,Web 进程推给浏览器

缓存模型回答要小心:只有输入完全相同、而且回答不依赖时间和用户的请求才能缓存;温度大于 0 时每次回答本来就不同,缓存会让用户拿到一模一样的”随机”结果。

4.9 AI 应用特有的几个后端问题

问题为什么特别做法
请求慢一次模型调用几秒到几十秒,Agent 任务几分钟流式输出;长任务异步化
请求贵每个请求都在花钱每用户 token 配额;单任务步数和金额上限(10 篇)
下游不稳定模型 API 会限流、超时、偶尔 5xx、余额不足超时、重试退避、降级到备用模型(02 篇);402 这类不可重试的错误要快速失败并告警
输出不确定同一个请求结果不同不做结果缓存,或者只缓存温度为 0 的;前端要能处理任意格式的输出
用户会中途离开看到一半关掉网页决定任务是继续还是取消;锁跟着任务走
有副作用Agent 会执行操作创建任务用 POST、支持幂等键;重要操作人工确认(07、11 篇)

第五部分 对照项目

本篇知识点项目里的位置做到了什么没做到或可以改进的
FastAPI 路由server/app.py 三个接口查询参数、路径参数自动校验;404、409、400 错误没用 Pydantic 请求体
同步接口def api_builddef api_runsAgent 和 SQLite 都是同步的,没有卡住事件循环
后台执行start_build 线程 + queue.Queue构建不在请求线程里跑Web 进程重启任务丢失;不能多进程
SSE 推送StreamingResponse + 手拼 data:实时推送每一步没有 id、没有心跳、不支持续传
反代缓冲响应头 X-Accel-Buffering: noCache-Control: no-cache部署到 Nginx 后不会被缓冲没有实际在 Nginx 后部署验证
并发控制threading.Lock + 409同时只允许一个构建,避免环境变量和目录冲突单进程有效;全局只能一个用户
锁的生命周期锁在 workerfinally 释放(1941ce9 修复)客户端断开不会提前放锁;有 ServerLock 回归测试
启动失败放锁api_buildtry / except启动线程出错时锁不会泄漏
前端接收BuildRunner.vueEventSource按事件类型更新界面;出错主动关闭防止自动重连重复启动GET 请求带副作用;断开后不能续看
历史查询/api/runs/api/runs/{id}RunHistory.vue断开后可以在历史记录里看结果没有分页、筛选
开发代理web/vite.config.js/api 转到 8765前后端分开开发
单进程托管前端app.mount("/assets", ...)构建前端后只需一个进程
任务队列、Redis、Docker没做上线必须补

第六部分 追问清单

你刚讲完下一个追问回答方向
async def vs def你项目为什么用 defAgent 用同步 OpenAI 客户端和 subprocess,SQLite 也同步
事件循环卡住怎么发现有接口在阻塞事件循环并发压测时延迟随并发线性增长;asyncio 调试模式会报慢回调
线程池def 接口并发很高会怎样线程池有上限,超了排队;长任务不该占着请求线程
SSE为什么不用 WebSocket单向推送够用;普通 HTTP;自动重连;实现简单
EventSource要带认证 token 怎么办不能加自定义头;用 Cookie,或先 POST 拿任务 ID 再订阅
自动重连你项目为什么在 onerror 里关闭/api/build 是启动构建,自动重连等于重复启动
断开连接用户关网页任务还跑吗我项目继续跑,结果进历史;聊天类可取消;锁跟任务走
锁的 bug怎么发现和验证的分析锁的生命周期;回归测试用慢构建检查请求返回后锁仍持有
409为什么是 409 不是 429和当前状态冲突,不是频率超限
多进程--workers 4 会怎样进程内锁失效,最多同时 4 个构建;要用外部锁
NginxSSE 部署在 Nginx 后面要注意什么关缓冲;读超时 60 秒是两次读取的间隔,要心跳
心跳你项目有吗没有;编译命令可能超过 60 秒没事件,会被断开
任务队列为什么要任务队列Web 进程重启不丢任务;工作进程独立扩容;重试
续传断线后怎么续看事件带 id 存 Redis;Last-Event-ID 重连续推
Redis分布式锁要注意什么设过期时间防死锁;释放时确认是自己持有的锁
缓存模型回答能缓存吗输入完全相同、不依赖时间和用户、温度为 0 才考虑

第七部分 闭卷自测

1. FastAPI 里 async def 接口和 def 接口分别在哪里执行?什么操作放在 async def 里会卡住整个服务?

答案

async def 直接在事件循环里执行;def 放到外部线程池里执行再等待结果。在 async def 里调用阻塞操作会卡住事件循环:time.sleep、同步的 requests、同步 OpenAI 客户端、同步数据库驱动、subprocess.run、CPU 密集计算。

2. 实验 1 里三种写法 5 个并发请求分别用了多久?为什么?

答案

async def + time.sleep:2.53 秒,每个请求卡住事件循环 0.5 秒,完全串行。def + time.sleep:0.53 秒,分到线程池同时睡。async def + await asyncio.sleep:0.51 秒,await 时让出控制权,事件循环同时推进。

3. SSE 消息格式里有哪几个字段?消息之间怎么分隔?心跳怎么发?

答案

data(内容,多行拼接)、event(事件类型)、id(事件编号,重连时作为 Last-Event-ID 发回)、retry(重连间隔毫秒)。消息以空行结束(\n\n)。冒号开头的行是注释,浏览器忽略,可以当心跳,比如 : ping\n\n

4. SSE 和 WebSocket 分别适合什么场景?浏览器的 EventSource 有哪些限制?

答案

SSE 适合服务器单向推送:模型逐字输出、Agent 进度;WebSocket 适合双向频繁通信:语音对话、协同编辑、游戏。EventSource 只能 GET、不能加自定义请求头;HTTP/1 下每个浏览器对同一域名最多 6 个连接。

5. 我项目前端为什么要在 onerror 里主动关闭 EventSource?这说明接口设计上有什么问题?

答案

EventSource 默认断线自动重连,而 /api/build 是启动构建,重连等于再次启动(得到 409 或启动新构建)。问题是用 GET 做有副作用的操作;更好的设计是 POST 创建任务返回 ID,GET 订阅该任务的事件流。

6. 我项目里 queue.QueueNone 各起什么作用?acquire(blocking=False) 为什么不用默认的阻塞方式?

答案

queue.Queue 是线程安全队列,后台线程 put 事件,响应生成器 get 事件,get 在空时阻塞等待。None 是哨兵值,表示没有更多事件,生成器收到就结束。blocking=False 拿不到锁立刻返回 False,于是直接返回 409;阻塞方式会让第二个请求一直挂着等前一个构建结束。

7. 修复前锁在哪里释放?客户端断开时会发生什么?修复后为什么没问题?

答案

修复前在响应生成器的 finally 里释放。客户端断开时服务器关闭生成器,finally 执行放锁,但后台线程的构建还在跑,新请求能拿到锁再开一个构建,两个构建同时改环境变量。修复后在 worker 线程的 finally 里释放,锁跟着任务结束释放,和连接无关。实验 2 里断开后立刻请求仍是 409。

8. api_build 里为什么在 start_build 外面包一层 try / except 并释放锁?

答案

锁已经拿到,如果启动线程这一步抛异常,worker 根本没运行,它的 finally 不会执行,锁就永远不会释放,之后所有请求都得到 409。所以启动失败时要在这里补释放再抛出。

9. SSE 部署在 Nginx 后面要注意哪几个配置?我项目解决了哪个、没解决哪个?

答案

一是 proxy_buffering 默认开启,实测 nginx 1.29 默认配置下事件仍实时到达,但行为取决于版本和配置,应显式关闭;二是 gzip,对事件流开压缩会把事件攒到最后一起发,实测三条事件在第 3 秒同时到达,事件流不能压缩;三是 proxy_read_timeout 默认 60 秒,两次读取之间超过 60 秒没数据就断开。项目加了 X-Accel-Buffering: no 显式关缓冲;没有配置 gzip 相关的排除,因为还没在 Nginx 后部署过;没有心跳,编译命令超过 60 秒没事件时会被断开。

10. 用 uvicorn --workers 4 启动我项目会出什么问题?

答案

4 个进程各有自己的 threading.Lock,锁只在进程内有效,最多可以同时跑 4 个构建,互相覆盖环境变量和工作目录。要多进程就得用 Redis 或数据库这类进程外的锁和事件存储。

11. 设计一个支持断线续看的长任务接口,说出接口划分和关键机制。

答案

POST /api/tasks 创建任务返回 202 和 task_id;工作进程从任务队列取任务执行,事件带递增 id 写入 Redis;GET /api/tasks/{id}/events 以 SSE 推送,先补发已有事件再订阅新事件;浏览器重连带 Last-Event-ID,服务端从下一条开始推;另有 GET /api/tasks/{id} 查状态、POST .../cancel 取消;定期发心跳。

12. Redis 在 AI 应用里有哪些常见用途?缓存模型回答要注意什么?

答案

缓存、限流计数、分布式锁、任务队列、事件流、会话和短期记忆、发布订阅。缓存模型回答要求输入完全相同、回答不依赖时间和用户,温度大于 0 的回答本来就应该每次不同,不适合缓存。

延伸阅读

  1. FastAPI:Concurrency and async / awaitdefasync def 的区别,什么时候用哪个
  2. FastAPI:Server-Sent EventsEventSourceResponse、自动心跳、Last-Event-ID 续传
  3. MDN:Using server-sent events — 事件流格式、EventSource、连接数限制
  4. Nginx:ngx_http_proxy_moduleproxy_bufferingX-Accel-Bufferingproxy_read_timeout

Redis、分布式锁、鉴权、MySQL 慢查询和 Docker 部署的补充见 14b AI 应用后端补课

下一篇:15 系统设计题——把前面所有篇的知识组合起来,回答”设计一个 XX Agent”。

Related · Agent 开发
⎇ main ai/agent开发 22 节 230 notes UTF-8