跳到主要内容

消费与责任链: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 拉一"小批"消息,反序列化,交给编排层;处理完提交 offsetlib/clients/kafkaConsumers/KafkaConsumer.tslib/clients/sqsConsumers/sqsConsumers.tslib/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 / orgParamsAuthenticationHandler认证后的组织身份
rawLogS3ReaderHandler从 S3 读回的原始请求/响应 body 字符串
processedLogRequestBody/ResponseBody/Prompt解析规整后的 body、模型名、properties
usage / legacyUsage / costBreakdownResponseBodyHandlertoken 用量与成本
timingMetrics每个 handler 进门时 push各环耗时(用来打 DataDog 指标)

一句话: HandlerContext 就是流水线托盘,空着进来,过一环补一格,到链尾时已被填满。

3.3 两阶段:逐条 handle 攒数据,批量 handleResults 落库

小问题: 落库如果每条消息都单独写一次,数据库会被打爆。

思路: 把加工和落库拆成两阶段:

  1. 逐条 handle 阶段:责任链对每条消息跑一遍,像 LoggingHandlerRateLimitHandler 这类只把"要写的东西"攒进自己的内部数组/批,不落库
  2. 批量 handleResults 阶段:整个小批的 handle 全跑完后,LogManager 再挨个调这些 handler 的 handleResults(),一次性把攒下的批刷出去。

RateLimitHandler 就懂:handle 里命中限流只是 this.rateLimitLogs.push(...)(RateLimitHandler.ts:58),真正插库在 handleResultsbatchInsertRateLimits(: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.handleRateLimitHandler.ts:57-65
body 不在 S3(免费额度超限/omit)S3Reader不报错,把 body 置空后继续 super.handleS3ReaderHandler.ts:47-57

第三种是"软熔断"的反例:S3 里没 body 是正常情况(用户开了 omit,或免费额度超了不存 body),所以它不熔断、只是带着空 body 往下走,元数据照样落库。


4. 深入实现:14 环各干什么、顺序为什么这么排

4.1 链的组装:一处声明,顺序即代码

整条链在 LogManager.processLogEntries 里一次性 new 好、setNext 串好(LogManager.ts:104-118)。顺序就是下面这张表的从上到下:

#Handler职责一句话熔断?源码符号
1AuthenticationHandler用 authorization 认证,拿到 authParams/orgParams(带 5 分钟缓存)失败即停AuthenticationHandler.ts:12 handle
2RateLimitHandler免费额度概率抽查 + percentLog 采样;命中就丢弃并记一条限流日志命中即停RateLimitHandler.ts:26 handle
3S3ReaderHandler按组织 ID+请求 ID 生成签名 URL,把请求/响应大 body 从 S3 读回 rawLogbody 缺失→带空 body 继续S3ReaderHandler.ts:15 handle
4RequestBodyHandler解析请求 body,推断请求侧模型名,清洗 properties(去 )RequestBodyHandler.ts:8 handle
5ResponseBodyHandler按 provider 选解析器解响应 body,抽 token 用量、算成本 costBreakdownResponseBodyHandler.ts:68 handle
6PromptHandler若带 Helicone 模板,sanitize 后挂到 processedLogPromptHandler.ts:8 handle
7OnlineEvalHandler在线评估旁路(详见 05)OnlineEvalHandler.ts:15 handle
8StripeIntegrationHandler按 token 用量生成 Stripe 计量事件,攒批StripeIntegrationHandler.ts:100 handle
9LoggingHandler主落库:攒请求/响应/资产等批(详见 04)LoggingHandler.ts:159 handle
10–14PostHog / 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-statushelicone-stripe-model 等标记(StripeIntegrationHandler.ts:268-273);而 LoggingHandler 落库时会把 properties 快照下来写进存储。

如果 Logging 排在 Stripe 前面,落库拿到的就是没打 Stripe 标记的旧 properties,那些计费状态标记就永远进不了库。

一句话规律: 凡是"改写 processedLog"的 handler,都必须排在 LoggingHandler(落库快照点)之前。 LoggingHandler 是这条链的"提交点",它之后的 handler(PostHog/Webhook 等)只读不改主日志。

4.4 耗时统计:每个 handler 进门先打卡

每个 handler 的 handle 第一件事,是往 context.timingMetrics push 一条 { constructor, start }(如 RateLimitHandler.ts:27-31)。链跑完后,LogManager相邻两条打卡的时间差反推出每环耗时,累加成全局指标推给 DataDog(LogManager.ts:131-141:210-218)。这是一种"埋点在链上、统计在链外"的轻量做法。


5. 消费入口:从队列到 consumeMiniBatch

5.1 小批(mini-batch):在"一条条"和"一整批"之间取平衡

Kafka 一次 eachBatch 可能给上千条消息。全塞进一次处理太重,一条条处理又太慢。jawn 的做法是把大批再切成小批(mini-batch),小批大小可由运行时 setting 动态调(KafkaConsumer.ts:110-117,常量 MESSAGES_PER_MINI_BATCH)。

Kafka 消费循环的骨架(KafkaConsumer.ts:75 eachBatch):

eachBatch(大批):
while 还有没切完的消息:
miniBatch = 从大批切出 miniBatchSize 条 # :125
mapKafkaMessageToMessage(miniBatch) # :140 反序列化
consumeMiniBatch(filteredMessages, ...) # :175 交给编排层
finally:
resolveOffset(lastOffset) # :190 处理完才认这段 offset
heartbeat(); commitOffsetsIfNecessary()

关键设计: eachBatchAutoResolve: false(:73)——关掉自动提交 offset,改成小批处理成功后手动 resolveOffset。这样某条消息处理挂了,offset 不会被误提交,重启后能重放。

5.2 反序列化的双层解包

队列里的消息值被包了两层 JSON,mapKafkaMessageToMessage 要解两次(consumer/helpers/mapKafkaMessageToMessage.ts:5):

const kafkaValue = JSON.parse(message.value.toString()); // 外层 { value: "..." }
const parsedMsg = JSON.parse(kafkaValue.value) as KafkaMessageContents; // 内层才是真消息
messages.push(mapMessageDates(parsedMsg));

为什么两层?因为生产侧写入时就是 JSON.stringify({ value: JSON.stringify(msg) })(见 §6 生产者)。解完后 mapMessageDates(:25)把 requestCreatedAt/responseCreatedAt 从字符串还原成 Date 对象——因为 JSON 不保留 Date 类型,后续 handler 又要拿它做时间运算。

5.3 consumeMiniBatch:薄薄一层,只负责"接住并兜底"

这个函数极短(consumer/consumeMiniBatch.ts:7,符号 consumeMiniBatch),就干三件事:

const logManager = new LogManager();
try {
await logManager.processLogEntries(messages, { batchId, partition, lastOffset, messageCount });
return ok(miniBatchId);
} catch (error) {
Sentry.captureException(error, { tags: { type: "ConsumeError", topic } }); // 兜底上报
return err(`Failed to process batch ${miniBatchId}, ...`);
}

它是队列世界和编排世界的接缝:new 一个 LogManager、把小批交出去、用 try/catch 兜住任何漏网异常送 Sentry。真正的活全在 LogManager.processLogEntries

5.4 两套消费者,一个编排层

Kafka 和 SQS 是两套并行的消费者实现,但都汇聚到同一个 LogManager.processLogEntries:

消费者入口怎么拉批落到编排层
KafkaKafkaConsumer.ts:75 eachBatch订阅 topic,大批切小批consumeMiniBatchprocessLogEntries
SQSsqsConsumers.ts:112 consumeRequestResponseLogswhile(true) 轮询,一次最多拉 10 条直接 processLogEntries

SQS 侧更朴素:ReceiveMessageCommand 拉一批(单次上限 10,sqsConsumers.ts:14 MAX_NUMBER_OF_MESSAGES),处理成功后 DeleteMessageBatchCommand 删消息(:92)——SQS 用"删消息"代替 Kafka 的"提交 offset"来标记已处理。选哪套由 SQS_ENABLED 决定(index.ts:127)。


6. 生产侧:双写、DLQ 与重试

消费侧也会反过来当生产者——主要是把处理失败的消息投进 DLQ(dead-letter queue,死信队列)。这套生产者抽象值得单讲。

6.1 统一接口 + 工厂:一个 sendMessages 三种后端

所有生产者实现同一个接口 MessageProducer.sendMessages({ msgs, topic })(producers/types.ts:9)。用哪个由环境变量 QUEUE_PROVIDER 在工厂里决定(clients/HeliconeQueueProducer.ts,符号 MessageProducerFactory.createProducer):

QUEUE_PROVIDER产出的生产者
"kafka"KafkaProducer
"sqs"SQSProducer
"dual"DualWriteProducer(kafka, sqs) —— 同时写两边
其它/未设null —— 退化成同进程直接处理(sendMessageHttp 直接 new LogManager 跑一遍)

6.2 双写(DualProducer):迁移期的"主备并行"

DualWriteProducer 用于队列后端迁移期(比如从 Kafka 迁到 SQS):两边都写,但只有一边的结果算数(producers/DualProducer.ts:16,符号 sendMessages):

try { await this.primary.sendMessages(queuePayload); } // 主:写失败只 log,不 fail
catch (error) { console.error(`Error sending to primary queue: ...`); }
return this.secondary.sendMessages(queuePayload); // 备:它的结果才是返回值

巧妙处: primary 写挂了只记日志、不影响返回;真正决定成败的是 secondary。这样迁移期新旧两条队列都有数据,又不会因为老队列抽风把整条流程带崩。

6.3 各生产者的重试与批处理

两个真实后端都自带最多 3 次重试 + 1 秒退避:

生产者批处理重试键(分区依据)源码
KafkaProducerproduceMany 一把发3 次,间隔 1smsg.log.request.idKafkaProducerImpl.ts:47-84
SQSProducer每 10 条一个 SendMessageBatchCommand3 次,间隔 1s同上作为 IdSQSProducer.ts:34-72

Kafka 生产者发的值同样是双层包裹(KafkaProducerImpl.ts:53-54),与 §5.2 的双层解包正好对应:

value: JSON.stringify({ value: JSON.stringify(msg) }),
topic: topic,
key: msg.log.request.id, // 用请求 ID 做 key → 同一请求落同一分区

6.4 DLQ:失败消息的"回收站"

责任链里任一条消息返回 err,LogManager 就把它单独投进 DLQ(LogManager.ts:143-205):

责任链 handle 返回 err

├─ Sentry.captureException(type: "HandlerError") # 先上报

├─ 若 err 是 "No API key found"(认证失败)
│ → console.log 后 return,不进 DLQ(重放也没用) # :160-167

└─ 否则 pushToDLQ = SQS_ENABLED || KAFKA_ENABLED # :174
→ new HeliconeQueueProducer()
→ sendMessages([logMessage], "request-response-logs-prod-dlq")

两个设计细节:

  • 认证失败不进 DLQ(:160):No API key 是永久性错误,重放一万次还是失败,进 DLQ 纯属占地方,直接丢。
  • DLQ 有独立消费者:request-response-logs-prod-dlq 这个 topic 有自己的消费循环(KafkaConsumer.ts consumeDlqsqsConsumers.ts:146 consumeRequestResponseLogsDlq),用同一条责任链重放失败消息。
  • 落库失败也进 DLQ:不只是 handle 阶段,连链尾 LoggingHandler.handleResults 批量插库失败,也会把整个小批投进 DLQ(LogManager.ts:309-333)。

6.5 熔断保护:15 分钟超时

每条消息跑整条链时,被 withTimeout 包了 15 分钟上限(LogManager.ts:122-128,符号 withTimeout 定义在 :43):

const result = await withTimeout(
authHandler.handle(handlerContext),
60_000 * 15 // 15 分钟,防止单条消息卡死整批
);

withTimeoutPromise.race 让"链执行"和"定时器"赛跑,超时就返回 err("Timeout")(:43-57)——超时的消息同样按失败逻辑进 DLQ。


7. 巧妙之处(可借鉴的技术)

  • 顺序即契约,还写进注释。 责任链最脆的地方是"谁在谁前面",Helicone 把这条约束直接钉在代码注释里(LogManager.ts:111 "needs to be before the logging handler since it is mutating the properties"),而不是散落在文档里。改写共享状态的 handler 必须排在快照 handler 之前,这是所有责任链/中间件管道的通用坑。

  • 两阶段 handle / handleResults。 逐条 handle 只攒批、批量 handleResults 才落库,把"N 次小写"合并成"1 次大写"(RateLimitHandler.ts:58 vs :168)。这是高吞吐日志管道压数据库连接数的标准手法。

  • 软熔断 vs 硬熔断。 不是所有"缺数据"都该失败:S3 里没 body 是合法场景(omit/免费超限),于是 S3Reader 带空 body 继续(S3ReaderHandler.ts:52-57);而认证失败是硬错,立即停。区分"可降级"和"必须停",让正常但不完整的数据也能落库。

  • 认证失败不进 DLQ。 DLQ 是给"重放能救"的错误用的;永久性错误(No API key)直接丢,避免 DLQ 被无解消息塞满(LogManager.ts:160-167)。

  • 双写迁移,主错不阻断。 DualWriteProducer 让队列后端迁移做到"新旧并行、以备为准",primary 抽风只记日志(DualProducer.ts:18-22)。

  • 关掉自动提交 offset。 eachBatchAutoResolve: false + 手动 resolveOffset,把"确认消费"的时机牢牢攥在自己手里,保证挂掉能重放(KafkaConsumer.ts:73:190)。


8. 边界与局限(诚实)

  • 责任链是"读写同一个可变对象":HandlerContext 被 14 个 handler 自由改写,顺序错了就会读到旧值(§4.3)。可维护性依赖开发者记得顺序约束,类型系统不保护这一点。

  • 小批内并发,批间串行:processLogEntriesPromise.all 让一个小批里的消息并发跑链(LogManager.ts:122),但一个小批必须整批处理完才提交 offset。单条消息卡住会拖到 15 分钟超时才放行。

  • DLQ 用同一条链重放:如果失败根因是链本身的 bug(不是瞬时故障),重放也会再次失败,消息在 DLQ 里来回打转,需要人工介入。

  • percentLog 采样丢弃是有损的:RateLimitHandler 命中采样会直接丢消息(RateLimitHandler.ts:57-65),这些请求的日志永久不落库——这是刻意的成本控制,不是 bug。

  • 本章不覆盖的部分:真正的 ClickHouse/S3/Postgres 落库逻辑在 LoggingHandler 里,见 04;OnlineEval/Webhook/PostHog/会话与打分等旁路见 05


9. 代码地图(导航索引)

主题文件路径符号名
消费入口(接住小批、兜底)valhalla/jawn/src/lib/consumer/consumeMiniBatch.tsconsumeMiniBatch
双层解包 + Date 还原valhalla/jawn/src/lib/consumer/helpers/mapKafkaMessageToMessage.tsmapKafkaMessageToMessage / mapMessageDates
Kafka 消费循环(切小批、手动 offset)valhalla/jawn/src/lib/clients/kafkaConsumers/KafkaConsumer.tseachBatch / consumeDlq
SQS 消费循环(轮询、删消息)valhalla/jawn/src/lib/clients/sqsConsumers/sqsConsumers.tsconsumeRequestResponseLogs / consumeRequestResponseLogsDlq
编排:组链 + 并发跑 + DLQ + 超时valhalla/jawn/src/managers/LogManager.tsprocessLogEntries / withTimeout
责任链抽象基类valhalla/jawn/src/lib/handlers/AbstractLogHandler.tsAbstractLogHandler / setNext / handle
共享上下文(流水线托盘)valhalla/jawn/src/lib/handlers/HandlerContext.tsHandlerContext / KafkaMessageContents
① 认证环valhalla/jawn/src/lib/handlers/AuthenticationHandler.tsAuthenticationHandler
② 限流/采样环valhalla/jawn/src/lib/handlers/RateLimitHandler.tsRateLimitHandler / handleResults
③ S3 读体环(软熔断)valhalla/jawn/src/lib/handlers/S3ReaderHandler.tsS3ReaderHandler
④ 请求体解析环valhalla/jawn/src/lib/handlers/RequestBodyHandler.tsRequestBodyHandler / processRequestBody
⑤ 响应体解析 + 算成本环valhalla/jawn/src/lib/handlers/ResponseBodyHandler.tsResponseBodyHandler / getBodyProcessor
⑥ Prompt 模板环valhalla/jawn/src/lib/handlers/PromptHandler.tsPromptHandler
⑧ Stripe 计费环(顺序约束点)valhalla/jawn/src/lib/handlers/StripeIntegrationHandler.tsStripeIntegrationHandler / handleResults
生产者统一接口 + topic 类型valhalla/jawn/src/lib/producers/types.tsMessageProducer / QueueTopics
生产者工厂 + DLQ 发送valhalla/jawn/src/lib/clients/HeliconeQueueProducer.tsMessageProducerFactory / HeliconeQueueProducer
双写生产者valhalla/jawn/src/lib/producers/DualProducer.tsDualWriteProducer
Kafka 生产者(重试+双层包裹)valhalla/jawn/src/lib/producers/KafkaProducerImpl.tsKafkaProducer / KAFKA_ENABLED
SQS 生产者(10 条一批)valhalla/jawn/src/lib/producers/SQSProducer.tsSQSProducer

接着读: 链尾 LoggingHandler 怎么把这块被填满的 processedLog 真正写进 ClickHouse / S3 / Postgres → 04-storage-clickhouse-s3.md;OnlineEval / Webhook / PostHog 等旁路怎么工作 → 05-sessions-scores-evals.md;消息最初怎么被边缘代理投进队列 → 02-capture-cost-queue.md