跳到主要内容

数据截至 (上游 commit 0261ea4f33d4)

上报层:从内存对象到后端的异步管道

30 秒导读: 上一章讲了 @track 怎么把一次函数调用变成 span 树(见 01-tracing-decorator.md)。这一章讲这些 span 怎么离开你的进程——从 client.span() 这一行调用,一直追到 httpx 把 JSON 发出去。核心答案:调用线程什么都不干,只造一个 dataclass 塞进队列;后台线程负责攒批、发送、限流退避、断网落盘重放。


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

一句话定义: 上报层是 Opik Python SDK 里的一条单向异步管道,把"内存里的 trace/span 对象"变成"后端收到的 HTTP 请求"。

它要解决的问题,用场景讲:

你在一个 LLM 应用里加了观测。一次用户请求会产生 1 个 trace + 20 个 span,每个 span 都带 prompt 和 response 的完整文本。如果每个 span 都同步发一次 HTTP:

  • 用户请求要多等 21 个 RTT——观测把业务拖慢了。
  • 后端要扛 21 倍的请求数——观测把后端打挂了。
  • 网断了怎么办?业务代码要不要因为"日志发不出去"而报错?——显然不应该。

上报层就是为这三件事存在的。它给出的三个答案分别是:异步(业务线程立刻返回)、攒批(21 个请求合成 1 个)、韧性(限流退避 + 断网落盘重放)。

用起来什么样。 从使用者角度看,它几乎是隐形的——你只会在两个地方感知到它:

import opik

client = opik.Opik()
trace = client.trace(name="chat", input={"q": "hi"}) # 立刻返回,什么都还没发出去
span = trace.span(name="llm-call", type="llm") # 同上
span.end(output={"a": "hello"})

client.flush(timeout=10) # 感知点 1:显式等一次,确保东西到了后端
client.end() # 感知点 2:进程退出前收尾(atexit 也会自动调)

一句话直觉: 把它当成应用日志的 syslog 客户端——logger.info() 从不阻塞你,因为真正的写盘/发网在别的线程;代价是进程被 kill -9 时尾巴上那点日志会丢。Opik 上报层是同一套权衡,只是把"写盘"换成了"发 REST",并且额外补了一层 SQLite 兜底。


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

怎么读这张图: 竖线左边是你的业务线程(调用 trace()/span() 的那个),右边是 SDK 自己拉起来的后台线程。消息一律从上往下流,只有"附件抽取"和"断线重放"两条回边会把消息重新塞回队列。

业务线程(你的代码) │ 后台线程(daemon)
───────────────────────────── │ ────────────────────────────────────
client.trace() / span() │
update_span() / 打分 ... │
│ 造一个 message │
▼ │
Streamer.put() │
├─① 附件预处理 │
└─② 批量预处理 ────────────┼──► 各类 Batcher(攒着,不入队)
│ 没被吃掉的才继续 │ │ 攒满 1000 条 / 到点(1~3s)
▼ │ ▼ 合成一条批消息
╔═══════════════════════════════════════════════╗
║ MessageQueue:有界 deque,满了丢最老的 ║
╚═══════════════════════════════════════════════╝
│ │ N 个 QueueConsumer 取
│ ▼
│ ChainedMessageProcessor
│ ├─ 附件抽取器 ──(回塞队列)──┐
│ ├─ Online 处理器 → REST → 后端
│ └─ 本地 emulator(评估时才开)
│ │ 发失败(断网)
│ ▼
│ ReplayManager + SQLite
│ └─ 连上了 ──(回塞队列)──────┘

部件一句话职责:

部件干什么在哪个文件
Opik 客户端唯一职责是造 message + put,从不发 HTTPsdks/python/src/opik/api_objects/opik_client.py
messages.*管道里流的数据类型,纯 dataclasssdks/python/src/opik/message_processing/messages.py
Streamer管道总控:预处理顺序、入队、flush/close 语义sdks/python/src/opik/message_processing/streamer.py
BatchManager + 各 Batcher按消息类型攒批,攒满或到点合成批消息sdks/python/src/opik/message_processing/batching/
MessageQueue线程安全有界队列 + 未完成任务计数sdks/python/src/opik/message_processing/message_queue.py
QueueConsumer后台线程:出队、判投递时间、处理限流sdks/python/src/opik/message_processing/queue_consumer.py
OpikMessageProcessor真正调 REST client 的那一层 + 错误分诊sdks/python/src/opik/message_processing/processors/online_message_processor.py
ReplayManager + DBManager断网时把消息落 SQLite,连上后重放sdks/python/src/opik/message_processing/replay/

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

  1. 业务线程调 client.span(...),SDK 造一个 CreateSpanMessage,调 Streamer.put()返回。整个过程无 IO。
  2. put() 里跑两级预处理。开了批量的话,CreateSpanMessage 会被 Batcher 吃掉(累积在内存列表里),此时队列里什么都没有。
  3. 攒够 1000 条、或 2 秒定时器到点,Batcher 把累积的 span 合成一条 CreateSpansBatchMessage 丢进队列。
  4. QueueConsumer 线程取出批消息,交给处理器链;OpikMessageProcessorrest_client.spans.create_spans(...) 发 HTTP。
  5. 后端返 429 → 抛 OpikCloudRequestsRateLimited → consumer 按 retry_after 推迟并把消息放回队尾保序。网络不通 → 消息在 SQLite 里被标记 failed,连接恢复后重放。

3. 管道里流的是什么:消息类型体系

这节讲什么: 整条管道只认一种东西——BaseMessage 的子类。理解了消息分类,后面所有分支逻辑都好懂。

3.1 BaseMessage 上那四个不起眼的字段

所有消息继承 BaseMessagesdks/python/src/opik/message_processing/messages.py:48-54),它自己不带任何业务字段,只带四个投递用的元数据

字段默认值谁写它干什么
delivery_time0.0QueueConsumer 遇到 429 时单调时钟意义上的"最早可投递时刻",0.0 表示随时可发
delivery_attempts1每次 _push_message_back 自增本地 emulator 用它识别"这是重试"
message_idNoneReplayManager.register_messageSQLite 主键,用来事后标记 delivered/failed
message_type类名字符串类属性,写死反序列化时的类型标签、鉴权黑名单的 key

这四个字段都是 field(init=False)——构造消息时不用传,也不会出现在发给后端的 payload 里:as_payload_dict() 会把它们连同附件标记一起 pop 掉(messages.py:56-69)。而 as_db_message_dict() 相反,原样保留全部字段,因为 SQLite 重放需要还原投递状态(messages.py:71-72)。

3.2 消息分三类

类别代表类型特点
单条实体消息CreateTraceMessageCreateSpanMessageUpdateSpanMessageUpdateTraceMessage一条 = 一个实体;有对应的单发 REST 端点,也可被 Batcher 攒成批
批消息CreateSpansBatchMessageCreateTraceBatchMessageAddSpanFeedbackScoresBatchMessageGuardrailBatchMessageCreateExperimentItemsBatchMessageAddAssertionResultsBatchMessage内部是一个 batch: List[...];对应后端的批量端点
纯载荷项FeedbackScoreMessageGuardrailBatchItemMessageExperimentItemMessageAssertionResultMessage源码注释写明"处理器里没有它的 handler,只作为 batch 的元素存在"(如 messages.py:239-243

Create*Message 有两个共同的小细节值得注意:

  • __post_init__ 会对 input/outputrecursive_shallow_copymessages.py:104-108)。因为消息进队列后要等好几秒才被序列化,期间业务代码可能还在改那个 dict——不拷一份就会读到"未来"的值。
  • as_payload_dict() 负责改名对齐 REST 契约:span_id → idtotal_cost → total_estimated_costmessages.py:187-191)。

3.3 supports_batching:谁决定一条消息要不要再攒一次

批消息上普遍有一个 supports_batching 布尔字段,它是给 BatchManager 看的路由开关(判定逻辑在 batching/batch_manager.py:42-49:有这个属性就用它的值,没有才看类型映射表)。

两个方向都有实例:

  • AddSpanFeedbackScoresBatchMessage 默认 supports_batching=Truemessages.py:259)——生产者已经切了一次,但仍允许 Batcher 把多个小批再合并成一个大批。合并完 Batcher 会把结果标成 supports_batching=Falsebatching/batchers.py:127-132),避免无限循环。
  • AddAssertionResultsBatchMessage 默认 supports_batching=False,源码注释直白说明原因:生产者 Opik.log_assertion_results 已经用 sequence_splitter 切好了,而且根本没给这个类型注册 batcher,让它进 BatchManager 会 KeyErrormessages.py:402-406)。

4. 入口层:Opik 类怎么"只塞不发"

这节讲什么:client.trace()Streamer.put() 之间到底发生了什么——答案是几乎什么都没发生,这正是设计意图。

4.1 所有公开入口的统一形状

不管是 trace、span、打分还是附件,Opik 上的写入口都是同一个三段式:补默认值 → 造 message → self._streamer.put(msg)

__internal_api__trace__ 是最典型的一个(opik_client.py:336-412):生成 id、取本地时间戳、解析 project_name,然后:

self._streamer.put(create_trace_message) # opik_client.py:405
self._display_trace_url(trace_id=id, project_name=project_name)

put 之后立刻返回一个 Trace 句柄,句柄里揣着同一个 message_streameropik_client.py:406),所以后续 trace.span(...)span.end(...) 也都是往同一个队列塞。span 的实际构造在 span_client.create_spansdks/python/src/opik/api_objects/span/span_client.py:371-393),update_span 同理(span_client.py:449 造消息、:479 入队)。

命名上带 __internal_api__ 前缀的方法(__internal_api__trace____internal_api__span__)是给 @track 装饰器和集成层调的低层入口,公开的 trace()/span() 是它们的薄封装。

4.2 打分入口多做一步:先切分再入队

log_spans_feedback_scores 是少数在入队前就做批量切分的入口(opik_client.py:843-852):

for batch in sequence_splitter.split_into_batches(
score_messages,
max_payload_size_MB=opik_config.MAX_BATCH_SIZE_MB, # 5 MB
max_length=constants.FEEDBACK_SCORES_MAX_BATCH_SIZE, # 1000
):
self._streamer.put(messages.AddSpanFeedbackScoresBatchMessage(batch=batch))

为什么这里要提前切?因为用户可能一次传进来 10 万条打分——如果整坨塞进队列,队列的"有界"保护就形同虚设(一条消息占的内存不受控)。log_traces_feedback_scores:854)和 log_assertion_results:904)用的是同一个套路。

4.3 流水线装配:共享资源租约

Opik.__init__opik_client.py:109)里唯一的重活是向全局 ConnectionResourceManager 按连接键租一套共享资源(_acquire_shared_resourcesopik_client.py:192-197)——同一连接键的多个 Opik 实例复用同一套部件,最后一个释放时才整体关闭。六部件的实际拼装在 create_connection_resourcesapi_objects/connection_resources.py:136-213):

步骤做的事关键配置
1建 httpx client + REST clientcheck_tls_certificateenable_json_request_compression
2算队列上限maximal_queue_size(1_000_000) ÷ maximal_queue_size_batch_factor(10)
3建文件上传管理器file_upload_background_workers = 16
4ReplayManager(含连接探针 + SQLite)replay_batch_size=50、replay_tick_interval=0.3
5建处理器链(online + 本地 emulator)emulator 默认 active=False
6construct_online_streamer 拼出 Streamerbackground_workers = 4 个 consumer

第 2 步的除法有点意思:calculate_max_queue_sizemessage_queue.py:128-135)在开批量时把队列上限除以 10。因为开批量后队列里躺的是"批消息",一条顶 1000 条,队列长度这个单位的含义变了,上限也得跟着缩。默认值是 1_000_000 / 10 = 100_000 条批消息。

第 6 步 construct_online_streamersdks/python/src/opik/message_processing/streamer_constructors.py:20-52)有一个先有鸡还是先有蛋的处理:附件抽取处理器需要持有 streamer 才能把消息回塞,所以它必须在 streamer 造好之后message_processor.add_first(...) 插到链首(streamer_constructors.py:50)。

4.4 收尾三件套:flush / end / drain_to_processors

方法等什么timeout 语义典型调用者
flush(timeout)队列 + 批缓冲 + 文件上传全部落地超时就返回 False不报错用户显式调、评估引擎
end(timeout, flush=True)同上,然后关掉所有线程None 时退回配置的 default_flush_timeoutatexitconnection_resources.py:441)、用户
__internal_api__drain_to_processors__(timeout)等消息被处理器链吃掉超时返回 False,调用方按"尽力而为"继续评估引擎的 agentic judge

endflush=False 是显式的"不要持久性"开关,文档字符串明说它是给测试清理用的(opik_client.py:1954-1959)。


5. Streamer:管道总控

这节讲什么: Streamer 是唯一一个既被业务线程调、又被后台线程调的部件,所以它的每个方法都要考虑并发。这一节把它的四个关键方法逐个拆开。

5.1 put():两级预处理的固定顺序

put() 的主体只有二十来行(streamer.py:52-82),但顺序是写死且注释加粗强调的:

message ──► ① AttachmentsPreprocessor ──► ② BatchingPreprocessor ──► 入队
(MUST ALWAYS BE DONE FIRST) (可能吞掉消息)

为什么①必须在前? 附件预处理器会把"可能内嵌 base64 附件"的消息包一层 AttachmentSupportingMessagepreprocessing/attachments_preprocessor.py:38-41)。如果先跑批量,消息已经被 Batcher 吃进累积列表了,就再没机会检查附件。判定条件也很克制(attachments_preprocessor.py:44-55):只有 Update 消息、以及 end_time is not None 的 Create 消息(也就是"这个 span 已经写完了")才会被包,且必须 input/output/metadata 至少有一个非空。

②这一级会让消息凭空消失。 BatchingPreprocessor.preprocess 如果把消息交给了 BatchManager,就返回 Nonepreprocessing/batching_preprocessor.py:28-35),put() 里的 if preprocessed_message is not None 于是跳过入队。这不是 bug——消息此刻躺在 Batcher 的内存列表里,稍后会以批消息的形式重新出现在队列里。

put() 还有一个 force 参数。self._drain 置位后普通 put 直接丢弃,但附件抽取处理器回塞消息时用 force=Trueprocessors/attachments_extraction_processor.py:80:131),保证关闭过程中已经在管道里的消息能走完。

5.2 _idle:一个只保护"半条消息"的标志

_idleput() 开头置 False、结尾置 Truestreamer.py:57:82),且整个 put() 持有 self._lockflush()close() 也拿同一把锁再 wait_for_done(lambda: self._idle)streamer.py:109-113:173-178)。

它保护的不是"队列空了",而是"没有半条消息卡在预处理中间"。 因为 putflush 争的是同一把 RLock,flush 一旦拿到锁,写入线程必然已经走完 put 全程——这个等待在语义上是"排干正在进行的 put"。真正判断"活儿干完没有"的是 _all_done()

5.3 close(flush=True/False):两种收尾策略

close(timeout, flush=?)

┌─────────────────┴─────────────────┐
flush=True flush=False
(生产默认) (测试清理)
│ │
等 _idle → _drain=True _drain=True
批缓冲 flush 后停 批缓冲直接停(丢)
join 重放线程 打警告 + 清空队列
flush() 等队列+上传 不 join,daemon 线程自生自灭
最后才关 consumer 立刻关 consumer

flush=False 分支的那段告警文案本身就是设计意图streamer.py:134-141):

"…discarding %d queued message(s) without flushing. Data that had not yet reached the backend will be lost. Use flush=True (the default) if you need durability — flush=False is intended for short-lived tests/teardowns, not production shutdown."

它把"这个开关会丢数据、只给测试用"直接写进了运行时日志,而不是只写在文档里——用错了的人在日志里就能看见。

另外 close 做了幂等处理:self._drain 已置位时直接返回(streamer.py:104-107),因为 atexit 注册的 end 很可能在用户手动 end() 之后再触发一次。

flush=True 分支里 consumer 是最后关的(streamer.py:126),注释解释得很清楚:队列排干的过程中必须有人在消费,先关 consumer 就永远排不干。

5.4 flush vs drain_to_processors:一次分工明确的拆分

维度flush(timeout)drain_to_processors(timeout)
等批缓冲
等队列被消费完
等文件上传是(阻塞,upload_sleep_time=5
触发断线重放是(仅在有连接时)
轮询间隔0.1s0.05s
关心的目标数据到了后端数据进了本进程的处理器

drain_to_processorsstreamer.py:146-172)是为评估引擎里的 agentic judge 专门加的。原因在文档字符串和调用点都写了:评估时链上会激活 LocalEmulatorMessageProcessor——一个把 trace/span 收在内存里的本地后端模拟器;agentic judge 要读"刚跑完那次任务的 span 树和 error_info"来打分,而 client.span() 只是提交到队列,consumer 可能还没处理完,judge 就会读到过期视图(sdks/python/src/opik/evaluation/engine/engine.py:235-247)。

它跳过文件上传和重放的理由也直白:那两件事关心的是后端投递,不是本地处理器状态;而 judge 每评一条就要 drain 一次,带上 5 秒粒度的上传轮询会慢到不可用。超时返回 False 时调用方不报错,只打 debug 日志然后拿现有状态继续(engine.py:248-253)——尽力而为是这条路径的明确契约。

5.5 _all_done():关掉"已出队未处理"的竞态窗口

def _all_done(self) -> bool: # streamer.py:211-219
return (
self._message_queue.all_tasks_done() and self._batch_preprocessor.is_empty()
)

注意它用的是 all_tasks_done() 而不是 empty()。区别是一个真实存在的窗口:consumer 已经把消息 pop 出来了,HTTP 还没发完——此时 empty() 返回 True,但数据其实还在天上飞。MessageQueue 因此维护了一个 _unfinished_tasks 计数(message_queue.py:32),只有 consumer 显式调 task_done() 才递减(message_queue.py:78-87)。flush 等的是这个计数归零,不是 deque 长度归零。


6. 批量层:batching/

这节讲什么: 攒批是这条管道最直接的性能收益来源。它由四个小部件组成,各自只做一件事。

6.1 BaseBatcher:三个方法讲完攒批

BaseBatcherbatching/base_batcher.py:11-52)的接口小到可以背下来:

方法做什么
add(message)追加到 _accumulated_messages长度 ≥ max_batch_size 就立刻 flush
flush()把累积列表交给子类合成批消息,逐个走 _flush_callback,然后清空、重置计时
is_ready_to_flush()now - 上次 flush 时间 >= flush_interval_seconds

_flush_callback 在构造时被绑成 queue.putbatching/batch_manager_constuctors.py:28 等),所以 Batcher 吐出的批消息是直接进队列的,不再走一遍 Streamer.put 的预处理——这也是为什么合成时要把 supports_batching 标成 False 只是双保险。

有个容易看漏的点:flush()清空 _accumulated_messages才逐个调回调(base_batcher.py:29-33)。顺序反了的话,回调里如果又触发一次 add,就会把新消息一起卷进这一批。

6.2 七个 Batcher,两种攒法

create_batch_manager 注册了 7 个消息类型 → batcher 的映射(batch_manager_constuctors.py:67-77):

消息类型Batcher 类flush 间隔攒法
CreateSpanMessageCreateSpanMessageBatcher2.0 s逐条累积
CreateTraceMessageCreateTraceMessageBatcher2.0 s逐条累积
AddSpanFeedbackScoresBatchMessageAddSpanFeedbackScoresBatchMessageBatcher1.0 s摊平批内元素
AddTraceFeedbackScoresBatchMessageAddTraceFeedbackScoresBatchMessageBatcher1.0 s摊平批内元素
AddThreadsFeedbackScoresBatchMessageAddThreadsFeedbackScoresBatchMessageBatcher1.0 s摊平批内元素
GuardrailBatchMessageGuardrailBatchMessageBatcher1.0 s摊平批内元素
CreateExperimentItemsBatchMessageCreateExperimentItemsBatchMessageBatcher3.0 s摊平批内元素

七个的 max_batch_size 全是 1000batch_manager_constuctors.py:6-19)。间隔分三档反映的是延迟容忍度:span/trace 是主数据流(2 秒),打分和 guardrail 用户等着看(1 秒),实验项是批处理场景、没人盯着(3 秒)。

两种攒法的区别add() 的实现:

  • 逐条累积(span/trace)走 BaseBatcher.add,一条消息就是列表里一项。
  • 摊平批内元素(其余五个)重写了 add,把 message.batch 里的元素拆出来累积,并在会超过 1000 时精确切一刀:先填满、flush、剩下的留给下一批(batching/batchers.py:104-118)。所以这类 batcher 的"1000"计的是打分条数,不是消息条数。

6.3 span batcher 的去重:把"开始"消息就地扔掉

CreateSpanMessageBatcher.add 有四行很划算的代码(batchers.py:39-44):

def add(self, message: messages.CreateSpanMessage) -> None:
# remove any duplicate start span message from the batch that was already added
if message.end_time is not None:
self._remove_matching_messages(lambda x: x.span_id == message.span_id)
return super().add(message)

背景:配置项 log_start_trace_span 默认为 True,一个 span 会产生两条消息——开始时一条(end_time=None)、结束时一条(带完整 output)。如果这个 span 在 2 秒内跑完,两条消息还在同一个累积列表里,那前一条就是纯浪费:后一条是它的超集。就地删掉,网络流量直接减半。trace batcher 同样处理(batchers.py:76-81)。

6.4 FlushingThread:一个只会踢腿的定时器

FlushingThreadbatching/flushing_thread.py:9-35)是全项目最简单的类:每 0.1 秒调一次传进来的 callable,捕获异常继续跑。它的文档字符串特意声明"Knows nothing about batchers or locks"——加锁和遍历是 BatchManager.flush_ready 的事(batching/batch_manager.py:69-85)。

flush_ready 有一个补丁级细节:每个 batcher 的 flush 各自包 try/except,注释说明是为了"一个 batcher 挂了不影响同一 tick 里剩下的"。

6.5 sequence_splitter:不做 JSON 编码,只估字节数

批消息合成后还有最后一道闸:单个 HTTP 请求体不能太大(span/trace batcher 传的是 BASE_BATCH_MEMORY_LIMIT_MB = 50base_batcher.py:8)。split_into_batchesbatching/sequence_splitter.py:68-112)负责按体积再切。

难点在"怎么知道这条 span 序列化后有多大"。真去 json.dumps 一遍是 CPU 和内存双重浪费——而且这份结果会被丢掉,只为量个尺寸。于是 _get_json_size 用递归估算sequence_splitter.py:25-65):

if isinstance(obj, str):
return len(obj.encode("utf-8")) + 2 # 加的 2 是那对引号
elif isinstance(obj, dict):
size = 2 # {}
for key, value in obj.items():
size += _get_json_size(key) + _get_json_size(value) + 1 + 1 # ':' 和 ','
return size - 1 # 减掉多算的最后一个逗号

三个值得留意的取舍:

  • 函数体的文档字符串明说前提是"只收到基础 Python 对象、无循环引用"——它不是通用工具。
  • 遇到认不出的类型退化成 len(str(obj)),只打 debug 日志。
  • 任何异常都返回 float("inf")sequence_splitter.py:62-65),注释写明"to be on the safe side"——因为 inf >= max_payload_size_MB 恒成立,这条消息会走独占一批的分支(sequence_splitter.py:92-94),不会连累同批的其他消息。

那条独占分支本身也是一个明确决策:单条就超限的消息照发不误(大概率被后端拒),而不是就地丢弃。SDK 不替用户做"你这条太大我不发了"的决定。


7. 队列与限流

这节讲什么: 消息进队列之后到出队之间,管道要回答两个问题——队列满了丢谁?后端说"太快了"怎么办?

7.1 有界队列:丢老不丢新

MessageQueue 底层是 collections.deque(maxlen=...)message_queue.py:30),写入用 appendleft、读取用 pop——标准 FIFO。maxlen 一满,appendleft静默地从另一端挤掉最老的那条。

所以 Streamer.put 在入队前先问一句(streamer.py:71-76):

if self._message_queue.accept_put_without_discarding() is False:
_logging.log_once_at_level(
logging.WARNING,
"The message queue size limit has been reached. The new message has been "
"added to the queue, and the oldest message has been discarded.",
logger=LOGGER,
)

两个设计点:

  1. 丢老留新。队列积压到 10 万条批消息,多半是后端挂了或长时间限流;这时最近的数据比几分钟前的数据有价值。
  2. log_once_at_level。队列一满就是持续满,逐条打 warning 会把用户日志淹掉——只打第一次。

队列被塞满时那个被挤掉的消息,会让 _unfinished_tasks 计数对不上(它永远等不到 task_done)。put/put_back/clear 三处都专门处理了这个记账(message_queue.py:40-43:57-64:105-113),否则 flush 会永远等一个幽灵任务。

7.2 QueueConsumer._loop:限流退避与保序

consumer 线程的主循环只做四件事(queue_consumer.py:36-75):

_loop() 每轮

now < next_message_time ? ──是──► sleep 0.1s,返回(全局退避中)
│否
从队列 get(超时 0.1s)

msg.delivery_time <= now ? ──否──► put_back,保序(这条还没到点)
│是
_process_message(msg)

抛 RateLimited ? ──是──► next_message_time = now + retry_after
msg.delivery_time = next_message_time
put_back,保序

限流处理有两层,都不可少:

  • 线程级self.next_message_time 让这个 consumer 在 retry_after 秒内完全不干活(queue_consumer.py:66)。
  • 消息级message.delivery_time 打在消息上(:68)。因为有 4 个 consumer,这条消息可能被别的 consumer 立刻捞走;不打在消息上的话退避就漏了。

put_back 用的是 append(写到 get 取的那一端,message_queue.py:59),效果是这条消息回到队首下一个被取的位置——注释里两次写了 "to keep an order in the queue"。限流退避不能打乱 trace/span 的到达顺序。

_push_message_back 顺手把 delivery_attempts 自增(queue_consumer.py:87)。这个计数的下游消费者是本地 emulator:LocalEmulatorMessageProcessor.process 看到 delivery_attempts > 1 且实体已记录,就跳过,避免重试在本地视图里造出重复 span(emulation/local_emulator_message_processor.py:33-60)。

7.3 _process_message 的三分支记账

task_done() 调多了会 raise ValueErrormessage_queue.py:83-84),调少了 flush 会挂死。所以 _process_message 把三种结局分得很清楚(queue_consumer.py:90-105):

结局task_done为什么
处理成功正常路径
OpikCloudRequestsRateLimited不调消息马上要被 put_back,任务仍在飞行中
抛其他异常调,然后重新抛永久失败,不会重试,不记账就会卡住 flush

8. 处理器层:真正发 HTTP 的地方

这节讲什么: 消息出队后进入处理器链。链上有三个成员,其中只有一个真的发网络请求。

8.1 ChainedMessageProcessor:一条消息喂给所有处理器

链本身很薄(processors/message_processors.py:25-92):按序把消息喂给每个处理器,单个处理器抛异常只记日志、不中断链:50-56)。唯一的例外是 OpikCloudRequestsRateLimited——它被暂存下来,等整条链跑完再重新抛出(:43:58-60),因为限流必须传到 QueueConsumer 才能触发退避。

链的组装在 create_message_processors_chainprocessors/message_processors_chain.py:19-66),初始是两个成员:

  1. OpikMessageProcessor——发 REST。
  2. LocalEmulatorMessageProcessor——本地模拟后端,默认 active=False

第三个成员 AttachmentsExtractionProcessor 在 streamer 造好后被 add_first 插到链首(streamer_constructors.py:41-50)。它专门处理 AttachmentSupportingMessage:抽出内嵌的 base64 附件、替换成引用、把原消息和新生成的 CreateAttachmentMessage 一起 put(force=True) 回队列(processors/attachments_extraction_processor.py:64-80:119-131)。回塞前给原消息打一个标记属性,附件预处理器见到标记就直接放行——这就是防无限递归的闸(attachments_preprocessor.py:34-36)。

本地 emulator 由 toggle_local_emulator_message_processor 在评估开始/结束时开关(message_processing/processors/message_processors_chain.py:69),详见 04-evaluation-engine.md。它存在的意义是:评估要读 trace 的完整结构来打分,如果每次都去查后端,评估就慢得没法用了。

8.2 OpikMessageProcessor:一张表 + 一个统一错误分诊

这个类的主体是一张 message_type → handler 的字典(processors/online_message_processor.py:62-77),14 个条目。每个 handler 都短得一眼看完,比如:

def _process_create_spans_batch_message(self, message): # online_message_processor.py:345-350
self._rest_client.spans.create_spans(spans=message.batch)

单条消息的 handler 会多做两步 payload 清洗:remove_none_from_dict 去掉 None 字段,encode_and_anonymize 做编码和脱敏(:216-223)。批消息不用——batcher 合成时已经洗过了(batchers.py:18-26)。

process() 的价值不在于分发,而在于那个统一的错误分诊online_message_processor.py:142-244):

情况处理
HTTP 409 冲突静默返回。注释解释:重试机制有时会重发同一请求,用户不该看到这个错
HTTP 429解析限流头,转成 OpikCloudRequestsRateLimited 抛出,交给 consumer 退避
HTTP 401把这个 message_type 记进未授权黑名单,以后同类消息直接不发
httpx.ConnectError / TimeoutException标记为待重放,交给 ReplayManager;日志是 warning 不是 error
pydantic.ValidationErrorerror 日志,不重试(数据本身有问题,重试也没用)
其他error 日志 + "检查 Opik 配置"的提示

401 那条是个自愈式降级:某些消息类型(比如企业版才有的 guardrail)在当前部署上没权限,与其每条都打一次 401,不如注册进 UnauthorizedMessageTypeRegistry,后续同类消息在 process 入口就被拦掉(:89-96)。注册表带重试间隔和最大重试次数配置,所以是临时拉黑而非永久。

处理器层不做本地存储——服务端 ClickHouse 怎么吞下这些乱序到达的 trace 与 span,见 03-storage-and-query.md


9. 断线续传:replay/

这节讲什么: 前面所有机制都假设网络是通的。这一节讲网断了之后消息去了哪、怎么回来。

9.1 三个部件的分工

部件是什么职责
OpikConnectionMonitor状态机定期 ping 后端,判定 ok / failed / restored
ReplayManagerdaemon 线程每 0.3 秒 tick 一次 monitor;一旦 restored 就触发重放
DBManagerSQLite 封装消息落盘、状态流转、按批取出失败消息

9.2 消息在 SQLite 里的三个状态

put → consumer → OpikMessageProcessor.process()

register_message

┌───────────────┴───────────────┐
有连接 无连接
│ │
status=registered status=failed
发 HTTP (根本不发)
│ │
┌─────┴─────┐ │
成功 连接错误 │
│ │ │
delivered failed ◄───────────────────┘
(删行) │
│ 连接恢复 → replay_failed_messages
└──► 改回 registered → Streamer.put() → 重走全流程

关键代码在 OpikMessageProcessor.process 的开头(online_message_processor.py:125-140):有连接就先登记再发,没连接就直接登记成 failed 并跳过发送。后一条很重要——已知断网时连试都不试,省掉每条消息一次超时等待。

ReplayManager 的方法名直接对应这三个状态:register_messagereplay/replay_manager.py:85)、unregister_message:104,实际是标 delivered)、message_sent_failed:118,同时通知 monitor 连接已断)。

重放的入口是 monitor 的状态跃迁(replay_manager.py:155-160):

status = self._monitor.tick()
if status == connection_monitor.ConnectionStatus.connection_restored:
self._replay_failed_messages()
self._monitor.reset()

connection_restored 只在"上次是断的、这次 ping 通了"时返回一次(healthcheck/connection_monitor.py:129-135),所以重放只在边沿触发,不会每 tick 都刷一遍 SQLite。

重放的 callback 就是 Streamer.put——在 Streamer.__init__ 里绑定(streamer.py:45)。所以被重放的消息会完整重走一遍管道:预处理、攒批、限流判定一个不少。

9.3 序列化:为什么不能直接 pickle

message_serialization 走的是 JSON(replay/message_serialization.py:38-57):as_db_message_dict()jsonable_encoder.encodejson.dumps。反序列化时按 message_type 字符串查 SUPPORTED_MESSAGE_TYPES 表定位类(replay/db_manager.py:702-742),认不出的类型抛 ValueError

反序列化有一个精心限制的细节:datetime 只对白名单字段还原(message_serialization.py:19):

DATETIME_FIELD_NAMES: Set[str] = {"start_time", "end_time", "last_updated_at"}

注释讲明了原因——用户的 input/output/metadata 里可能有长得像 ISO 时间的普通字符串,无差别转换会把用户数据改掉。

messages.from_db_message_dictmessages.py:22-45)则处理另一个坑:delivery_time/message_id 这些是 init=False 字段,不能当构造参数传,得先建对象再 setattr 回去。

9.4 DB 挂了怎么办:降级,不是崩溃

DBManager 有四个状态:undefined / initialized / closed / errordb_manager.py:55-65)。任何 SQLite 操作失败都会走 _mark_as_db_faileddb_manager.py:693-699):

def _mark_as_db_failed(self, message: str) -> None:
self.status = DBManagerStatus.error
LOGGER.error(
"Due to an internal error, some network resiliency features were disabled "
"which could lead to data loss. Contact us at [email protected]. ..."
)

进入 errorinitialized 属性为 Falsereplay_failed_messages 开头就直接返回 0(db_manager.py:446-448)。整条上报管道照常工作,只是没有断网兜底了。 告警文案把后果(可能丢数据)和补救途径都写清楚,而不是抛异常打断用户的业务代码——这是这套上报层贯穿始终的态度。

replay_failed_messages 里还有两处并发处理值得看:

  • 一把非阻塞的 _replay_mutex:并发调用时第二个直接返回 0,避免两个线程捞到同一批失败消息(db_manager.py:450-457)。
  • 主锁只在最小临界区里持有(关闭检查、把一批标成 in-progress),重放回调和批间 sleep 都在锁外执行(db_manager.py:502-513)。因为回调是 Streamer.put,持锁调它会把所有业务线程的 register_message 全堵住。

10. 巧妙之处(可以带走的)

  1. 不真做 JSON 编码,只估字节数。 为了量尺寸而 json.dumps 一遍再丢掉,是纯浪费;_get_json_size 递归估算,估不准就返回 inf 让那条消息独占一批(batching/sequence_splitter.py:25-65)。

  2. all_tasks_done() 而不是 empty() 队列空 ≠ 活干完——消息可能已出队但 HTTP 还没回。引入未完成任务计数,把"已出队未处理"这个窗口关掉(message_queue.py:78-91streamer.py:216-224)。

  3. 告警文案当设计文档用。 close(flush=False) 丢消息时打的那段警告,把"会丢数据、只给测试用、生产别碰"直接写进运行时日志(streamer.py:134-141);DB 降级的告警同理。用错的人在日志里就看得到,不用去翻文档。

  4. span 的"开始"消息就地去重。 同批里出现同一 span 的开始+结束两条消息时,直接删掉开始那条——后者是前者的超集,流量减半,四行代码(batching/batchers.py:39-44)。

  5. 限流退避打两个地方。 只设 consumer.next_message_time 会漏(还有 3 个 consumer),只设 message.delivery_time 会空转。两个都设,且回塞用 append 保序(queue_consumer.py:57-70)。

  6. 无连接时连试都不试。 已知断网就把消息直接标 failed 落盘、跳过 HTTP,省掉每条消息一次超时(online_message_processor.py:133-137)。

  7. 降级而非崩溃。 SQLite 挂了、消息类型 401 了、附件抽取失败了——每一条路径都是"记日志 + 关掉这个功能 + 继续跑"。观测组件把宿主应用搞崩是最不可接受的失败模式。

  8. 为一个特定读者拆出一个轻量 drain。 drain_to_processors 只等本地处理器、不等文件上传和重放,就为了让评估里的 agentic judge 能高频调用而不被拖垮(streamer.py:146-172)。


11. 边界与局限

它刻意不做的事:

  • 不保证送达。 队列满了丢最老的、close(flush=False) 丢全部、进程被 kill -9 丢内存里的一切。SQLite 只兜"发送时连接失败"这一种情况,不是 WAL 式的全量持久化。
  • 不做加密/重签名。 payload 里有什么就发什么,脱敏交给 encode_and_anonymize 的字段白名单(fields_to_anonymize() 只覆盖 input/output/metadata)。
  • 不替用户拒发超大消息。 单条超过 50 MB 也照发,大概率被后端拒(sequence_splitter.py:92-94)。

已知的坑:

  • 开批量 + update_* 有竞态。 Opik.__init__ 的文档字符串明写:开批量时 update_span/update_trace 可能在批量的 create 请求 flush 之前到达服务端,导致更新丢失(opik_client.py:128-131)。SDK 提供了 warn_if_batching_update 这个显式提醒(opik_client.py:709-713),但没有在协议层解决。
  • 投递顺序只在队列内保证。 攒批把 span 和 trace 分到了不同的 batcher、不同的 flush 间隔,多个 consumer 又并发发送——到达服务端的顺序基本无序。这个包袱是服务端接的,见 03-storage-and-query.md
  • _idle 不是"队列空"。 名字容易误导,它只表示"没有 put 正在预处理中"。
  • 重放会重发。 重放的消息重走整条管道,服务端可能收到重复;靠后端幂等和 409 静默兜底(online_message_processor.py:142-148)。

12. 与其他章的关系

想知道什么去哪章
span 树是怎么被构造出来的、@track 干了什么01-tracing-decorator.md
乱序到达的 trace/span 服务端怎么存03-storage-and-query.md
本地 emulator 在评估里怎么用、drain_to_processors 的调用方04-evaluation-engine.md
打分数据从哪来(本章只负责运走它们)05-metrics-and-judges.md
在线评估规则怎么触发 guardrail 类消息06-online-scoring-and-optimizer.md

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

所有路径相对克隆根,前缀统一为 sdks/python/src/opik/

主题文件路径符号名
客户端写入口(造消息 + put)api_objects/opik_client.pyOpik.__internal_api__trace__Opik.__internal_api__span__Opik.log_spans_feedback_scores
流水线组装api_objects/opik_client.pyOpik._initialize_streamerOpik._create_replay_manager
收尾语义api_objects/opik_client.pyOpik.flushOpik.endOpik.__internal_api__drain_to_processors__Opik.__internal_api__failed_uploads__
span 消息构造api_objects/span/span_client.pycreate_spanupdate_span
消息类型体系message_processing/messages.pyBaseMessageCreateSpanMessageUpdateSpanMessageAddSpanFeedbackScoresBatchMessageGuardrailBatchMessageCreateExperimentItemsBatchMessagefrom_db_message_dict
管道总控message_processing/streamer.pyStreamer.putStreamer.closeStreamer.flushStreamer.drain_to_processorsStreamer._all_done
流水线构造message_processing/streamer_constructors.pyconstruct_online_streamerconstruct_streamer
两级预处理message_processing/preprocessing/AttachmentsPreprocessor.preprocessBatchingPreprocessor.preprocess
攒批基类message_processing/batching/base_batcher.pyBaseBatcher.addBaseBatcher.flushBaseBatcher.is_ready_to_flushBASE_BATCH_MEMORY_LIMIT_MB
各类 batchermessage_processing/batching/batchers.pyCreateSpanMessageBatcherBaseAddFeedbackScoresBatchMessageBatcherGuardrailBatchMessageBatcherCreateExperimentItemsBatchMessageBatcher
批量参数与注册表message_processing/batching/batch_manager_constuctors.pycreate_batch_managerCREATE_SPANS_MESSAGE_BATCHER_FLUSH_INTERVAL_SECONDS
批管理与定时 flushmessage_processing/batching/batch_manager.pyflushing_thread.pyBatchManager.flush_readyBatchManager.message_supports_batchingFlushingThread
体积切分与尺寸估算message_processing/batching/sequence_splitter.pysplit_into_batches_get_json_size
有界队列与任务记账message_processing/message_queue.pyMessageQueue.putMessageQueue.put_backMessageQueue.task_doneMessageQueue.all_tasks_doneaccept_put_without_discardingcalculate_max_queue_size
消费线程与限流message_processing/queue_consumer.pyQueueConsumer._loopQueueConsumer._process_messageQueueConsumer._push_message_back
处理器链message_processing/processors/message_processors.pymessage_processors_chain.pyChainedMessageProcessor.processcreate_message_processors_chaintoggle_local_emulator_message_processor
REST 发送与错误分诊message_processing/processors/online_message_processor.pyOpikMessageProcessor.process_process_create_spans_batch_message
附件抽取message_processing/processors/attachments_extraction_processor.pyAttachmentsExtractionProcessor.process
本地模拟处理器message_processing/emulation/local_emulator_message_processor.pyLocalEmulatorMessageProcessor_should_skip_retry
断线重放线程message_processing/replay/replay_manager.pyReplayManager.register_messagemessage_sent_failed_replay_failed_messages
SQLite 落盘message_processing/replay/db_manager.pyDBManager.replay_failed_messagesMessageStatusDBManagerStatus_mark_as_db_failedSUPPORTED_MESSAGE_TYPES
消息序列化message_processing/replay/message_serialization.pyserialize_messagedeserialize_messageDATETIME_FIELD_NAMES
连接状态机healthcheck/connection_monitor.pyOpikConnectionMonitor.tickConnectionStatus
相关配置项config.pybackground_workersmaximal_queue_sizemaximal_queue_size_batch_factorreplay_batch_sizeconnection_monitor_ping_interval