跳到主要内容

冷热双存储与自研实时引擎

30 秒导读: Laminar 把每条 span 同时写进两套数据库——ClickHouse 存海量 span 明细 与全文搜索用的高吞吐列存,Postgres 存项目 / API key / trace 级别的关系型元数据。写库 的同一时刻,后端还会通过一套自研的 SSE(服务器推送)+ Pub/Sub 引擎,把这条 span 即时 推到正在看这条 trace 的浏览器上。本章讲清「数据落在哪」和「怎么实时到屏幕」两件事。

本章聚焦存储与推送这条主线。span 字段怎么解析、语义怎么抽取属于 02-span-llm-extraction; 去重与内容存储属于 03-dedup-content-storage;消息怎么排队进到这里,见 01-ingestion-pipeline


1. 这是什么(零基础也能懂)

  • 一句话定义: 一个「写两次 + 推一次」的存储与实时层——同一批 span,一份进 ClickHouse 当明细/搜索,一份的统计汇总进 Postgres 当元数据;写完立刻通过 SSE 推给前端。

  • 解决什么问题 / 给谁用: 假设你在网页上打开某条 trace 的实时视图,你的 agent 正在后台 一边跑一边吐 span。你希望不刷新页面就能看到新的 span 一条条冒出来。要做到这点,后端 必须:(1) 把 span 可靠地存下来(还要能被搜索、被 SQL 查询);(2) 在存的同一瞬间把它推到 你的屏幕;(3) 就算后端有好几个副本(pod),你连到的那个也得收到别的 pod 处理的 span。

  • 它能做什么:

    • 把 span/trace 明细高吞吐写入 ClickHouse(列存,适合海量 + 全文搜索)。
    • 把 trace 级别的汇总(token、成本、tag、状态、metadata)写入 Postgres(关系型,适合 信号触发器做过滤查询、项目/权限校验)。
    • 通过 SSE 长连接把 trace 更新、span 更新实时推给前端。
    • 通过 Pub/Sub(Redis 或进程内)把消息扇出到集群里所有 pod 的本地连接。
  • 用起来什么样: 前端对着这个 URL 开一条 SSE 长连接就能收流:

    GET /api/v1/projects/{project_id}/realtime?key=trace_{trace_id}

    后端每来一批 span,就往这条连接推一段文本事件(SSE 线格式):

    event: span_update
    data: {"spans":[{"spanId":"...","name":"llm.chat","spanType":"LLM", ...}]}

    event: heartbeat
    data: {}

    key 决定你订阅什么:traces 收全项目的 trace 列表更新,trace_<id> 收某条 trace 的 span 更新(见 routes/realtime.rsRealtimeQuery)。

  • 一句话直觉:ClickHouse 当仓库(能塞很多、能翻找,但不适合频繁改一行),把 Postgres 当账本(一条 trace 一行、可被关系查询和事务约束),把 SSE+Pub/Sub 当广播站 (数据一进仓库就吼一嗓子,全城的收音机都能收到)。


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

一批 span 从消费者进来后,存储与实时这一层同时干三件事。怎么读这张图:从左到右是时间 顺序;中间那步 tokio::join! 把「写 trace」和「写 span」两条支路并发跑;最右边是推屏。

┌─────────────────────────────────────────────┐
一批 Span ───────▶│ process_span_messages (traces/processor.rs) │
(来自队列) └───────────────┬──────────────────┬────────────┘
│ tokio::join! │
┌───────────────▼──────┐ ┌───────▼───────────────┐
│ trace_branch │ │ span_branch │
│ (冷·关系型) │ │ (热·明细,顺序写) │
│ │ │ │
│ ① PG traces upsert │ │ ① CH shared_content │
│ ② CH traces 汇总插入 │ │ ② CH spans 插入 │
│ ③ 推 trace 更新 ─────┼─┐ │ ③ mark_seen (Redis) │
└──────────────────────┘ │ └───────────┬───────────┘
│ │
join 完成 ─────────┴─────────────┴──▶ ④ 推 span 更新

┌────────────────────────────▼─────────────┐
│ send_to_key → PubSub.publish("sse:...") │
└────────────────────┬──────────────────────┘
Redis / 进程内 │ 扇出到每个 pod
┌────────────────────────────────────┼───────────────────────┐
▼ ▼ ▼
pod A: start_redis_subscriber pod B: 同左 pod C: 同左
│ send_to_local_connections │ │
▼ ▼ ▼
pod A 上的 SSE 长连接 pod B 上的连接 (无本地连接则丢弃)


浏览器收到 event: span_update

各部件一句话职责:

部件干什么在哪
ClickhouseService按部署模式把批量插入路由到直连 CH 或数据面app-server/src/ch/service.rs:19
CHSpan / CHTracespan/trace 的 ClickHouse 行结构与转换ch/spans.rs:74ch/traces.rs:30
TraceAggregation把一批 span 汇总成 trace 级统计ch/traces.rs:121
upsert_trace_statistics_batch把汇总 upsert 进 Postgres tracesdb/trace.rs:234
SseConnectionMap本进程内「订阅键 → 连接列表」的注册表realtime/mod.rs:24
create_sse_response建立一条 SSE 长连接、挂心跳、返回流realtime/mod.rs:52
send_to_key把一条 SSE 消息发布到 Pub/Sub 频道realtime/mod.rs:202
PubSub(Redis / InMemory)跨 pod 扇出的发布订阅抽象pubsub/mod.rs:59
start_redis_subscriber每个 pod 订阅频道、转发给本地连接realtime/mod.rs:229

3. 为什么冷热分离(先讲直觉)

问题: trace 数据有两副面孔。span 明细是海量、只追加、要全文搜索和聚合分析的;而 trace 级别的东西(它属于哪个项目、有没有触发信号、状态是不是 error)是关系型、要被 过滤查询和事务约束的。一套库很难同时把这两件事做好。

Laminar 的选择: 各挑各擅长的。

维度ClickHouse(热)Postgres(冷/关系)
存什么span 明细、去重内容、trace 明细副本项目、API key、trace 级汇总、信号/告警等元数据
数据形态海量、只追加、宽行列存一 trace 一行、可被 JOIN/事务约束
典型访问全文搜索、SQL 聚合、图表权限校验、信号触发器过滤、后台关系查询
写入特性高吞吐异步批插,可容忍瞬时重复upsert(ON CONFLICT),强一致
代表表spanstraces_replacingdeduped_contentprojectsproject_api_keystraces

一个关键点:span 明细只进 ClickHouse,不进 Postgresdb/spans.rs 里的 Span(:73) 只是内存里的领域模型,没有任何往 PG spans 表插入的代码。Postgres 侧只保留 trace 这一层 (upsert_trace_statistics_batch),因为信号触发器要对 traces.span_namestraces.status 这类字段做关系过滤(见 01-ingestion-pipeline 里的触发器评估)。

注意:同一份 trace 汇总既写 PG 又写 CH——PG 的 traces 供关系查询/触发器,CH 的 traces_replacing 供图表/SQL 分析。这不是重复浪费,而是「同一事实、两种读法」(inferred: 从 trace_branch 里先 upsert_trace_statistics_batchCHTrace::from_db_trace 插 CH 推断)。


4. ClickHouse 侧:热存储怎么写

4.1 模块清单(先看地图)

ch/mod.rs 是这一层的门面,pub mod 出一串按数据类型分的子模块:

模块负责的表/数据
ch/spans.rsspans 表(CHSpan)、按 span_id 追加 tag、调试缓存查询
ch/traces.rstraces_replacing 表(CHTrace)、TraceAggregation 汇总
ch/deduped_content.rs去重内容表(见 03)
ch/service.rsClickhouseService——按部署模式路由插入
ch/cloud.rs / ch/data_plane.rs两种 ClickhouseTrait 实现:直连 / 走数据面
ch/notifications.rsch/datapoints.rs其余各类数据

4.2 三个核心抽象:Table / ClickhouseInsertable / ClickhouseTrait

写入侧靠三层抽象解耦「什么数据」和「怎么写」。

Table 枚举 把逻辑名映射到真实表名(注意 Traces 映射到 traces_replacing):

// ch/mod.rs:66 Table::as_str — 逻辑名 → 真实 CH 表名
Table::Spans => "spans",
Table::Traces => "traces_replacing", // ReplacingMergeTree
Table::DedupedContent => "deduped_content",

ClickhouseInsertable trait(ch/mod.rs:78)是每个行类型要实现的:声明自己属于哪张 TABLE,并可选地 configure_insert 定制插入参数。热表(spanstraces_replacing)重写 了 configure_insert 去封顶服务端的异步插入等待时间(下面 4.4 讲)。

ClickhouseTrait trait(ch/mod.rs:95)是「怎么写」的多态入口——只有一个方法 insert_batch。它有两个实现:CloudClickhouse(直连 CH)和 DataPlaneClickhouse(HYBRID 部署下把数据 POST 到远端数据面)。

// ch/mod.rs:95 ClickhouseTrait — 插入的多态入口
#[async_trait]
pub trait ClickhouseTrait: Send + Sync {
async fn insert_batch<T: ClickhouseInsertable>(
&self,
items: &[T],
config: Option<&WorkspaceDeployment>,
) -> Result<()>;
}

4.3 ClickhouseService:按部署模式路由

ClickhouseService(ch/service.rs:19)是对上层暴露的门面,内部同时持有 clouddata_plane 两个实现,查一次工作区的部署模式再决定走哪条路:

// ch/service.rs:55 按 DeploymentMode 分发
match config.mode {
DeploymentMode::CLOUD => self.cloud.insert_batch(items, None).await,
DeploymentMode::HYBRID => self.data_plane.insert_batch(items, Some(&config)).await,
}
  • CLOUD:自托管和 Laminar 云的默认,直接 client.insert(...) 写 CH。
  • HYBRID:客户自己在本地跑 CH,app-server 把批数据发到对方数据面(数据不出客户网络)。

insert_batch_for_workspace(ch/service.rs:65)是同一路由的变体,给拿不到 project_id 的工作区级操作(如报表)用。

4.4 CHSpan / CHTrace:行结构与「顺序即真源」

CHSpan(ch/spans.rs:74)的字段顺序必须和 CH spans 表的列顺序一致,这样 SELECT * 才能正确反序列化(结构体注释明确点了这条)。from_db_span(ch/spans.rs:149) 把内存 Span 转成行:抽取 session/user/path、把过大的 input/output 换成 <lmnr_payload_url> 占位、把 metadata 序列化成字符串。去重相关的 hash 字段(input_message_hashes 等)在这里先留空,由消费者在别处填(见 03)。

CHTrace(ch/traces.rs:30)是 trace 的 CH 行。它的数据来自 TraceAggregation::from_spans (ch/traces.rs:150):遍历一批 span,按 trace_id 归组,累加 token/成本、取最早 start_time / 最晚 end_time、收集去重的 tag 和 span 名、认定 top span(parent_span_id 为 空的那个)。这里有个和 PG 一致的细节——trace_type 只在还是 0(默认)时才被覆盖 (ch/traces.rs:235),保证一旦被标成 EVALUATION 就不会被后到的子 span 冲掉。

为什么 traces 表是 ReplacingMergeTree: 同一条 trace 会被多批 span 反复 upsert, 每次都插一整行新版本,CH 靠 ReplacingMergeTree(num_spans) 在后台合并时留下版本号最大 (num_spans 最大)的那行。这就是为什么 Table::Traces 映射到 traces_replacing

4.5 异步插入的等待封顶(热路径调优)

热表重写 configure_insert 塞进一个设置:

// ch/spans.rs:236 给 spans 插入封顶服务端异步等待
insert.with_setting(
"async_insert_busy_timeout_max_ms",
SPANS_CH_ASYNC_INSERT_BUSY_TIMEOUT_MAX_MS.as_str(),
)

全局 CH 客户端开了 async_insert=1, wait_for_async_insert=1。CH 的自适应 busy-timeout 会 在「按时间触发 flush」时把每缓冲等待抬到默认上限(~1s)。而 Laminar 的消费者上游已经把行 攒成批了,每次 flush 的字节数远低于大小阈值,于是每次都是时间触发,自适应就一直停在上限, 客户端被 wait_for_async_insert=1 卡满 1 秒。把上限压到默认 400ms,p50 客户端等待从 ~1s 降到 ~400ms(代价:该表 parts/秒约 2.5×,CHTrace 同样处理,见 ch/traces.rs:108)。

另外 CloudClickhouse::insert_batch(ch/cloud.rs:31)给整个插入请求包了个 with_timeouts(None, INSERT_END_TIMEOUT)(ch/mod.rs:50):一个静默的 CH 端点(死 pod / 黑洞连接)会让 end() 永远挂着;超时一到,任务被 abort、返回 TimedOut,进而变成 transient 错误让 RabbitMQ 重投——一个卡死的消费者就地自愈,而不是永久静默。


5. Postgres 侧:冷/关系型元数据

Postgres 连接在 db/mod.rs(DB::connect_from_env,:33)里建池,用 search_path 指向 可配置的 schema(默认 public),所有查询用不带 schema 前缀的表名。

它存的是关系型元数据,典型三类:

  • 项目与工作区:db/projects.rs——get_projects_for_workspaceproject_has_member(权限校验)、get_project_and_workspace_billing_info(计费/限额)。
  • API key:db/project_api_keys.rs::get_api_key(:16)按哈希查项目 key,每次 ingest 请求鉴权都要读它。
  • trace 级汇总:db/trace.rs::upsert_trace_statistics_batch(:234)。

trace upsert 的 SQL(db/trace.rs:254)值得一看,它把「多批 span 增量汇总」用一条 INSERT ... ON CONFLICT (project_id, id) DO UPDATE 表达:

-- db/trace.rs:281 冲突时增量合并(节选)
input_token_count = traces.input_token_count + EXCLUDED.input_token_count,
status = CASE
WHEN traces.status = 'error' OR EXCLUDED.status = 'error' THEN 'error'
ELSE COALESCE(EXCLUDED.status, traces.status)
END,
-- `||` 合并 span_names 对象,保留全部唯一名字(供信号触发器过滤)
span_names = COALESCE(traces.span_names || EXCLUDED.span_names, ...),
  • token/成本用 + 累加;status 一旦 error 就锁死 error;span_names 用 JSONB || 合并 跨批的唯一名字——这正是信号触发器 matches_filters 依赖的累积状态。
  • type(trace 类型)用 CASE WHEN ... = 0 THEN EXCLUDED ELSE traces.type 的「首个非零胜出」 策略,和 CH 侧 TraceAggregation 的规则一致(见 CLAUDE.md「Trace Type Upsert」)。

这个 upsert 的返回值(RETURNING ...)是热路径的关键:它返回合并后的完整 Trace 行, 既用来构造 CHTrace 插 CH,又用来喂给实时推送和信号评估。下一节接上。


6. 自研实时引擎:一条连接怎么建、怎么收

6.1 连接模型:SseConnectionMap

每个 app-server 进程在内存里维护一张全局注册表:

// realtime/mod.rs:24 (project_id, 订阅键) → 连接列表
pub type SseConnectionMap =
Arc<DashMap<(Uuid, SubscriptionKey), Vec<SseConnection>>>;

SseConnection(realtime/mod.rs:14)= 一个 Uuid + 一个 mpsc::UnboundedSender<SseMessage>。 往这个 sender 发消息,就等于往对应浏览器的长连接推一段。SseMessage(:31)只有两字段: event_typedata

6.2 建立连接:create_sse_response

前端请求打到 routes/realtime.rs::sse_endpoint,它从路径取 project_id、从 query 取 订阅 key,交给 create_sse_response(realtime/mod.rs:52)。这个函数干四件事:

  1. 建一个 mpsc 无界通道,把 sender 包成 SseConnection 登记进 map(:74)。
  2. 若有初始消息就先塞进去(建连时通道刚建,发送必成功)。
  3. 起一个每 5 秒的心跳任务(:90):往 sender 发一条 heartbeat;一旦发送失败,说明 浏览器关了连接,就把这条连接从 map 里摘掉,空了就删整个 entry(:111)。心跳既是保活, 也是唯一的连接清理机制(SSE 没有 TCP FIN 的可靠信号,靠发送失败探测)。
  4. create_sse_stream(:38)把接收端包成 HTTP 流,套上 text/event-streamCache-Control: no-cacheX-Accel-Buffering: no(禁 Nginx 缓冲)等头返回。

create_sse_stream 的循环就是把每条 SseMessage 格式化成 SSE 线格式:

// realtime/mod.rs:45 SSE 线格式:event: <类型>\ndata: <json>\n\n
let sse_data = format!("event: {}\ndata: {}\n\n", message.event_type, json_data);

6.3 推给本地连接:send_to_local_connections

send_to_local_connections(realtime/mod.rs:153)是「把消息塞进本进程所有匹配连接」的函数: 按 (project_id, key) 找到连接列表,retain 时顺便清掉发送失败的死连接,列表空了就删 entry。注意它只碰本进程的 map——跨进程靠下一节的 Pub/Sub。

6.4 SSE 服务跑在哪

SSE 端点同时挂在两处:消费者进程(CONSUMER_PORT,默认 8002,main.rs:1609)和主 HTTP 服务(生产者模式,main.rs:1755)。连接 map(main.rs:851)每个进程一份;订阅者 (start_redis_subscriber,main.rs:856)也每进程一个。所以浏览器连到哪个 pod 无所谓—— 下面的扇出保证每个 pod 都收得到消息。


7. 跨实例扇出:Pub/Sub 抽象

7.1 为什么需要它

集群有多个 pod。处理某条 span 的是 pod A,但看这条 trace 的浏览器可能连在 pod B 上。 pod A 的 SseConnectionMap 里没有那条连接,直接 send_to_local_connections 推不到。 解法:pod A 不直接推,而是发布到一个频道;每个 pod 都订阅这个频道,收到后各自往 自己的本地连接转发。

7.2 PubSub trait 与两种实现

// pubsub/mod.rs:59 两种后端,enum_dispatch 静态分发
pub enum PubSub {
InMemory(InMemoryPubSub),
Redis(RedisPubSub),
}

PubSubTrait(pubsub/mod.rs:65)两个方法:publish(channel, message)subscribe(pattern, callback)(阻塞直到订阅结束)。

  • RedisPubSub(pubsub/redis.rs):publish 走可复用的 ResilientRedisConnection (:28);subscribe(:41)必须用独立的非复用连接(get_async_pubsub),因为 Redis 把订阅状态钉在单个 socket 上。流结束就退出,由调用方决定要不要重订。
  • InMemoryPubSub(pubsub/in_memory.rs):Redis 没配时的回退,进程内一张 HashMap<pattern, Vec<Sender>>,matches_pattern(:25)手写 * 通配匹配。单进程部署下 它就够了(反正只有一个 pod)。

选哪个在 main.rs:267 定:有 REDIS_URL 用 Redis,否则 InMemory。

7.3 频道命名:SseChannel

频道字符串格式是 sse:<project_id>:<subscription_key>,由 SseChannel(pubsub/mod.rs:14) 强类型封装,to_string(:44)和 from_str(:28)负责编解码。订阅用的通配 pattern 是 常量 SSE_CHANNEL_PATTERN = "sse:*:*"(pubsub/keys.rs:1)——一个订阅覆盖所有项目所有键。

7.4 发布与订阅的闭环

  • 发布端 send_to_key(realtime/mod.rs:202):把 SseMessage 序列化成 JSON, pubsub.publish("sse:<pid>:<key>", payload)
  • 订阅端 start_redis_subscriber(realtime/mod.rs:229):对 sse:*:* 订阅,回调里 把频道解析回 SseChannel、把 payload 反序列化回 SseMessage,再调 send_to_local_connections 转给本地连接。

至此闭环:任一 pod publish → 所有 pod 的 subscriber 收到 → 各自推给本地连接


8. 端到端:一条新 span 如何一边写库一边上屏

把前面的零件串起来,看 traces/processor.rs::process_span_messages 的主干。

第一步——并发写两套存储。tokio::join!(processor.rs:602)同时跑两条支路:

  • trace_branch(processor.rs:422):先 upsert_trace_statistics_batch 写 PG(拿回合并 后的 Trace 行),再把它们转成 CHTrace 插 CH,然后在这条支路里就地dispatch_trace_realtime_updates 推 trace 更新(processor.rs:507)。
  • span_branch(processor.rs:557):严格顺序 shared_content → spans → mark_seen。 为什么顺序而非并发?因为 spans 是普通 MergeTree(非 RMT),要是先插 spans 再插 shared_content 而后者失败,重投会把每条 span 存两遍。mark_seen(写 Redis 去重标记) 放最后,保证只有两套 CH 表都落库了才盖戳(processor.rs:558 的注释详述)。

第二步——推 span 更新。 join 完、span_result? 确认 span 已落库后(processor.rs:603), 才调 send_span_updates(processor.rs:622)。顺序很关键:先落库、后推屏,这样前端收到 span_update 时,数据在库里已经查得到。

为什么 trace 更新在支路里推、span 更新在 join 后推?(inferred)trace 更新只依赖 PG/CH 的 trace 行,在 trace_branch 内部推最快;而 span 更新要等 spans 真正落库(不能让用户看到 一条查不到的 span),所以必须等 join 完成。

第三步——扇出到屏幕。 send_span_updates(traces/realtime.rs:81)把 span 按 (project_id, trace_id) 归组成轻量 RealtimeSpan(from_span,:244——故意不带 input/output 重字段,只带结构信息),然后对每组 send_to_key(pubsub, project_id, "trace_<id>", msg)。 接着 §7 的闭环把它送到每个 pod 的本地连接,最终变成浏览器里的 event: span_update

多频道路由。 一条 trace 可能同时属于多个订阅频道,channels_for_trace (traces/realtime.rs:157)决定推给谁:

频道订阅键什么时候
TraceChannel::Projecttraces普通 trace,进项目 trace 列表
TraceChannel::Evaluation(id)evaluation_<id>top span 名为 evaluation 的评测 trace
TraceChannel::RolloutDebugger(sid)rollout_session_<sid>metadata 带 rollout.session_id 的调试 trace

dispatch_trace_realtime_updates(processor.rs:770)据此把 trace 分桶,分别推到对应键。 span 侧同理——除了 trace_<id>,带 rollout.session_id 属性的 span 还会额外推到 rollout_session_<sid>(traces/realtime.rs:97)。


9. 巧妙之处(可带走的技术)

  • 写库与推屏共用同一个 Trace 行。 PG upsert 的 RETURNING 一次拿回合并后的行,同时 喂给 CH 插入、实时推送、信号评估三处——不做二次查询(db/trace.rs:309)。
  • 实时消息是「瘦」的。 RealtimeSpan 刻意剔掉 input/output(traces/realtime.rs:244 注释「Excludes heavy input/output fields for performance」),推屏只带结构,详情等前端按需拉。
  • 心跳兼做连接 GC。 没有可靠的 SSE 断开信号,就用「发心跳失败 = 连接已死」来清理 map (realtime/mod.rs:103),一个机制解决保活 + 清理两件事。
  • 一层抽象吃下单机与集群。 PubSub enum 让同一份推送代码在「无 Redis 的单 pod」(进程内) 和「多 pod」(Redis)下都跑得通,业务代码完全无感(pubsub/mod.rs:59)。
  • 静默端点会自愈。 CH 插入超时 → transient → Rabbit 重投,把「永久卡死」转成「短暂重试」 (ch/mod.rs:50ch/cloud.rs:45)。

10. 边界与局限

  • 实时推送尽力而为,不保证送达。 publish 失败只记日志(realtime/mod.rs:218);浏览器 在 pod 间切换、或订阅期间断连,那段时间的更新就丢了——前端要靠刷新/重新拉取补齐(inferred)。
  • InMemoryPubSub 不能跨进程。 多 pod 部署没配 Redis 时,pod A 的消息到不了 pod B 的 连接。生产多副本必须配 Redis(pubsub/in_memory.rs 只在单进程内匹配)。
  • span 明细只在 CH。 Postgres 没有 span 行,任何要翻 span 内容的功能都得走 ClickHouse; PG 只回答 trace 级和元数据级问题。
  • 热表是普通 MergeTree,可容忍瞬时重复。 spans 靠「顺序写 + 失败重投」保证不重, 但 traces_replacing 依赖后台合并去重,查询未合并窗口时可能读到多版本(靠 RMT 语义收敛)。

11. 横向对比

同属 ai-agent-reference 的可观测项目里,「双存储」是常见取舍:明细列存(ClickHouse/ 类似 OLAP)扛吞吐与搜索,关系库(Postgres)扛元数据与权限。Laminar 的特色在于实时层是 自研的 SSE + 可插拔 Pub/Sub,而非依赖第三方消息总线——单 pod 用进程内、多 pod 用 Redis, 一套代码通吃。数据从这里流出去被查询的部分,见 05-sql-query-engine; 写进来之前的排队,见 01-ingestion-pipeline


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

主题文件路径符号名
CH 插入路由(CLOUD/HYBRID)app-server/src/ch/service.rsClickhouseServiceinsert_batch
插入多态入口app-server/src/ch/mod.rsClickhouseTraitClickhouseInsertable
逻辑表名 → 真实表名app-server/src/ch/mod.rsTableTable::as_str
插入超时封顶app-server/src/ch/mod.rsINSERT_END_TIMEOUT
直连 CH 写入app-server/src/ch/cloud.rsCloudClickhouse::insert_batch
span 行结构与转换app-server/src/ch/spans.rsCHSpanfrom_db_spanconfigure_insert
span tag 追加 / 调试缓存查询app-server/src/ch/spans.rsappend_tags_to_spanquery_debug_cache_spans_page
trace 行与批量汇总app-server/src/ch/traces.rsCHTraceTraceAggregation::from_spans
PG trace 增量 upsertapp-server/src/db/trace.rsupsert_trace_statistics_batch
PG 连接与 schemaapp-server/src/db/mod.rsDB::connect_from_env
PG API key 查询app-server/src/db/project_api_keys.rsget_api_key
SSE 连接注册表app-server/src/realtime/mod.rsSseConnectionMapSseConnection
建立 SSE 连接 + 心跳app-server/src/realtime/mod.rscreate_sse_responsecreate_sse_stream
推给本地连接app-server/src/realtime/mod.rssend_to_local_connections
发布到 Pub/Subapp-server/src/realtime/mod.rssend_to_key
每 pod 订阅转发app-server/src/realtime/mod.rsstart_redis_subscriber
SSE HTTP 端点app-server/src/routes/realtime.rssse_endpointRealtimeQuery
Pub/Sub 抽象与频道app-server/src/pubsub/mod.rsPubSubPubSubTraitSseChannel
Redis 发布订阅app-server/src/pubsub/redis.rsRedisPubSub
进程内发布订阅app-server/src/pubsub/in_memory.rsInMemoryPubSubmatches_pattern
订阅通配 patternapp-server/src/pubsub/keys.rsSSE_CHANNEL_PATTERN
span/trace 实时构造与分频道app-server/src/traces/realtime.rssend_span_updateschannels_for_traceRealtimeSpan
端到端编排(并发写 + 推屏)app-server/src/traces/processor.rsprocess_span_messagesdispatch_trace_realtime_updates