跳到主要内容

接入管线:从 gRPC 到队列到落库

30 秒导读: Laminar 是 AI agent 的可观测平台,SDK 把 OpenTelemetry span 发过来。 这一章只讲骨架——一条 span 从进程入口(gRPC / HTTP)开始,经过鉴权、限流、配额检查, 被生产者塞进一条消息队列,再由后台消费者 worker 批量拉出、处理、写库;中途失败怎么重投。 不讲 span 里的字段怎么解析、LLM 语义怎么抽取(那是 02 的事), 也不讲落库时的去重与存储细节(0304)。


1. 这章讲什么:一条 span 的旅程

Laminar 的 app-server(Rust 后端)在同一个进程里同时扮演两个角色:

  • 生产者(producer): 对外开 HTTP + gRPC 端口收 span,做完轻量校验就扔进队列、立刻回包。
  • 消费者(consumer): 后台 worker 从队列里批量捞 span,做重活(落 ClickHouse、Postgres、Quickwit)。

中间隔着一条消息队列。这样设计的好处很直白:收 span 的接口要快(SDK 在等回包),而真正 落库慢且可能失败——把两者用队列解耦,收得快、写得稳,写失败还能重投而不丢数据。

这章要让你能回答三个问题:

  1. 一条 span 经过哪些跳?(入口 → 生产者 → 队列 → 消费者 → 库)
  2. 谁在写队列、谁在消费?(生产者 publish_span_messages 写,worker SpanHandler 读)
  3. 失败怎么办?(worker 返回 requeue/reject,RabbitMQ 决定重投还是丢弃)

生产者和消费者虽然默认同进程,但可以靠一个环境变量 OPERATION_MODE 拆成两批独立部署 (enable_producer / enable_consumer,features/mod.rs:104-116):不设=两者都跑;设成 producer=只开 HTTP/gRPC;设成 consumer=只跑 worker。队列就是它俩之间唯一的接头。


2. 顶层全景(它大概怎么转)

怎么读这张图: 从左到右是一条 span 的流向,竖虚线是"队列"这道分界——左边是快路径(收下就回), 右边是慢路径(后台慢慢消化)。

SDK (OTLP)
│ gRPC :8001 ┌───────── 生产者进程 ─────────┐
▼ │ │
┌──────────────────────┐ HTTP :8000 │ ① 鉴权 authenticate_request │
│ ProcessTracesService │◄─── /v1/traces ──┤ ② 限流(按 project) │
│ ::export() │ │ ③ workspace 字节配额 │
└──────────┬───────────┘ └──────────────┬───────────────┘
│ push_spans_to_queue │
▼ │
publish_span_messages ─── 按 workspace 部署模式选队列 ─┘

│ ┌─────────────── 消息队列(二选一后端)───────────────┐
├──►│ CLOUD → observations_exchange → observations_queue │
└──►│ HYBRID → spans_data_plane_exchange → …_queue │
│ 后端: RabbitMQ(生产) 或 进程内 tokio mpsc(本地) │
└───────────────────────┬──────────────────────────────┘
│ worker 拉取(prefetch)
┌─────────── 消费者进程 ───────────▼──────────────┐
│ SpanHandler / DataPlaneSpanHandler(攒批) │
│ └─► process_span_messages(flush) │
│ ├─► Postgres traces(聚合 upsert) │
│ ├─► ClickHouse spans / 内容表 │
│ └─► Quickwit / realtime / 计费 │
│ 失败 → requeue(重投) / reject(丢弃) │
└──────────────────────────────────────────────────┘

部件一句话职责:

部件干什么在哪(文件:符号)
gRPC 入口收 OTLP/gRPC 的 span 导出请求,鉴权+限流+配额traces/grpc_service.rs:49 ProcessTracesService::export
HTTP 入口收 OTLP/HTTP 的 /v1/traces,内部同样调 push_spans_to_queueapi/v1/traces.rs(装配见 main.rs:1776)
生产者预处理、按部署模式挑队列、发布到 exchangetraces/producer.rs:180 push_spans_to_queue / :114 publish_span_messages
队列抽象统一 RabbitMQ 与进程内 mpsc 两种后端mq/mod.rs:15 MessageQueue / :88 MessageQueueTrait
Cloud 消费者攒批 → 落库traces/consumer.rs:27 SpanHandler
数据面消费者按 project 攒批 → 落库(HYBRID 部署)traces/data_plane_consumer.rs:29 DataPlaneSpanHandler
处理器一批 span 的落库主流程traces/processor.rs:117 process_span_messages
worker 循环从队列 receive、调 handler、按结果 ack/requeue/rejectbatch_worker/worker.rs:72 process_inner

主线走一遍(高层,不进代码):

  1. SDK 通过 gRPC(8001)或 HTTP(8000)把一批 span 发来。
  2. 入口鉴权、按 project 限流、查 workspace 是否超字节配额,任一不过就直接拒。
  3. 生产者把每条 span 转成队列消息,查这个 workspace 是 CLOUD 还是 HYBRID 部署,选对应 exchange 发布。
  4. 后台 worker 从队列拉消息,攒到一定数量或超时就整批调 process_span_messages 落库。
  5. 落库成功 → ack 消息;瞬时失败 → requeue(RabbitMQ 重投);永久失败(如反序列化失败)→ reject 丢弃。

3. gRPC 入口:ProcessTracesService::export

它要解决的小问题: 收下 SDK 发来的一批 span,但在花力气之前先把好三道关——你是谁、你是不是发太快、你的额度还够不够。

入口是一个 tonic gRPC service,实现 OTel 标准的 TraceService::export(traces/grpc_service.rs:49)。 按顺序做四件事:

① 鉴权。 从 gRPC metadata 里取 Authorization: Bearer <key>,查出对应的 project (grpc_service.rs:53auth::authenticate_request,auth/mod.rs:99)。取 header 时同时兼容 小写 authorization(OTel 默认)和大写 Authorization(自定义 exporter),失败返回 unauthenticated。拿到的 api_key.project_id 是后面一切的 scope。

② 按 project 限流。 只有当 grpc_rate_limiter 存在、且这个 project 被标记为"需要限流"时才计数 (grpc_service.rs:64-77)。计数用 Redis key grpc_ratelimit:<project_id>,超限返回 resource_exhausted。这里有个易踩的诚实点:

// grpc_service.rs:64-76(节选)
if let Some(ref limiter) = self.rate_limiter {
if is_project_id_rate_limited(self.cache.clone(), project_id).await {
let key = format!("grpc_ratelimit:{}", project_id);
match limiter.count(key).await {
Ok(_) => {}
Err(LimiterError::LimitExceeded(_)) =>
return Err(Status::resource_exhausted("Rate limit exceeded")),
Err(e) => log::error!("Rate limiter error, allowing request: {:?}", e),
}
}
}

注意末尾:Redis 出错时是"放行"(fail-open),而不是拦截——一次 Redis 抖动不该把 ingestion 打死。 另外(inferred)HTTP 与 gRPC 用的是两个独立的 limiter:HTTP 侧 Feature::RateLimiter (main.rs:1622,key 前缀 ratelimit:),gRPC 侧 Feature::GrpcRateLimiter(main.rs:1660,key 前缀 grpc_ratelimit:),各有各的 HTTP_LIMIT / GRPC_LIMIT 环境变量,并非共用一把配额。

③ workspace 字节配额。 若开了 Feature::UsageLimit,查这个 workspace 是否已超累计字节上限 (grpc_service.rs:79-97)。查询本身出错也是 fail-open(记日志、放行);只有明确"已超"才返回 resource_exhausted("Workspace data limit exceeded")

④ 入队。 三关都过,调 push_spans_to_queue(grpc_service.rs:99)把请求交给生产者;它内部出错统一包成 Status::internal("Failed to process traces")

HTTP 入口(/v1/traces)是并列的另一条入口,解码 protobuf/JSON 后汇入同一个 push_spans_to_queue。 本章聚焦控制流骨架,HTTP 的 Content-Type 分派细节不展开。


4. 生产者:把 span 塞进哪条队列

它要解决的小问题: 把一批已鉴权的 span 变成队列消息,并且——关键——发对队列

push_spans_to_queue(traces/producer.rs:180)先把 OTLP 请求摊平:每个 resource_span → scope_span → span,逐条 Span::from_otel_span,用 span.should_save() 过滤掉不该存的 (producer.rs:196-198),包成 RabbitMqSpanMessage。然后交给 publish_span_messages

4.1 一次预处理,然后挑队列

publish_span_messages(producer.rs:114)做两步:

  1. producer 侧预处理(producer.rs:126-135):对每条还没处理过的消息跑 preprocess_for_queue ——解析富化属性、provider 转换、算 prompt hash、并在消息进队列前就查 Redis 把已见过的 LLM 输入消息剥成 32 字节哈希(LAM-1608)。这一步是为减小队列载荷,细节属于 02/03; 骨架上你只需知道:去重在生产者侧做一次,消费者信任这个结果、不再重算

  2. 按 workspace 部署模式挑队列(producer.rs:149-174):

部署模式exchangerouting key
CLOUDOBSERVATIONS_EXCHANGEOBSERVATIONS_ROUTING_KEY
HYBRIDSPANS_DATA_PLANE_EXCHANGESPANS_DATA_PLANE_ROUTING_KEY

这两组常量都定义在 traces/mod.rs:20-26为什么分两条? CLOUD 部署的 span 直接进 Laminar 自己的 处理链;HYBRID(自托管数据面)的 span 走单独的 data-plane 队列,由 DataPlaneSpanHandler 消费、 写到客户自己的 ClickHouse(见 §7.2)。分队列 = 分处理链、分存储目标。

4.2 一道保险:载荷过大直接拒

发布前先量一下序列化后的字节数,超过 mq_max_payload()整批拒收并返回一个 partial-success 响应(producer.rs:139-147:217-227):

// producer.rs:137-147(节选)
let mq_message = serde_json::to_vec(&messages).unwrap();
if mq_message.len() >= mq_max_payload() {
log::warn!("[SPANS] MQ payload limit exceeded. Project ID: [{}], ...");
return Ok(span_count); // 返回被拒条数;上层据此报 partial_success
}

mq_max_payload()(mq/utils.rs:4)只有在用真 RabbitMQ 时才是有限值;用进程内 mpsc(本地开发)时是 usize::MAX,即不设限。


5. 消息队列抽象:一套接口,两种后端

它要解决的小问题: 生产环境要用 RabbitMQ(持久、跨节点、能重投),但本地开发不想装 RabbitMQ。 于是队列被抽象成一个 trait,底下挂两种实现。

抽象在 mq/mod.rs。核心是一个 enum_dispatch 的枚举 + trait:

// mq/mod.rs:14-18
#[enum_dispatch]
pub enum MessageQueue {
Rabbit(RabbitMQ),
TokioMpsc(TokioMpscQueue),
}

MessageQueueTrait(mq/mod.rs:88)只暴露三件事:publish(发)、get_receiver(拿一个消费者)、 is_healthy(就绪探针用)。两种后端的取舍:

维度RabbitMQ(mq/rabbit.rs)进程内 mpsc(mq/tokio_mpsc.rs)
用于生产本地开发(未配 RABBITMQ_URL 时回退)
持久化delivery_mode=2 持久消息(rabbit.rs:160)无,进程内内存 channel
重投靠 broker redeliver(nack/reject requeue)ack/nack 都是 no-op(mod.rs:49/66)
消息 TTL支持 per-message TTL(rabbit.rs:161-164)忽略(tokio_mpsc.rs:85)
连接韧性通道池 + 退避重试 + setup 超时(rabbit.rs:20)不涉及
is_healthy检查连接状态(rabbit.rs:306)true(tokio_mpsc.rs:147)

RabbitMQ 侧的健壮细节(骨架相关):

  • 发布走通道池(rabbit.rs:152 publish):从 deadpool 取一个 channel,用完还池,避免每条消息新建通道; 发布失败按指数退避重试(最长 60s,rabbit.rs:217-221)。
  • 消费 setup 有整链超时(rabbit.rs:257-301):create_channel → basic_qos → queue_bind → basic_consume 整个包在 tokio::time::timeout 里。因为 lapin 在半死连接上可能卡在 basic_consume 里不返回,没有这个超时,worker 外层的退避重连永远等不到下一次尝试。

进程内 mpsc 侧用一个 DashMap<key, Vec<Sender>> 维护每条队列的发送端,publish挑当前容量最大 (最闲)的那个 receiver 投递(tokio_mpsc.rs:105-117)——一个粗糙的负载均衡。它只在本地开发用, 所以 ack/持久化/TTL 全部省略。


6. main.rs 里的装配:声明队列、拼 handler、双服务器并存

它要解决的小问题: 进程启动时,谁来把 exchange/queue 在 broker 上声明出来?worker 怎么和 handler 绑上?gRPC 和 HTTP 两个服务器怎么同时跑?

6.1 声明 exchange 与 queue

启动时(仅当用 RabbitMQ),main.rs 用一个临时 channel 把所有 exchange/queue 声明一遍。span 相关的两条:

// main.rs:318-349(节选)
let mut quorum_queue_args = FieldTable::default();
quorum_queue_args.insert("x-queue-type".into(),
lapin::types::AMQPValue::LongString("quorum".into())); // 仲裁队列:跨节点存活

channel.exchange_declare(OBSERVATIONS_EXCHANGE.into(),
ExchangeKind::Fanout, /* durable */ ..).await.unwrap(); // main.rs:327
channel.queue_declare(OBSERVATIONS_QUEUE.into(),
/* durable */ .., quorum_queue_args.clone()).await.unwrap(); // main.rs:340

两个要点:exchange 是 Fanout 类型(main.rs:329),意味着 routing key 其实被忽略,消息广播给所有绑定队列; 队列声明为 quorum 队列(x-queue-type=quorum,main.rs:318-322),所以即便某个 RabbitMQ 节点挂了, 消息数据仍在。data-plane 那条(SPANS_DATA_PLANE_EXCHANGE/QUEUE,main.rs:353-375)是同样的模板。 用进程内 mpsc 时则改为 queue.register_queue(...)(main.rs:782-784)登记 key。

6.2 拼 worker 与 handler

消费者线程里,用 BatchWorkerPool::spawn 把"worker 类型 + 数量 + 一个 handler 工厂 + 队列配置"绑一起 (main.rs:1135-1156):

// main.rs:1135-1156(节选)
batch_worker_pool_clone.spawn(
BatchWorkerType::Spans,
num_spans_workers as usize,
move || SpanHandler { db, cache, queue, clickhouse, ch: ch_cloud, pubsub,
pii_redactor, config: BatchingConfig { size, flush_interval } },
QueueConfig::new(OBSERVATIONS_QUEUE, OBSERVATIONS_EXCHANGE, OBSERVATIONS_ROUTING_KEY),
);

data-plane 的一份在 main.rs:1178-1199,绑 DataPlaneSpanHandler 到 data-plane 队列。worker 数量、批大小、 flush 间隔都来自 env(env::workers::* / env::batching::*)。

6.3 gRPC 与 HTTP 服务器并存

两个服务器跑在各自的 OS 线程上,各自 block_on 一个 tokio 运行时(main.rs:1696 HTTP 线程 / main.rs:1891 gRPC 线程)——因为 actix-web 的 HttpServer 和 tonic 的 Server 各自要 own 一个运行时。

服务器绑定端口起点
HTTP(actix-web)PORT,默认 8000main.rs:1721 HttpServer::new:1883 .bind(port)
gRPC(tonic)GRPC_PORT,默认 8001main.rs:1912 Server::builder:1925 serve_with_shutdown
Realtime SSECONSUMER_PORT,默认 8002(见 04)

端口常量在 env/server.rs:6-11。gRPC service 组装在 main.rs:1897:ProcessTracesService::new(db, cache, clickhouse, queue, grpc_rate_limiter)——注意 grpc_rate_limiter这里被 move 进 gRPC 线程, 而 HTTP 侧的 rate_limiter 早已 move 进 HTTP 闭包(main.rs:1762),两者互不干扰。gRPC service 还开了 gzip 压缩收发和 max_decoding_message_size(main.rs:1915-1917)。


7. 消费者:攒批,然后整批落库

它要解决的小问题: 一条条 span 单独写 ClickHouse 太碎(part 太多)。所以 worker 不是收一条写一条, 而是攒一批再一次性落库。

7.1 Cloud 侧:SpanHandler

SpanHandler(traces/consumer.rs:27)实现 BatchMessageHandler。它的 state 就是一个 Vec<MessageDelivery<...>>——攒着还没落库的投递。逻辑很简单:

handle_message(每次来一个投递):
空批 → 直接 ack
否则 push 进 state,累计 span 数
累计数 ≥ config.size ──► flush_batch (consumer.rs:51-77)
handle_interval(定时器每 flush_interval 触发):
state 非空 ──► flush_batch (consumer.rs:79-85)

即**"攒够数量"或"超时"二者先到者**触发 flush。flush_batch(consumer.rs:91)把所有攒着的投递摊平成 一个 Vec<RabbitMqSpanMessage>,调 process_span_messages(consumer.rs:121),再按返回结果把整批 投递 ack / requeue / reject(consumer.rs:100-104)。

7.2 数据面侧:DataPlaneSpanHandler

HYBRID 部署用 DataPlaneSpanHandler(traces/data_plane_consumer.rs:29)。关键差别:它的 state 是 按 project 分桶的 HashMap<Uuid, Vec<...>>(data_plane_consumer.rs:43),因为不同 project 的 span 要写到各自的数据面 ClickHouse。攒批阈值是每个 project 独立算的(data_plane_consumer.rs:63-85), flush 前还会 get_workspace_deployment 拿到该 project 的部署配置,把它作为 Some(&config) 传给 process_span_messages(data_plane_consumer.rs:149-163)——Cloud 侧则传 None

7.3 落库主流程 process_span_messages(骨架)

process_span_messages(traces/processor.rs:117)是一批 span 的落库总控。本章只勾勒它的控制流跳数, 不展开去重/计费/PII(那些在 03 讲):

  1. 富化属性:未预处理的消息补跑 parse_and_enrich_attributes(processor.rs:132-140)。
  2. 摘出 metadata-only 虚拟 span:POST /v1/traces/metadata 打进来的补丁消息在这里被拆走,走单独的 trace 元数据合并路径,不进常规落库(processor.rs:146-176)。
  3. 补 usage、构建 trace 聚合、构建去重批次(processor.rs:178-294)。
  4. 两条分支并行(processor.rs:602 tokio::join!):
    • trace 分支(processor.rs:422):把 span 聚合成 trace,upsert 进 Postgres,再写 ClickHouse traces、 推 realtime。
    • span 分支(processor.rs:557):严格顺序 shared_content → spans → mark_seen——因为 spans 是普通 MergeTree(非幂等),先写 spans 后写内容表若失败,重投会把每行 span 写两遍;mark_seen (盖 Redis 去重标记)必须在两张表都落定后才做。
  5. 信号、Quickwit 索引、autocomplete、字节计费(processor.rs:606-765)收尾。

骨架上记住一点:span 分支返回 Result,任一 ClickHouse 插入失败就返回 HandlerError::transient (processor.rs:577-593),这个"瞬时错误"会一路传回 worker,触发下面的重投。


8. 失败与重投:requeue / reject 怎么决定

它要解决的小问题: 落库失败了,这批消息是该重试,还是该丢掉别再浪费?

答案取决于错误类型。HandlerError(worker/mod.rs:35)分两种:

错误类型语义队列动作
Transient瞬时(如 ClickHouse 暂时插入失败)requeue → RabbitMQ 重投
Permanent永久(如消息反序列化失败,重试也没用)reject 无重投 → 丢弃

handler 把结果表达成 HandlerResult { to_ack, to_reject, to_requeue }(batch_worker/message_handler.rs:23)。 worker 循环 process_inner(batch_worker/worker.rs:72)负责翻译成实际的 broker 操作:

receiver.receive() 拿到一条投递 (worker.rs:98)
反序列化失败 → acker.reject(requeue=false) 丢弃 (worker.rs:112, 193)
否则调 handler.handle_message → HandlerResult
handle_result: (worker.rs:201-230)
to_ack → acker.ack()
to_reject → reject(requeue=false) 永久失败,丢
to_requeue → reject(requeue=true) 瞬时失败,重投

几个骨架要点:

  • 反序列化失败当场 reject 且不重投(worker.rs:193-195)——坏消息重投一万次还是坏的,直接扔。
  • 连接断了会重连,未 ack 的消息由 broker 自动重投:process_inner 每次进来先清空 state 和 acker (worker.rs:74-75),注释点明"unacked messages will be redelivered"。这正是 quorum 队列 + 手动 ack 的意义——worker 崩了或没 ack,消息不丢。
  • 进程内 mpsc 后端的 ack/reject 都是 no-op(mq/mod.rs:49/66/76),所以本地开发没有真正的重投; 重投语义只在 RabbitMQ 后端成立。

9. 采样:sampling.rs(诚实边界)

traces/sampling.rs 提供按用户的采样能力:核心是 should_sample_trace(sampling.rs:144)——按 p = min(1.0, sample_rate/100 * base_factor) 掷一次随机数决定这条 trace 收不收,其中 base_factor 来自对每个用户历史 trace 量的均衡(compute_sampling_factors,sampling.rs:80),用 ClickHouse traces_replacing 的 3 天窗口统计、缓存 1 小时(get_sampling_factors_cached,sampling.rs:26)。 用意是给高频用户降采样、低频用户不动,避免大户淹没数据。

诚实说明: 这个模块整文件带 #![cfg_attr(not(feature = "signals"), allow(dead_code))] (sampling.rs:1)——它只在 signals 特性下被消费,而 signals 是企业版功能(OSS 里 app-server/src/signals 只是公开桩)。所以在开源默认构建里,采样并不挂在本章描述的主 ingestion 控制流上;ProcessTracesService::exportprocess_span_messages 都不调用它。它属于接入相关的 边角能力,放在这里让你知道"有这么个东西、在哪",但别把它当成每条 span 必经的一跳。


10. 边界与坑(骨架层面)

  • 两把限流器,别当一把。 HTTP 与 gRPC 用独立的 limiter、独立的 feature flag 和环境变量,Redis key 前缀也不同(ratelimit: vs grpc_ratelimit:)。想真正卡死一个 project 的两条入口,两边都得配。
  • 限流/配额都是 fail-open。 Redis 或配额查询出错时放行(grpc_service.rs:73/91)。安全性让位于 可用性——宁可漏限一次,不可因依赖抖动黑洞掉 ingestion。
  • 载荷超限=整批丢,只回 partial-success。 超过 mq_max_payload 的批被直接拒(producer.rs:139), SDK 需要自己拆小重发;服务端不会替你分批。
  • 本地 mpsc 无持久化、无重投。 只用于开发。任何"消息丢没丢/重投几次"的行为验证,必须在 RabbitMQ 后端做。
  • span 分支的落库顺序是硬约束。 shared_content → spans → mark_seen 不能重排,否则重投会造成 span 重复行或去重标记悬空(processor.rs:557-599)。改这段前先读 03

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

主题文件符号
gRPC 入口 / 鉴权限流配额app-server/src/traces/grpc_service.rsProcessTracesService::export
Bearer 鉴权app-server/src/auth/mod.rsauthenticate_request
生产者:摊平+发布app-server/src/traces/producer.rspush_spans_to_queue / publish_span_messages
生产者预处理app-server/src/traces/producer.rspreprocess_for_queue
队列常量(两条链)app-server/src/traces/mod.rsOBSERVATIONS_* / SPANS_DATA_PLANE_*
队列抽象app-server/src/mq/mod.rsMessageQueue / MessageQueueTrait
RabbitMQ 后端app-server/src/mq/rabbit.rsRabbitMQ::publish / get_receiver
进程内 mpsc 后端app-server/src/mq/tokio_mpsc.rsTokioMpscQueue
载荷上限app-server/src/mq/utils.rsmq_max_payload
Cloud 消费者app-server/src/traces/consumer.rsSpanHandler
数据面消费者app-server/src/traces/data_plane_consumer.rsDataPlaneSpanHandler
落库主流程app-server/src/traces/processor.rsprocess_span_messages
worker 循环 / ack 逻辑app-server/src/batch_worker/worker.rsprocess_inner / handle_result
错误类型 → 重投判定app-server/src/worker/mod.rsHandlerError
exchange/queue 声明、双服务器app-server/src/main.rsmain(见 :318/:1135/:1721/:1912)
producer/consumer 角色开关app-server/src/features/mod.rsenable_producer / enable_consumer
端口默认值app-server/src/env/server.rsPORT / GRPC_PORT / CONSUMER_PORT
按用户采样(signals 专用)app-server/src/traces/sampling.rsshould_sample_trace