Skip to content

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=Truestop_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_timeusage_detailsoutput
进行中NoneNone持续 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() 会:

  1. match_subagent() 找到 agent-*.jsonl
  2. tools[tool_id]["subagent_stream"] 中维护子代理的 offset/buffer
  3. 增量 read_new_jsonl() 读子代理新行
  4. build_turns(include_pending=True) 切分子 Turn
  5. 为每个子 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