LLM Agent 服务化改造:基于 JSON-RPC 与事件流的实时反馈实践
2026/9/24 20:00:30 网站建设 项目流程

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 接口层,把runstopresume这些操作暴露成远程方法。外部调用方不需要关心 Agent 内部跑在哪台机器、用的是哪个模型,只要按 RPC 协议发请求、读事件流就行。

1.2 RPC 模式解决什么问题

RPC,也就是远程过程调用,本质上是让调用方像调本地函数一样调远程服务。对 Agent 这个场景来说,它解决的不只是“跨进程”,更关键的是把“一次 Agent 任务”从一次函数调用语义,变成一段可以被外部观察、控制、中断的交互过程。

在 PI 的 RPC 设计里,核心方法是agent.invoke。调用方提交session_idinput_textconfig,服务端创建一个运行上下文,然后开始执行。执行过程中,所有关键节点都会产生事件:模型开始生成、生成完成、工具被选中、工具返回结果、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 消息结构简单,methodparamsid三个核心字段就能表达调用意图;
  • 服务端可以使用StreamingResponse把事件流一点点推出去,不打破 HTTP 语义;
  • 跨语言友好,任何语言都能解析 JSON 和按行读取流式数据;
  • 不需要像 gRPC 那样引入 IDL、代码生成,团队心智负担低。

gRPC 的双向流我也认真考虑过,高性能、强类型确实好,但 PI 的调用方里有很多是前端和脚本,让他们接 protobuf 二进制反而麻烦。JSON 事件流慢是慢一点,胜在直观、可调试、随处可用。对我们这种中小规模 Agent 服务来说,性能瓶颈主要在模型推理和工具调用延迟上,序列化开销根本排不上号。

2. JSON 事件流:把 Agent 思考变成看得见的时序数据

2.1 事件流的基本形态

所谓 JSON 事件流,简单说就是服务端在响应过程中,持续向 body 里写入一条条以换行符分隔的 JSON 对象。每一条 JSON 都是一个独立事件,比如“模型输出了内容片段”“工具开始执行”“工具执行完成”。调用方像读日志文件一样,一行一行读,解析 JSON,然后分发处理。

这里要区分一个概念:它不是 SSE(Server-Sent Events)。SSE 要求每行以data:前缀开头,而 NDJSON 就是裸的 JSON 行。但两者可以互相转换,客户端如果想走浏览器原生EventSource,服务端只要在每条事件前加上data:前缀,再保持格式合法即可。PI 内部直接输出裸 JSON 行,一是不想被 SSE 的event:字段限制住事件类型,二是方便各种语言直接按行读。

事件流的最小约定是:每个事件必须在一行内结束,事件内部不允许出现裸换行符。这个约定看似琐碎,实际是排障时最容易出问题的地方。比如模型返回文本里带了\n,如果你直接拼接 JSON,不把换行转义成\\n,客户端按readline读到的就是残缺的半条事件,解析必然失败。

2.2 事件类型与状态机

PI 的事件模型最开始只定义了三种:startmessageend。用下来发现粒度太粗,工具调用过程根本看不清,后来扩展成一套更完整的事件类型。

事件名用途关键 payload 字段
session_started会话创建成功session_idinput_text
model_call一次模型推理开始/结束model_nameinput_tokensoutput_tokens
message_beforeAgent 文本输出的增量片段message_iddelta
message_after一条完整消息生成完毕contentfinish_reason
tool_callAgent 决定调用工具tool_namearguments
tool_result工具执行完成tool_call_idoutput
error异常信息codemessage
session_finished整个会话结束finish_reasontotal_tokens

事件类型背后要跟一个状态机。PI 里每个 session 的状态流转是:pending → running → waiting_tool → running → completed/failed。当模型决定调用工具时,会话进入waiting_tool,工具结果回来后再回到running。事件类型必须能反映状态迁移,客户端才好画状态图、做超时判断。例如超过 30 秒没收到任何tool_result,客户端就知道工具卡住了,可以主动发agent.cancel

每一事件还都带两个关键 ID:event_idtrace_idevent_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 +asyncio,HTTP 层我选了 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 客户端示例

依赖就四个:

fastapi==0.110.0 uvicorn[standard]==0.29.0 pydantic==2.6.0 httpx==0.27.0

3.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_factory=lambda: str(uuid.uuid4())) trace_id: str session_id: str type: EventType payload: dict[str, Any] = Field(default_factory=dict) created_at: float = Field(default_factory=time.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_type="application/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 客户端:逐行读取并分发事件

客户端我用httpxstream模式,读取每一行,然后按事件类型分发。关键步骤有三个:构造 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, json=payload, timeout=None) 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 在想什么”的体验,比原来干等一个返回值强太多。我在实际使用中还发现,timeout=None是必须的,因为一次复杂 Agent 任务可能跑好几分钟,默认超时 5 秒根本不够用。

3.5 密钥与鉴权信息防泄漏的落地细节

做 RPC 模式的 Agent,密钥管理是个绕不开的话题。事件流里会带上各种上下文,包括模型请求参数、工具入参、外部 API 返回,稍不注意就会把 API Key、Authorization 头、内部 token 写进事件流,然后被前端直接看到。

PI 里我做了两道防护。第一道是在 Agent 核心层,所有工具调用参数在传给大模型或工具之前,经过一个sanitize_payload函数,把api_keysecretauthorizationx-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_ascii=False, 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 都发一条事件,光序列化和传输开销就够受的。我目前的做法是普通回答按句聚合,只有工具调用和关键状态变更才单独发事件,整体负担可控。下一步我打算在事件流之上加一个背压控制,等实现出效果了再回来填坑。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询