主题
6. 流式上报机制
默认开启:
CC_LANGFUSE_STREAMING=true(Watcher 方案)
关闭后回退到批量模式(Turn 完整结束才上报)
1. 为什么需要流式?
批量模式下,一整轮对话可能持续几分钟(多次 Tool 调用、子代理执行),Langfuse UI 里什么都看不到,直到 end_turn。
流式模式的目标:
| 时机 | 批量模式 | 流式模式 |
|---|---|---|
| 用户发消息 | 无 | 立即创建 Trace |
| 模型开始回复 | 无 | 创建 Generation(output 持续更新) |
| 工具调用发起 | 无 | 立即创建 Tool Span |
| 工具返回结果 | 无 | 结束 Tool Span |
| 子代理执行中 | 无 | 增量监听 subagent jsonl |
| end_turn | 一次性上报全部 | finalize Trace |
生活类比:批量 = 做完菜才拍照;流式 = 边做边直播。
2. 流式状态机
每个进行中的 Turn 在 SessionState.streaming 中维护:
python
streaming = {
"trace_id": "uuid", # 已创建的 Trace ID
"turn_num": 3,
"generations": {
"msg_01ABC": {
"id": "gen-uuid-1",
"done": False,
"last_output": "正在分析..."
}
},
"tools": {
"toolu_01XYZ": {
"id": "span-uuid-1",
"ended": False,
"subagent": False,
"subagent_stream": { ... } # 子代理增量状态
}
},
"used_subagents": ["agent-id-1"]
}3. 核心函数调用链
process_turn_streaming()
├── _streaming_open_trace() # 首次:创建 Trace;后续:resume
├── _streaming_sync_subagent_open_tools() # 子代理增量
└── _process_turn_in_container()
├── _streaming_sync_generations()
└── _streaming_sync_tools()finalize=True(stop_reason == end_turn)时额外:
- 结束未完成的 Generation(写入 usage)
- 结束未返回结果的 Tool Span(标记 incomplete)
trace.update(output=..., metadata={finalized: true})- 清空
ss.streaming
4. Trace 创建与 Resume
python
def _streaming_open_trace(langfuse, streaming, session_id, turn_num, turn, ...):
trace_id = streaming.get("trace_id")
if trace_id:
return _StreamingTraceHandle(langfuse, trace_id, session_id), False
# 首次:创建新 Trace
trace_id = str(uuid.uuid4())
langfuse.trace(id=trace_id, name=f"Claude Code - Turn {turn_num}", ...)
streaming["trace_id"] = trace_id
return _StreamingTraceHandle(...), True_StreamingTraceHandle 包装后续调用,始终带 trace_id:
python
def generation(self, **kwargs):
return self._langfuse.generation(trace_id=self.trace_id, **kwargs)
def span(self, **kwargs):
return self._langfuse.span(trace_id=self.trace_id, **kwargs)关键:轮询时不会重复创建 Trace,而是在已有 Trace 下追加 observation。
5. Generation 增量同步
python
def _streaming_sync_generations(trace, turn, streaming, finalize=False):
calls = _group_assistant_calls(turn.assistant_msgs)
gens = streaming.setdefault("generations", {})
for gi, g in enumerate(calls):
mid = g["mid"]
entry = gens.get(mid)
if entry is None:
# 新模型调用 → 创建 Generation(end_time=None 表示进行中)
trace.generation(id=gen_id, ..., end_time=None if not finalize else ...)
gens[mid] = {"id": gen_id, "done": finalize}
elif not entry["done"]:
if finalize:
gen.end(output=..., usage_details=..., end_time=...)
entry["done"] = True
else:
# 进行中 → 更新 output 文本
gen.update(output=..., metadata=...)进行中 vs 完成
| 状态 | end_time | usage_details | output |
|---|---|---|---|
| 进行中 | None | None | 持续 update |
| 完成 | 真实时间戳 | 完整 token | 最终文本 |
6. Tool Span 增量同步
python
def _streaming_sync_tools(trace, turn, streaming, sub_files, used_subagents, finalize=False):
for tc in tool_calls:
if tool_id not in tools:
# 新工具调用 → 立即创建 Span(只有 start_time)
trace.span(id=span_id, name=f"Tool: {name}", start_time=...)
tools[tool_id] = {"id": span_id, "ended": False}
if tool_id in turn.tool_results_by_id:
# 有结果 → 结束 Span
_streaming_end_tool_span(...)
elif finalize and not tools[tool_id]["ended"]:
# Turn 结束但工具无结果 → 标记 incomplete
span.end(output=None, metadata={"incomplete": True})时间线效果:
10:00:02 Tool: Bash 开始(Langfuse 立刻出现)
10:00:05 Tool: Bash 结束(填入 output)7. 子代理流式监听
当 Tool: Agent / Tool: Task 的 Span 尚未结束时,_streaming_sync_subagent_open_tools() 会:
match_subagent()找到agent-*.jsonl- 在
tools[tool_id]["subagent_stream"]中维护子代理的 offset/buffer - 增量
read_new_jsonl()读子代理新行 build_turns(include_pending=True)切分子 Turn- 为每个子 Turn 创建嵌套 Span,递归
_process_turn_in_container()
python
subagent_stream = {
"path": "/path/to/agent-abc.jsonl",
"agent_id": "abc",
"agent_type": "Explore",
"offset": 12345,
"buffer": "",
"messages": [...],
"turn_count": 1,
"turns": {
"0": {"turn_span_id": "...", "streaming": {}, "done": True},
"1": {"turn_span_id": "...", "streaming": {}, "done": False}
}
}效果:主会话的 Agent 工具还在执行时,Langfuse 里就能展开看到子代理内部的 Generation 和 Tool 调用。
8. 流式时序图
9. 关闭流式
bash
export CC_LANGFUSE_STREAMING=false回退行为:
- Watcher:
completed_turns()+emit_main_turn()批量上报 - Hook:本来就不走流式(始终批量)
10. 流式模式的代价
| 方面 | 说明 |
|---|---|
| 状态复杂度 | streaming 字段较大,需正确 finalize |
| API 调用次数 | 轮询 + update 比批量多几次 HTTP |
| 中途失败 | 未 finalize 的 Turn 下次轮询继续 resume |
| UI 体验 | 可实时观察 AI 执行过程,值得 |
下一章:7_setup_guide.md