
LLM Agent 开发做到一定规模通信协议迟早要摆到桌面上。我维护的 PI Agent 框架一直是用本地函数调用的方式跑任务输入一段文本跑完返回最终答案中间过程完全关在进程里。直到前段时间需要把 PI 接入到 Web 端和另一个 Java 服务才意识到这种单体写法有多难用。这一篇我记录把 PI 的调用入口改造成 RPC 模式、并且用 JSON 事件流做实时反馈的完整过程包括设计取舍、代码结构和踩过的几个坑希望能给同样在做 Agent 框架、或者想把自己 Agent 服务化的同学一个可直接参考的模板。1. 为什么需要 RPC 模式从本地函数到跨进程服务1.1 单体 Agent 的硬伤最开始 PI 的形态很简单一个 Python 类内部维护 prompt 模板、模型客户端、工具注册表调用方通过agent.run(查一下周五北京到上海的航班)这种方式同步拿到结果。这种写法在单机脚本、Jupyter Notebook 里非常顺手但随着接入方变多问题开始冒头。一是进程边界太硬。Web 后端想拿到 Agent 中间状态只能靠日志文件或者轮询数据库前端要展示“当前在调用哪个工具”根本没有通道。二是不好扩容。模型调用是 I/O 密集工具执行也可能要等外部 API单进程内并发一上去asyncio事件循环被某个慢工具卡住整个服务就堵死了。三是对接语言受限。PI 的核心是 Python但业务方有 Java 服务、有前端页面不能要求大家都进 Python 进程里调一个内部对象。所以我把 PI 的服务端拆成两层核心层继续跑 LLM 推理、工具循环、上下文管理这层我习惯叫它 harness只负责任务执行外层新增一个 RPC 接口层把run、stop、resume这些操作暴露成远程方法。外部调用方不需要关心 Agent 内部跑在哪台机器、用的是哪个模型只要按 RPC 协议发请求、读事件流就行。1.2 RPC 模式解决什么问题RPC也就是远程过程调用本质上是让调用方像调本地函数一样调远程服务。对 Agent 这个场景来说它解决的不只是“跨进程”更关键的是把“一次 Agent 任务”从一次函数调用语义变成一段可以被外部观察、控制、中断的交互过程。在 PI 的 RPC 设计里核心方法是agent.invoke。调用方提交session_id、input_text、config服务端创建一个运行上下文然后开始执行。执行过程中所有关键节点都会产生事件模型开始生成、生成完成、工具被选中、工具返回结果、Agent 决定下一步行动等等。这些事件不是攒到最后一块返回而是边生成边写入响应流调用方拿到的是持续的 JSON 对象流而不是一个大 JSON 包。这样做的好处很明显。第一实时性前端可以逐条渲染 Agent 的思考过程用户不会再面对一个转圈圈的黑盒。第二可观测性RPC 事件流天然就是结构化日志每一条都可以直接落库排查问题的时候回放一遍事件流比看大段文本日志高效太多。第三可控制性外部调用方可以根据事件类型决定是否继续、是否中断这是单体函数调用很难做到的。1.3 为什么不选 REST 而是 JSON-RPC有人会问直接做 REST 接口不行吗POST/api/agent/run返回一个文本看起来更常规。但 REST 面向的是资源Agent 运行过程不是资源它更像一个需要持续“执行”的方法调用。如果用 REST你要么设计一堆回调地址要么让客户端轮询状态接口轮询间隔定长了延迟高定短了又是无谓开销。我最终选了 JSON-RPC 2.0 over HTTP响应体用application/x-ndjson流式输出。理由很直接JSON-RPC 消息结构简单method、params、id三个核心字段就能表达调用意图服务端可以使用StreamingResponse把事件流一点点推出去不打破 HTTP 语义跨语言友好任何语言都能解析 JSON 和按行读取流式数据不需要像 gRPC 那样引入 IDL、代码生成团队心智负担低。gRPC 的双向流我也认真考虑过高性能、强类型确实好但 PI 的调用方里有很多是前端和脚本让他们接 protobuf 二进制反而麻烦。JSON 事件流慢是慢一点胜在直观、可调试、随处可用。对我们这种中小规模 Agent 服务来说性能瓶颈主要在模型推理和工具调用延迟上序列化开销根本排不上号。2. JSON 事件流把 Agent 思考变成看得见的时序数据2.1 事件流的基本形态所谓 JSON 事件流简单说就是服务端在响应过程中持续向 body 里写入一条条以换行符分隔的 JSON 对象。每一条 JSON 都是一个独立事件比如“模型输出了内容片段”“工具开始执行”“工具执行完成”。调用方像读日志文件一样一行一行读解析 JSON然后分发处理。这里要区分一个概念它不是 SSEServer-Sent Events。SSE 要求每行以data:前缀开头而 NDJSON 就是裸的 JSON 行。但两者可以互相转换客户端如果想走浏览器原生EventSource服务端只要在每条事件前加上data:前缀再保持格式合法即可。PI 内部直接输出裸 JSON 行一是不想被 SSE 的event:字段限制住事件类型二是方便各种语言直接按行读。事件流的最小约定是每个事件必须在一行内结束事件内部不允许出现裸换行符。这个约定看似琐碎实际是排障时最容易出问题的地方。比如模型返回文本里带了\n如果你直接拼接 JSON不把换行转义成\\n客户端按readline读到的就是残缺的半条事件解析必然失败。2.2 事件类型与状态机PI 的事件模型最开始只定义了三种start、message、end。用下来发现粒度太粗工具调用过程根本看不清后来扩展成一套更完整的事件类型。事件名用途关键 payload 字段session_started会话创建成功session_id、input_textmodel_call一次模型推理开始/结束model_name、input_tokens、output_tokensmessage_beforeAgent 文本输出的增量片段message_id、deltamessage_after一条完整消息生成完毕content、finish_reasontool_callAgent 决定调用工具tool_name、argumentstool_result工具执行完成tool_call_id、outputerror异常信息code、messagesession_finished整个会话结束finish_reason、total_tokens事件类型背后要跟一个状态机。PI 里每个 session 的状态流转是pending → running → waiting_tool → running → completed/failed。当模型决定调用工具时会话进入waiting_tool工具结果回来后再回到running。事件类型必须能反映状态迁移客户端才好画状态图、做超时判断。例如超过 30 秒没收到任何tool_result客户端就知道工具卡住了可以主动发agent.cancel。每一事件还都带两个关键 IDevent_id和trace_id。event_id是全局唯一用于日志关联和去重trace_id是整个会话共享用于把多轮模型调用、多次工具调用串联起来。后面排查事件乱序、重复投递都靠它们。2.3 兼容 SSE 和 WebSocket 的扩展方式虽然 PI 默认用 HTTP 长连接 NDJSON但实际集成时不同客户端偏好不同前端往往想要 SSE移动端可能更想用 WebSocket。事件模型设计成独立于传输层就能平滑兼容。我们的做法是在事件序列化之前先转成统一的AgentEvent对象然后由不同的 Writer 负责输出。NDJSON Writer 直接json.dumps之后加换行SSE Writer 在前面拼data:前缀再用空行分隔WebSocket Writer 则直接发一条 JSON 文本消息。底层 Agent 生成的逻辑完全不用改。这样折腾一遍之后我发现事件模型本身才是核心传输层其实就是个壳别让壳反过来绑架核心。3. 实操落地在 PI 中实现 RPC 服务与 JSON 事件流3.1 环境准备与项目结构PI 核心层是 Python 3.10 asyncioHTTP 层我选了 FastAPI主要是它处理流式响应特别顺手StreamingResponse天然支持异步生成器。客户端演示用了httpx它同样支持流式读取响应。项目结构大概长这样pi/ core/ agent.py # Agent 运行时harness 层 events.py # AgentEvent 模型定义 tools.py # 工具注册与调度 rpc/ server.py # FastAPI 应用RPC 方法注册 ndjson_writer.py # 事件序列化与写出 client/ demo.py # Python 客户端示例依赖就四个fastapi0.110.0 uvicorn[standard]0.29.0 pydantic2.6.0 httpx0.27.03.2 定义事件模型我用 Pydantic 定义事件对象既方便校验又能在出问题时直接拿到友好报错。这里的关键点是要把type做成字符串枚举而不是直接用类名因为事件流要跨语言客户端不一定知道 Python 类名是什么。# rpc/events.py import time import uuid from enum import Enum from typing import Any, Optional from pydantic import BaseModel, Field class EventType(str, Enum): session_started session_started model_call model_call message_before message_before message_after message_after tool_call tool_call tool_result tool_result error error session_finished session_finished class AgentEvent(BaseModel): event_id: str Field(default_factorylambda: str(uuid.uuid4())) trace_id: str session_id: str type: EventType payload: dict[str, Any] Field(default_factorydict) created_at: float Field(default_factorytime.time) class RpcMeta(BaseModel): jsonrpc: str 2.0 id: Optional[str] None method: str 这里有个设计细节trace_id必须在创建 session 时生成然后跟随整个会话的所有事件。event_id则每生成一个事件就创建一个新的。这样即便多个工具并发返回结果客户端也能通过trace_id归组通过event_id去重。3.3 服务端把 invoke 方法变成流式 RPC服务端的作用是把 HTTP 请求转换成内部 Agent 调用。我定义了一个装饰器rpc_method把所有 RPC 方法统一注册到一张表里然后由入口函数根据请求体里的method字段分发。核心方法是agent.invoke。需要注意这里的响应不是一个普通 JSON 对象而是一个流式生成器。生成器要先输出一条session_started事件再进入 Agent 事件循环最后输出session_finished。客户端读到session_finished就认为整个调用结束。# rpc/server.py import asyncio import json from fastapi import FastAPI, Request from fastapi.responses import StreamingResponse from pi.core.agent import run_agent app FastAPI() rpc_handlers {} def rpc_method(name): def deco(fn): rpc_handlers[name] fn return fn return deco rpc_method(agent.invoke) async def agent_invoke(params: dict, trace_id: str): session_id params[session_id] input_text params[input_text] config params.get(config, {}) async def event_gen(): # 先返回 session_started让客户端第一时间拿到会话标识 started { jsonrpc: 2.0, type: session_started, session_id: session_id, trace_id: trace_id, payload: {input_text: input_text}, } yield json.dumps(started) \n # 进入 Agent 核心循环逐条产生事件 async for event in run_agent(session_id, input_text, config): yield json.dumps(event.model_dump()) \n finished { jsonrpc: 2.0, type: session_finished, session_id: session_id, trace_id: trace_id, payload: {finish_reason: completed}, } yield json.dumps(finished) \n return StreamingResponse(event_gen(), media_typeapplication/x-ndjson) app.post(/rpc) async def rpc_entry(request: Request): body await request.json() method body.get(method) handler rpc_handlers.get(method) if not handler: return JSONResponse( {jsonrpc: 2.0, error: {code: -32601, message: method not found}} ) trace_id body.get(params, {}).get(trace_id) or str(uuid.uuid4()) return await handler(body.get(params, {}), trace_id)这段代码里最容易忽略的是media_type必须设置成application/x-ndjson。如果忘了某些客户端会因为嗅探不到类型把它当成普通 text 处理导致流式解析逻辑不生效。另外yield的每条数据末尾一定要有换行符这是客户端按行读取的基础。3.4 客户端逐行读取并分发事件客户端我用httpx的stream模式读取每一行然后按事件类型分发。关键步骤有三个构造 JSON-RPC 请求、持续读行、处理事件。# client/demo.py import json import httpx URL http://127.0.0.1:8000/rpc def dispatch(event: dict): etype event.get(type) if etype session_started: print(f[session] {event[session_id]} started) elif etype message_before: print(event[payload][delta], end) elif etype tool_call: print(f\n[tool] {event[payload][tool_name]} args{event[payload][arguments]}) elif etype tool_result: print(f\n[tool result] {event[payload][output][:200]}) elif etype error: print(f\n[error] {event[payload]}) elif etype session_finished: print(f\n[session] finished: {event[payload][finish_reason]}) def invoke(session_id: str, input_text: str): payload { jsonrpc: 2.0, id: session_id, method: agent.invoke, params: {session_id: session_id, input_text: input_text}, } with httpx.stream(POST, URL, jsonpayload, timeoutNone) as resp: for line in resp.iter_lines(): if not line: continue event json.loads(line) dispatch(event) if __name__ __main__: invoke(session-001, 帮我查一下北京明天的天气然后根据结果写一封提醒邮件)第一次跑通的时候终端里会看到文字一个一个字蹦出来工具调用记录一行一行打出来那种“看得到 Agent 在想什么”的体验比原来干等一个返回值强太多。我在实际使用中还发现timeoutNone是必须的因为一次复杂 Agent 任务可能跑好几分钟默认超时 5 秒根本不够用。3.5 密钥与鉴权信息防泄漏的落地细节做 RPC 模式的 Agent密钥管理是个绕不开的话题。事件流里会带上各种上下文包括模型请求参数、工具入参、外部 API 返回稍不注意就会把 API Key、Authorization 头、内部 token 写进事件流然后被前端直接看到。PI 里我做了两道防护。第一道是在 Agent 核心层所有工具调用参数在传给大模型或工具之前经过一个sanitize_payload函数把api_key、secret、authorization、x-api-token这类字段全部替换成***。这个函数要递归处理嵌套字典不能只查顶层。第二道是在日志层加了一个过滤器任何包含敏感字段的日志记录都不写盘也不进入事件流。# core/sanitize.py SENSITIVE_KEYS {api_key, secret, authorization, x-api-token, password, token} def sanitize_payload(data): if isinstance(data, dict): return { k: sanitize_payload(v) if k not in SENSITIVE_KEYS else *** for k, v in data.items() } if isinstance(data, list): return [sanitize_payload(item) for item in data] return data这个工作一定要前置而不是事后清洗。我在早期版本里试过在输出事件前再原样脱敏结果漏掉很多边角字段比如工具返回的 JSON 里嵌套了一个authorization字段前端页面直接给渲染出来了。把脱敏收敛到 Agent 核心层之后RPC 层和展示层都不用再担心。4. 常见问题与排查技巧实录4.1 RPC 调用超时与连接被对端关闭接入过程中最典型的报错大概有两类。一类是调用方等不到结果直接提示 30 秒超时另一类是客户端日志里出现连接被对端关闭像rpc failed; curl 56 schannel: server closed abruptly这类错误本质上都是同一个原因服务端或者中间传输层主动断了连接。出现超时先分清是客户端超时还是服务端没数据。如果客户端给的是 30 秒 RPC 超时而 Agent 任务要跑两分钟那要调的是客户端timeout参数这是最直接的解法。如果是服务端长时间没有事件输出可能是模型 API 卡住了也可能是工具在等一个永远不返回的外部请求。PI 里的经验是Agent 核心循环每 15 秒发一条heartbeat事件事件 payload 里带上当前状态和已耗时这样客户端不会因为长时间没数据而误判超时。如果是反向代理层把连接关闭多半是缓冲问题。Nginx 默认会缓冲上游响应流式数据在里面攒着不往下发。PI 在做 HTTP 流式响应时会显式设置X-Accel-Buffering: no响应头告诉中间层不要缓冲。同时反向代理的proxy_read_timeout也要调大否则一条长事件流会让代理以为自己被挂起。症状可能原因处理方式30 秒后调用失败客户端 RPC 超时调大客户端超时时间/设置 no timeout响应头已返回但 body 没数据反向代理缓冲设置X-Accel-Buffering: no事件流中间断掉读超时/无心跳服务端加heartbeat事件网络层连接重置TLS/网络不稳定客户端开启重试服务端幂等去重4.2 JSON 事件流被截断或读到一半断掉事件流被截断十有八九是换行符和 JSON 序列化的问题。我刚开始实现时事件 payload 里有模型返回的原始文本里面天然包含\n我图省事直接用了json.dumps(event)结果把多行 JSON 输出了出去。客户端按行读的时候一条事件被拆成半条JSON 解析直接抛异常。正确做法是序列化时用json.dumps(event, ensure_asciiFalse, separators(,, :))把键值之间的空格去掉、逗号合并关键是确保整个 JSON 对象在一行内。如果 payload 内部确实有换行符JSON 库会自动把它转成\\n不会真正输出换行这一点可以放心。另外Java 集成时尤其要注意。很多 Java 库解析 JSON 时默认用BufferedReader.readLine()它要求服务端每一行都是完整 JSON。如果服务端图省事把事件 JSON 格式化成了多行readLine只会读到第一行后面的JSON.parse就会报错。这时候要么让服务端统一输出单行 JSON要么 Java 客户端改用 Jackson 的readValues按流式读或者用专门的 JSON Lines 库。4.3 事件乱序与重复消费PI 里有多个工具并行执行时工具完成回调可能同时触发。如果不做控制两条tool_result事件可能交错写入同一个响应流造成前面的会话状态还没更新后面的结果就到了。解决方式是在服务端用一个asyncio.Queue做事件汇流所有并发的工具回调都把事件放进队列由唯一的 writer 协程按先进先出的顺序写出。这样事件顺序严格和入队顺序一致不会出现两个协程同时写流。重复消费更隐蔽。客户端侧如果出现网络抖动HTTP 层可能自动重试同一个请求服务端就会生成两条相同 trace 的事件流Agent 被重复执行。PI 的处理是在 RPC 入口按trace_id做幂等判断如果trace_id对应的 session 已经在执行中新请求直接返回一条error事件而不是重新跑一遍。客户端侧也要做一层去重按event_id维护一个最近的 Redis 集合或者内存集合重复事件直接丢弃。4.4 工具返回体太大把事件流撑爆Agent 的 RPC 响应流不是为超大 payload 设计的。我试过让一个工具直接把一张 10 万行的 SQL 查询结果传给大模型结果模型还没开始思考事件流先膨胀到几十兆客户端解析都开始卡。更典型的是很多低代码平台里SQL 查询内容太多导致 LLM 上下文溢出、返回不稳定本质是同一个问题。PI 里给工具结果设了一个硬上限默认单条tool_result事件不超过 512KB超过就会被截断并在 payload 里标记truncated: true。同时在把工具结果送给 LLM 前先做一个摘要只保留前几行和统计信息。这个限制要可配置不同场景差异很大检索型工具可以返回长文本结构化查询工具最好只返回聚合结果。大 payload 还有一个隐患网络层不一定能保证一次读完。NDJSON 是流式的客户端要不断消费否则 TCP 缓冲区满了服务端 writer 就会阻塞。如果客户端解析速度跟不上事件流看起来就像“卡住”了。我的建议是事件流不要塞大段完整文本工具的结果尽量落库或落到对象存储事件里只放引用地址和摘要需要全文再单独拉取。最后再分享一个小体会。这套 RPC JSON 事件流的方案在 PI 里跑了将近三周我最大的感受是调试效率提升非常明显因为每一条事件都是结构化日志随便 grep 一个trace_id就能把一次完整 Agent 执行过程回放出来。但也要注意别把事件粒度切得太细如果模型每个 token 都发一条事件光序列化和传输开销就够受的。我目前的做法是普通回答按句聚合只有工具调用和关键状态变更才单独发事件整体负担可控。下一步我打算在事件流之上加一个背压控制等实现出效果了再回来填坑。