数据截至 (上游 commit 67df07b8d807)
消费与责任链:jawn 如何把一条队列消息加工成结构化日志
30 秒导读: 边缘代理把每次 LLM 调用塞进队列后(见 02),后端服务 jawn 要把这条又生又乱的原始消息,加工成能落库、能算钱、能触发 webhook 的结构化日志。它的做法是一条 14 环的责任链(chain-of-responsibility):消息像流水线上的工件,依次经过认证、限流、读体、算成本……每个工位只补一块字段,任何一环判定"此消息不该继续"就地熔断。本章讲清这条流水线的形状和每个工位干什么;真正把数据写进数据库的
LoggingHandler细节留给 04,在线评估/webhook/PostHog 等旁路留给 05。
1. 这是什么(零基础也能懂)
一句话定义
jawn 的消费侧 = "队列 → 结构化日志"的加工车间。 它一头连着队列(Kafka 或 SQS),把边缘代理投进来的原始调用 记录成批拉出来,送进一条责任链逐环加工,最后成批写进三处存储。
它解决什么问题
边缘代理(worker)为了不阻塞用户请求,只做了最少的事:把"这次调用长什么样"打包丢进队列就返回了(见 02)。于是队列里的每条消息都是半成品——只有一个 API key 字符串、一份 heliconeMeta、一份 log 元数据,请求/响应的大 body 还躺在 S3 里没读回来,成本没算、模型名没规整、限流没判。
把这些半成品补全成"可查询、可计费、可告警"的成品,就是消费侧的活。
为什么用"责任链"而不是一个大函数
因为加工步骤多、且彼此有顺序约束:得先认证拿到组织身份,才能去 S3 按组织 ID 找 body;得先把 body 读回来解析,才能算 token 和成本;得先算完成本,计费旁路才有数可上报。
把每一步写成一个独立 handler、用 setNext 串起来,好处是:
- 每个 handler 只关心自己那一小块,易读易测(每个都有独立单测)。
- 顺序在一个地方声明清楚(
LogManager),调整流程就是挪一行。 - 任一环都能就地熔断:比如认证失败、限流命中,直接返回、不往下走。
用起来什么样(它在系统里的位置)
worker(边缘代理) jawn(后端服务,本章主角)
───────────── ───────────────────────────────
一次 LLM 调用
│ 把记录投进队列
▼
┌──────────────┐ 拉批 ┌───────────────┐ 逐条 ┌─────────────┐
│ Kafka / SQS │ ───────▶ │ 消费入口 │ ─────▶ │ 责任链 │
│ 队列 │ 小批 │ consumeMiniBatch│ 加工 │ 14 个 handler│
└──────────────┘ └───────────────┘ └──────┬──────┘
│ 链尾批量落库
▼
ClickHouse / S3 / Postgres(第 04 章)
2. 顶层全景(它大概怎么转)
三段式:拉批 → 编排 → 加工
整个消费侧可以切成三层,每层一个主角文件:
| 层 | 干什么 | 主角文件 |
|---|---|---|
| ① 队列消费 | 从 Kafka/SQS 拉一"小批"消息,反序列化,交给编排层;处理完提交 offset | lib/clients/kafkaConsumers/KafkaConsumer.ts、lib/clients/sqsConsumers/sqsConsumers.ts、lib/consumer/consumeMiniBatch.ts |
| ② 编排 | 组装那条 14 环责任链,把小批里每条消息并发喂进链头,统计耗时、失败进 DLQ,最后触发链尾批量落库 | managers/LogManager.ts |
| ③ 加工 | 责任链本体:每个 handler 补一块字段或熔断 | lib/handlers/*.ts |
怎么读下面这张主线图
从上到下是一条消息的一生。实线=正常往下流;虚线=出错时的旁路(进 DLQ)。链条部分只画了"骨架顺序",各 handler 细节见 §4。
队列(Kafka 或 SQS)
│ 拉一小批(mini-batch)
▼
consumeMiniBatch(messages, ...) consumeMiniBatch.ts:7
│ new LogManager()
▼
LogManager.processLogEntries(...) LogManager.ts:71
│ 组装责任链(setNext 链)
│ Promise.all 并发,每条消息一个 HandlerContext
▼
┌───────────────── 责任链(逐环 handle)──────────────────┐
│ Auth → RateLimit → S3Reader → RequestBody → ResponseBody │
│ → Prompt → OnlineEval → StripeIntegration → Logging │
│ → PostHog → Lytix → Webhook → Segment → StripeLog │
└──────────────────────────┬───────────────────────────────┘
成功 │ │ 任一环返回 err
▼ ▼(虚线旁路)
链尾各 handler.handleResults() 进 DLQ 死信队列
批量落库 / 上报(第 04、05 章) request-response-logs-prod-dlq
一句话把主线走一遍
一小批消息拉出来 → LogManager 组装责任链、给每条消息并发跑一遍链 → 每个 handler 往共享的 HandlerContext 上补字段 → 全跑完后,链尾几个 handler 各自把攒下的批一次性落库/上报 → 任何一条消息中途出错,就把它扔进 DLQ 等重放。
3. 核心原理(由浅入深)
3.1 责任链的骨架:setNext 串珠子,super.handle 往下传
它要解决的小问题: 怎么让"一串加工步骤"既能顺序执行,又能任一步喊停?
思路: 每个 handler 记住"我的下一环是谁",自己干完活后主动调用下一环。要熔断,就不调下一环、直接返回。
抽象基类只有两个方法,极简:
setNext(handler) 记住下一环,并把它返回(好让链式写法接着 .setNext)
handle(context) 默认实现:有下一环就调它,没有就"链结束"
真实代码就这么短(lib/handlers/AbstractLogHandler.ts:13、:18,符号 setNext / handle):
public setNext(handler: LogHandler): LogHandler {
this.nextHandler = handler;
return handler; // 返回 handler → 支持 a.setNext(b).setNext(c) 链写
}
public async handle(context: HandlerContext): PromiseGenericResult<string> {
if (!this.nextHandler) return ok("Chain complete.");
return await this.nextHandler.handle(context); // 把控制权交给下一环
}
关键点: 具体 handler 覆写 handle,先干自己的活,再用 return await super.handle(context) 把接力棒传下去。要熔断就别调 super.handle。 比如认证成功才 return await super.handle(context)(AuthenticationHandler.ts:31),失败直接 return err(...)(:15),后面的环一个都不会跑。
3.2 共享上下文:HandlerContext 是流水线上的那块"工件"
小问题: 14 个 handler 之间怎么传数据?
思路: 不用返回值层层传,而是共享一个可变对象。每条消息 new 一个 HandlerContext,从链头传到链尾,谁补的字段就挂在它身上。
HandlerContext 的关键字段(lib/handlers/HandlerContext.ts:12,类 HandlerContext):
| 字段 | 谁写的 | 装什么 |
|---|---|---|
message | 构造时传入 | 队列原始消息(KafkaMessageContents:authorization + heliconeMeta + log) |
authParams / orgParams | AuthenticationHandler | 认证后的组织身份 |
rawLog | S3ReaderHandler | 从 S3 读回的原始请求/响应 body 字符串 |
processedLog | RequestBody/ResponseBody/Prompt | 解析规整后的 body、模型名、properties |
usage / legacyUsage / costBreakdown | ResponseBodyHandler | token 用量与成本 |
timingMetrics | 每个 handler 进门时 push | 各环耗时(用来打 DataDog 指标) |
一句话: HandlerContext 就是流水线托盘,空着进来,过一环补一格,到链尾时已被填满。
3.3 两阶段:逐条 handle 攒数据,批量 handleResults 落库
小问题: 落库如果每条消息都单独写一次,数据库会被打爆。
思路: 把加工和落库拆成两阶段:
- 逐条 handle 阶段:责任链对每条消息跑一遍,像
LoggingHandler、RateLimitHandler这类只把"要写的东西"攒进自己的内部数组/批,不落库。 - 批量 handleResults 阶段:整个小批的 handle 全跑完后,
LogManager再挨个调这些 handler 的handleResults(),一次性把攒下的批刷出去。
看 RateLimitHandler 就懂:handle 里命中限流只是 this.rateLimitLogs.push(...)(RateLimitHandler.ts:58),真正插库在 handleResults 里 batchInsertRateLimits(:163,符号 handleResults)。
LogManager 在链跑完后依次触发这些"收尾"(LogManager.ts:220-229):
await this.logRateLimits(rateLimitHandler, logMetaData);
await this.logHandlerResults(loggingHandler, logMetaData, logMessages); // 主落库,第04章
await this.logStripeMeter(stripeLogHandler, logMetaData);
await this.logStripeIntegration(stripeIntegrationHandler, logMetaData);
// BEST EFFORT LOGGING —— 下面几个失败不影响主流程
this.logPosthogEvents(...); this.logLytixEvents(...);
this.logSegmentEvents(...); this.logWebhooks(...);
注意分界:await 的是必须成功的落库(限流、日志、Stripe);不 await、标注 BEST EFFORT 的是尽力而为的旁路(PostHog/Lytix/Segment/Webhook),它们的细节在 05。
3.4 熔断:不是每条消息都要走完全程
责任链的价值在于能提前退出。三种典型熔断:
| 场景 | 在哪一环 | 怎么退 | 依据 |
|---|---|---|---|
| 认证失败/查不到组织 | Authentication | 返回 err,后续全不跑 | AuthenticationHandler.ts:15、:24 |
采样丢弃(percentLog 抽样命中) | RateLimit | 攒一条 rate_limit 记录,return ok("Rate limited.") 不调 super.handle | RateLimitHandler.ts:57-65 |
| body 不在 S3(免费额度超限/omit) | S3Reader | 不报错,把 body 置空后继续 super.handle | S3ReaderHandler.ts:47-57 |
第三种是"软熔断"的反例:S3 里没 body 是正常情况(用户开了 omit,或免费额度超了不存 body),所以它不熔断、只是带着空 body 往下走,元数据照样落库。
4. 深入实现:14 环各干什么、顺序为什么这么排
4.1 链的组装:一处声明,顺序即代码
整条链在 LogManager.processLogEntries 里一次性 new 好、setNext 串好(LogManager.ts:104-118)。顺序就是下面这张表的从上到下:
| # | Handler | 职责一句话 | 熔断? | 源码符号 |
|---|---|---|---|---|
| 1 | AuthenticationHandler | 用 authorization 认证,拿到 authParams/orgParams(带 5 分钟缓存) | 失败即停 | AuthenticationHandler.ts:12 handle |
| 2 | RateLimitHandler | 免费额度概率抽查 + percentLog 采样;命中就丢弃并记一条限流日志 | 命中即停 | RateLimitHandler.ts:26 handle |
| 3 | S3ReaderHandler | 按组织 ID+请求 ID 生成签名 URL,把请求/响应大 body 从 S3 读回 rawLog | body 缺失→带空 body 继续 | S3ReaderHandler.ts:15 handle |
| 4 | RequestBodyHandler | 解析请求 body,推断请求侧模型名,清洗 properties(去 �) | 否 | RequestBodyHandler.ts:8 handle |
| 5 | ResponseBodyHandler | 按 provider 选解析器解响应 body,抽 token 用量、算成本 costBreakdown | 否 | ResponseBodyHandler.ts:68 handle |
| 6 | PromptHandler | 若 带 Helicone 模板,sanitize 后挂到 processedLog | 否 | PromptHandler.ts:8 handle |
| 7 | OnlineEvalHandler | 在线评估旁路(详见 05) | 否 | OnlineEvalHandler.ts:15 handle |
| 8 | StripeIntegrationHandler | 按 token 用量生成 Stripe 计量事件,攒批 | 否 | StripeIntegrationHandler.ts:100 handle |
| 9 | LoggingHandler | 主落库:攒请求/响应/资产等批(详见 04) | 否 | LoggingHandler.ts:159 handle |
| 10–14 | PostHog / Lytix / Webhook / Segment / StripeLog | 各类旁路上报,攒批,尽力而为 | 否 | 各 *Handler.ts |
4.2 为什么 body 必须先读、成本必须后算
这条链的顺序不是随意的,存在硬依赖:
Auth ─▶ 有了 orgParams,S3Reader 才知道去哪个组织的桶里找 body
S3Reader ─▶ 有了 rawLog(原始 body),Request/ResponseBody 才有东西可解析
ResponseBody ─▶ 解析出 token 用量 + 算出成本,StripeIntegration 才有数可计费
所以认证在最前、读体在解析前、算成本在计费前——每一步都为后一步铺路。
4.3 顺序约束的典型:stripeIntegration 为何必须排在 logging 之前
这是最容易踩的一条顺序约束。代码里专门留了注释(LogManager.ts:111-112):
.setNext(onlineEvalHandler)
// note this needs to be before the logging handler since it is mutating the properties
.setNext(stripeIntegrationHandler)
.setNext(loggingHandler)
原因: StripeIntegrationHandler 会改写 processedLog.request.properties——往里塞 helicone-stripe-integration-status、helicone-stripe-model 等标记(StripeIntegrationHandler.ts:268-273);而 LoggingHandler 落库时会把 properties 快照下来写进存储。
如果 Logging 排在 Stripe 前面,落库拿到的就是没打 Stripe 标记的旧 properties,那些计费状态标记就永远进不了库。
一句话规律: 凡是"改写 processedLog"的 handler,都必须排在 LoggingHandler(落库快照点)之前。 LoggingHandler 是这条链的"提交点",它之后的 handler(PostHog/Webhook 等)只读不改主日志。