目 录CONTENT

文章目录

大模型调用链路追踪:OpenTelemetry + LangFuse 实践

PySuper
2025-06-28 / 0 评论 / 0 点赞 / 3 阅读 / 0 字
温馨提示:
本文最后更新于2026-05-25,若内容或图片失效,请留言反馈。 所有牛逼的人都有一段苦逼的岁月。 但是你只要像SB一样去坚持,终将牛逼!!! ✊✊✊

在大模型应用调试中,你是否遇到过这样的困境:明明 prompt 很简单,响应却莫名其妙地慢;Token 消耗异常高,却找不到原因;线上出问题,只能靠猜……这些问题都指向一个核心诉求:LLM 调用需要可观测性。本文将介绍如何通过 OpenTelemetry + LangFuse 实现完整的大模型调用链路追踪。


一、LLM 调用链路追踪的特殊性

1.1 为什么传统 APM 不够用?

传统 API 监控(APM)工具如 Jaeger、Zipkin 主要针对微服务架构设计,它们擅长追踪 HTTP 请求、数据库查询、消息队列等。但 LLM 调用有其独特性:

┌─────────────────────────────────────────────────────────────────────────────┐
│                    传统 APM vs LLM 追踪的差异                                │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│  传统 APM 追踪内容:                                                         │
│  ┌──────────────────────────────────────────────────────────────────────┐   │
│  │  [HTTP Request] → [Auth] → [DB Query] → [Cache] → [HTTP Response]  │   │
│  │   ↑                                      ↑                           │   │
│  │   └────── Span: 50ms ────── Span: 100ms ─┘                          │   │
│  └──────────────────────────────────────────────────────────────────────┘   │
│  特点:调用链路清晰,耗时分布均匀                                            │
│                                                                             │
│  LLM 追踪内容:                                                              │
│  ┌──────────────────────────────────────────────────────────────────────┐   │
│  │  [Preprocess] → [LLM Call] → [Postprocess] → [Response]            │   │
│  │                     ↑                                                  │   │
│  │   ┌──────────────────────────────────────────────────────────────┐   │   │
│  │   │  • Model: gpt-4o                                              │   │   │
│  │   │  • Input Tokens: 2,847 (巨大!)                                │   │   │
│  │   │  • Output Tokens: 523                                         │   │   │
│  │   │  • Cache Hit: false                                           │   │   │
│  │   │  • Reasoning Steps: 12 (如果有)                               │   │   │
│  │   │  • Time to First Token: 1.2s (流式特有)                        │   │   │
│  │   │  • Total Duration: 4.8s (可能很慢)                             │   │   │
│  │   └──────────────────────────────────────────────────────────────┘   │   │
│  └──────────────────────────────────────────────────────────────────────┘   │
│  特点:LLM 调用是主要耗时点,且有独特的度量维度                               │
│                                                                             │
└─────────────────────────────────────────────────────────────────────────────┘

1.2 LLM 追踪的核心维度

维度

说明

重要性

Token 分布

输入/输出/缓存 Token 数量

成本分析、性能优化

模型版本

gpt-4-0613 vs gpt-4-0125-preview

问题定位、A/B 测试

缓存命中

是否命中语义/精确缓存

成本优化

推理时间

TTFT、Total Duration

性能调优

Token 速率

tokens/second

实时性能监控

错误类型

Rate Limit、Timeout、Content Filter

稳定性分析

Prompt 版本

哪个版本的 Prompt

回归分析

会话上下文

多轮对话的 Turn 数

上下文管理

1.3 链路追踪的典型应用场景

  1. 性能问题定位

  • 哪个环节耗时最长?是 LLM 调用还是后处理?

  • Token 生成速度是否符合预期?

  • 是否遇到限流?

  1. 成本异常排查

  • 为什么今天的 Token 消耗是昨天的 3 倍?

  • 哪个用户的 Token 消耗异常高?

  • 缓存命中率为什么下降了?

  1. Prompt 迭代分析

  • 新版 Prompt 比旧版多用多少 Token?

  • 输出质量变化与 Token 消耗的关系?

  1. 模型对比实验

  • GPT-4 vs Claude-3.5 的响应质量与成本对比

  • 不同模型在不同任务上的表现差异


二、OpenTelemetry 基础与 LLM 适配

2.1 OpenTelemetry 核心概念

┌─────────────────────────────────────────────────────────────────────────────┐
│                         OpenTelemetry 核心概念                               │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│  Trace (追踪)                                                               │
│  ├── Span (跨度) ← 一次操作的基本单元                                        │
│  │   ├── name: 操作名称                                                    │
│  │   ├── trace_id: 追踪ID                                                  │
│  │   ├── span_id: 跨度ID                                                   │
│  │   ├── parent_span_id: 父跨度ID                                          │
│  │   ├── start_time / end_time: 时间范围                                   │
│  │   ├── attributes: 键值属性 ← LLM 特有信息放这里                          │
│  │   ├── events: 子事件(流式输出的每个 token)                              │
│  │   └── status: 状态 (OK, ERROR)                                          │
│  │                                                                            │
│  │   ├── [Preprocess Span]                                                  │
│  │   │   └── [LLM Call Span] ← 核心 LLM 追踪                               │
│  │   │       ├── input_tokens: 2847                                        │
│  │   │       ├── output_tokens: 523                                        │
│  │   │       ├── model: gpt-4o                                             │
│  │   │       └── events: [Token 1] [Token 2] ... [Token 523]               │
│  │   │                                                                            │
│  │   └── [Postprocess Span]                                                  │
│  │                                                                            │
│  Metrics (指标) ← 聚合的数值                                                  │
│  └── Events (事件) ← 离散的发生事件                                           │
│                                                                             │
│  ┌───────────────────────────────────────────────────────────────────────┐ │
│  │                        Collector Pipeline                              │ │
│  │                                                                        │ │
│  │   App → OTLP → Collector → Backend (Jaeger/Tempo/LangFuse)            │ │
│  │         │                   │                                          │ │
│  │         │                   └──可同时发送多个后端                       │ │
│  │         │                                                              │ │
│  │         └── 支持 HTTP/gRPC 协议                                        │ │
│  └───────────────────────────────────────────────────────────────────────┘ │
│                                                                             │
└─────────────────────────────────────────────────────────────────────────────┘

2.2 LLM 语义约定 (Semantic Conventions)

OpenTelemetry 正在制定 LLM 相关的语义约定,以下是推荐使用的属性:

# otel_llm_semconv.py
"""
OpenTelemetry LLM 语义约定
参考: https://opentelemetry.io/docs/specs/semconv/gen-ai/
"""

from opentelemetry import trace
from opentelemetry.trace import SpanKind, Status, StatusCode


class LLMTraceAttributes:
    """
    LLM 调用追踪属性定义
    
    这些属性遵循 OpenTelemetry 语义约定规范
    """
    
    # ============================================================
    # 必需属性 (Required)
    # ============================================================
    
    # 生成式 AI 系统类型
    GEN_AI_SYSTEM = "gen_ai.system"  # 固定值: "openai"
    
    # 生成式 AI 供应商
    GEN_AI_PROVIDER = "gen_ai.provider"
    
    # 模型名称
    GEN_AI_MODEL = "gen_ai.model"
    
    # 请求类型
    GEN_AI_REQUEST_MODEL = "gen_ai.request.model"
    GEN_AI_RESPONSE_MODEL = "gen_ai.response.model"
    
    # ============================================================
    # Token 相关属性
    # ============================================================
    
    # 输入 Token 数
    GEN_AI_REQUEST_TOKEN_COUNT = "gen_ai.request.token_count"
    
    # 输出 Token 数
    GEN_AI_RESPONSE_TOKEN_COUNT = "gen_ai.response.token_count"
    
    # Token 使用统计(结构化)
    GEN_AI_USAGE_PROMPT_TOKENS = "gen_ai.usage.prompt_tokens"
    GEN_AI_USAGE_COMPLETION_TOKENS = "gen_ai.usage.completion_tokens"
    GEN_AI_USAGE_TOTAL_TOKENS = "gen_ai.usage.total_tokens"
    
    # ============================================================
    # Prompt 和 Response 内容
    # ============================================================
    
    # Prompt 内容(可选,需要注意隐私)
    GEN_AI_PROMPT = "gen_ai.prompt"
    
    # Response 内容(可选,需要注意隐私)
    GEN_AI_COMPLETION = "gen_ai.completion"
    
    # Prompt 角色
    GEN_AI_PROMPT_ROLE = "gen_ai.prompt.role"
    
    # ============================================================
    # 缓存相关
    # ============================================================
    
    # 缓存命中
    GEN_AI_RESPONSE_CACHE_HIT = "gen_ai.response.cache_hit"
    
    # 缓存创建 Token 数(首次生成)
    GEN_AI_USAGE_CACHE_CREATION_TOKENS = "gen_ai.usage.cache_creation_tokens"
    
    # 缓存读取 Token 数(命中缓存)
    GEN_AI_USAGE_CACHE_READ_TOKENS = "gen_ai.usage.cache_read_tokens"
    
    # ============================================================
    # 速率限制
    # ============================================================
    
    # 速率限制剩余
    GEN_AI_RESPONSE_RATE_LIMIT_REMAINING = "gen_ai.response.rate_limit.remaining"
    
    # 速率限制重置时间
    GEN_AI_RESPONSE_RATE_LIMIT_RESET = "gen_ai.response.rate_limit.reset"
    
    # ============================================================
    # 推理相关 (Reasoning Models)
    # ============================================================
    
    # 推理 Token 数(Thinking Models)
    GEN_AI_USAGE_REASONING_TOKENS = "gen_ai.usage.reasoning_tokens"
    
    # ============================================================
    # 自定义业务属性
    # ============================================================
    
    # 项目/用户 ID(业务维度)
    PROJECT_ID = "project.id"
    USER_ID = "user.id"
    SESSION_ID = "session.id"
    
    # 请求类型
    REQUEST_TYPE = "llm.request.type"  # chat, completion, embedding
    TASK_TYPE = "llm.task.type"  # qa, summarization, translation
    
    # 版本标识
    PROMPT_VERSION = "llm.prompt.version"
    
    # 自定义追踪属性
    TRACE_ID = "llm.trace.id"
    SPAN_ID = "llm.span.id"
    
    # 错误信息
    ERROR_TYPE = "error.type"
    ERROR_MESSAGE = "error.message"


def create_llm_span(
    tracer,
    name: str,
    model: str,
    provider: str,
    input_tokens: int = None,
    output_tokens: int = None,
    **kwargs
) -> trace.Span:
    """
    创建 LLM 调用 Span 的辅助函数
    
    示例:
        with tracer.start_as_current_span(
            "llm.chat",
            kind=SpanKind.CLIENT
        ) as span:
            set_llm_span_attributes(
                span,
                model="gpt-4o",
                provider="openai",
                input_tokens=100,
                output_tokens=50
            )
    """
    span = tracer.start_span(name, kind=SpanKind.CLIENT)
    
    # 设置必需属性
    span.set_attribute(LLMTraceAttributes.GEN_AI_SYSTEM, "gen_ai")
    span.set_attribute(LLMTraceAttributes.GEN_AI_PROVIDER, provider)
    span.set_attribute(LLMTraceAttributes.GEN_AI_MODEL, model)
    
    # 设置 Token 属性
    if input_tokens is not None:
        span.set_attribute(LLMTraceAttributes.GEN_AI_REQUEST_TOKEN_COUNT, input_tokens)
    
    if output_tokens is not None:
        span.set_attribute(LLMTraceAttributes.GEN_AI_RESPONSE_TOKEN_COUNT, output_tokens)
    
    # 设置其他属性
    for key, value in kwargs.items():
        span.set_attribute(key, value)
    
    return span

2.3 追踪器初始化配置

# config/otel.yaml
# OpenTelemetry 配置

service:
  name: "llm-gateway"
  version: "1.0.0"
  environment: "production"

# OTLP 导出器配置
otlp:
  # gRPC 导出器(推荐生产环境)
  grpc:
    enabled: true
    endpoint: "localhost:4317"
    insecure: false  # 生产环境设为 true 并配置 TLS
    
  # HTTP 导出器(简单部署)
  http:
    enabled: false
    endpoint: "http://localhost:4318/v1/traces"

# 本地开发:打印到控制台
console:
  enabled: true
  pretty_print: true

# 采样配置
sampling:
  # 采样策略: always_on, always_off, trace_id_ratio, parent_based
  strategy: "parent_based"
  
  # 父级基础采样
  parent_based:
    # 如果有父span,继承父span的采样决策
    # 如果没有父span,使用以下配置
    root:
      sampler: "trace_id_ratio"
      fraction: 0.1  # 采样 10% 的请求
  
  # 固定采样(开发环境)
  always_on:
    sampler: "always_on"
    
  # 尾部采样(捕获所有错误请求)
  tail:
    enabled: false
    errors_fraction: 1.0  # 100% 保留错误请求

# 资源属性
resource:
  attributes:
    service.name: "llm-gateway"
    service.version: "1.0.0"
    deployment.environment: "production"
    host.name: "${HOSTNAME}"

# Baggage 配置(跨进程传递上下文)
baggage:
  enabled: true
  # 允许的 baggage 键(安全限制)
  allowed_keys:
    - "project.id"
    - "user.id"
    - "prompt.version"

# 日志集成
logging:
  enabled: true
  level: "INFO"
  format: "json"

三、LangFuse 架构与部署

3.1 LangFuse 架构图

┌─────────────────────────────────────────────────────────────────────────────┐
│                              LangFuse 架构                                    │
├─────────────────────────────────────────────────────────────────────────────┤
│                                                                             │
│  ┌─────────────────────────────────────────────────────────────────────────┐│
│  │                           LangFuse Server                                 ││
│  │  ┌─────────────────────────────────────────────────────────────────────┐ ││
│  │  │                          Next.js Frontend                           │ ││
│  │  │  • 项目管理                                                           │ ││
│  │  │  • Trace 查看器                                                       │ ││
│  │  │  • 指标 Dashboard                                                     │ ││
│  │  │  • Prompt 管理                                                        │ ││
│  │  │  • Eval 评测                                                          │ ││
│  │  └─────────────────────────────────────────────────────────────────────┘ ││
│  │                                    │                                      ││
│  │  ┌─────────────────────────────────────────────────────────────────────┐ ││
│  │  │                          API Server (Python/FastAPI)                 │ ││
│  │  │  • Trace/span ingestion                                              │ ││
│  │  │  • Dataset management                                                 │ ││
│  │  │  • Eval execution                                                     │ ││
│  │  │  • REST API + WebSocket                                               │ ││
│  │  └─────────────────────────────────────────────────────────────────────┘ ││
│  │                                    │                                      ││
│  │  ┌─────────────────────────────────────────────────────────────────────┐ ││
│  │  │                          PostgreSQL Database                         │ ││
│  │  │  • Traces (id, project_id, name, user_id, metadata...)             │ ││
│  │  │  • Observations (spans, generation,.ChatCompletion...)              │ ││
│  │  │  • Datasets, Prompts, Evals...                                       │ ││
│  │  └─────────────────────────────────────────────────────────────────────┘ ││
│  │                                                                         ││
│  │  ┌─────────────────────────────────────────────────────────────────────┐ ││
│  │  │                          Langfuse JS SDK                             │ ││
│  │  │  • decorator @observe                                                │ ││
│  │  │  • LangchainCallbackHandler                                           │ ││
│  │  │  • Direct API calls                                                   │ ││
│  │  └─────────────────────────────────────────────────────────────────────┘ ││
│  │                                                                         ││
│  └─────────────────────────────────────────────────────────────────────────┘│
│                                    │                                          │
│    ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─│─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ ─ │
│                                    │                                            │
│  Your Application                                                                 │
│  ┌─────────────────┐  ┌─────────────────┐  ┌─────────────────┐                │
│  │   FastAPI App    │  │  LangChain App   │  │    LlamaIndex   │                │
│  │  + Python SDK    │  │  + Callback     │  │  + Callback     │                │
│  └────────┬────────┘  └────────┬────────┘  └────────┬────────┘                │
│           │                    │                    │                        │
│           └────────────────────┼────────────────────┘                        │
│                                │                                              │
│                         Langfuse Python SDK                                   │
│                         (trace, observe, callback)                            │
│                                                                                 │
└─────────────────────────────────────────────────────────────────────────────┘

3.2 Docker Compose 部署

# docker-compose.yml
# LangFuse 完整部署配置

version: '3.8'

services:
  # LangFuse 主服务
  langfuse:
    image: langfuse/langfuse:latest
    container_name: langfuse
    restart: unless-stopped
    ports:
      - "3000:3000"      # Frontend
      - "3001:3001"      # API Server
    environment:
      # 数据库配置
      DATABASE_URL: "postgresql://langfuse:langfuse_secret@postgres:5432/langfuse"
      
      # Redis 配置(会话存储)
      REDIS_URL: "redis://redis:6379"
      
      # NextAuth 配置
      NEXTAUTH_SECRET: "your-secret-key-change-in-production"
      NEXTAUTH_URL: "http://localhost:3000"
      
      # OAuth 配置(可选)
      # GOOGLE_CLIENT_ID: ""
      # GOOGLE_CLIENT_SECRET: ""
      
      # Salt key(用于数据加密)
      SALT: "your-salt-key-change-in-production"
      
      # LangFuse API Keys(创建后填入)
      # 从 LangFuse UI 的 Settings > API Keys 获取
      # 或者通过环境变量配置
      # LANGFUSE_PUBLIC_KEY: "pk-lf-..."
      # LANGFUSE_SECRET_KEY: "sk-lf-..."
      # LANGFUSE_HOST: "http://localhost:3000"
      
      # 控制台日志级别
      LOG_LEVEL: "INFO"
      
      # S3 存储配置(可选,用于存储 LLM 输入输出)
      # S3_ACCESS_KEY_ID: ""
      # S3_SECRET_ACCESS_KEY: ""
      # S3_BUCKET_NAME: "langfuse-storage"
      # S3_REGION: "us-east-1"
      
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_started
    volumes:
      # 持久化 LangFuse 数据
      - langfuse_data:/data
    networks:
      - langfuse_network
    healthcheck:
      test: ["CMD", "wget", "--no-verbose", "--tries=1", "--spider", "http://localhost:3000/api/public/health"]
      interval: 30s
      timeout: 10s
      retries: 3
      start_period: 60s

  # PostgreSQL 数据库
  postgres:
    image: postgres:15-alpine
    container_name: langfuse_postgres
    restart: unless-stopped
    environment:
      POSTGRES_DB: langfuse
      POSTGRES_USER: langfuse
      POSTGRES_PASSWORD: langfuse_secret
    volumes:
      - postgres_data:/var/lib/postgresql/data
    networks:
      - langfuse_network
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U langfuse"]
      interval: 10s
      timeout: 5s
      retries: 5

  # Redis 缓存
  redis:
    image: redis:7-alpine
    container_name: langfuse_redis
    restart: unless-stopped
    volumes:
      - redis_data:/data
    networks:
      - langfuse_network
    command: redis-server --appendonly yes

  # pgvector 扩展(用于向量存储和相似度搜索,可选)
  # 如果需要 LangFuse 的语义搜索功能
  # 注意:需要使用支持 pgvector 的 PostgreSQL 镜像
  # postgres_vector:
  #   image: pgvector/pgvector:pg15
  #   ...其他配置类似...

volumes:
  langfuse_data:
  postgres_data:
  redis_data:

networks:
  langfuse_network:
    driver: bridge

3.3 部署后初始化

# 1. 启动服务
docker-compose up -d

# 2. 查看日志确认启动成功
docker-compose logs -f langfuse

# 3. 访问 http://localhost:3000
# 首次访问会要求创建管理员账户

# 4. 创建 API Keys
# Settings > API Keys > Create new key
# 保存生成的公钥和私钥

# 5. 配置环境变量
export LANGFUSE_PUBLIC_KEY="pk-lf-..."
export LANGFUSE_SECRET_KEY="sk-lf-..."
export LANGFUSE_HOST="http://localhost:3000"

四、完整的 Python 追踪代码

4.1 OpenTelemetry + LangFuse 集成

# tracing/llm_tracing.py
"""
LLM 调用追踪完整实现
集成 OpenTelemetry + LangFuse
"""

import os
import time
import json
from typing import Optional, Any, Callable
from datetime import datetime
from contextlib import contextmanager

from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import (
    BatchSpanProcessor,
    ConsoleSpanExporter,
    SpanExporter
)
from opentelemetry.sdk.resources import Resource
from opentelemetry.propagate import set_global_textmap
from opentelemetry.trace import SpanKind, Status, StatusCode
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter

# LangFuse SDK
from langfuse import Langfuse
from langfuse.api.resources.commons.dparams import Dparams
from langfuse.api.resources.observe.params.observation_level import ObservationLevel

# LLM 语义约定
from otel_llm_semconv import LLMTraceAttributes


class LLMOtelTracer:
    """
    LLM OpenTelemetry 追踪器
    
    功能:
    1. 自动创建 Trace/Span
    2. 记录 Token 使用、延迟等指标
    3. 支持流式响应追踪
    4. 集成 LangFuse
    """
    
    def __init__(
        self,
        service_name: str,
        otlp_endpoint: Optional[str] = None,
        langfuse_config: Optional[dict] = None,
        console_export: bool = True
    ):
        self.service_name = service_name
        
        # 初始化 OpenTelemetry
        self._setup_opentelemetry(
            otlp_endpoint=otlp_endpoint,
            console_export=console_export
        )
        
        # 初始化 Langfuse
        self.langfuse = None
        if langfuse_config:
            self._setup_langfuse(langfuse_config)
        
        # 创建 tracer
        self.tracer = trace.get_tracer(service_name)
    
    def _setup_opentelemetry(
        self,
        otlp_endpoint: Optional[str],
        console_export: bool
    ):
        """配置 OpenTelemetry"""
        
        # 创建资源
        resource = Resource.create({
            "service.name": self.service_name,
            "service.version": "1.0.0",
            "deployment.environment": os.getenv("ENV", "development")
        })
        
        # 创建 Provider
        provider = TracerProvider(resource=resource)
        
        # 添加 Console 导出器(本地调试)
        if console_export:
            console_processor = BatchSpanProcessor(ConsoleSpanExporter())
            provider.add_span_processor(console_processor)
        
        # 添加 OTLP 导出器(发送到后端)
        if otlp_endpoint:
            otlp_exporter = OTLPSpanExporter(
                endpoint=otlp_endpoint,
                insecure=True
            )
            otlp_processor = BatchSpanProcessor(otlp_exporter)
            provider.add_span_processor(otlp_processor)
        
        # 设置全局 Provider
        trace.set_tracer_provider(provider)
    
    def _setup_langfuse(self, config: dict):
        """配置 Langfuse"""
        self.langfuse = Langfuse(
            public_key=config.get("public_key"),
            secret_key=config.get("secret_key"),
            host=config.get("host", "http://localhost:3000"),
            timeout=config.get("timeout", 30)
        )
    
    @contextmanager
    def trace_llm_request(
        self,
        name: str,
        model: str,
        provider: str,
        project_id: Optional[str] = None,
        user_id: Optional[str] = None,
        metadata: Optional[dict] = None
    ):
        """
        LLM 请求追踪上下文管理器
        
        用法:
            with tracer.trace_llm_request(
                name="chat_completion",
                model="gpt-4o",
                provider="openai",
                project_id="project-123",
                metadata={"session_id": "sess-456"}
            ) as span:
                span.set_input("Hello, how are you?")
                
                # 调用 LLM
                response = call_llm(...)
                
                span.set_output(response)
                span.set_token_usage(
                    input_tokens=10,
                    output_tokens=50
                )
        """
        with self.tracer.start_as_current_span(
            name=name,
            kind=SpanKind.CLIENT
        ) as otel_span:
            # Langfuse tracking
            langfuse_trace = None
            langfuse_span = None
            
            if self.langfuse:
                langfuse_trace = self.langfuse.trace(
                    name=name,
                    metadata=metadata,
                    tags=["llm", provider]
                )
                langfuse_span = langfuse_trace.span(
                    name=f"{provider}.{name}",
                    input=None,  # 先设置为 None,后面更新
                    metadata={
                        "model": model,
                        "provider": provider,
                        "project_id": project_id,
                        "user_id": user_id
                    }
                )
            
            # 创建追踪上下文
            context = LLMTraceContext(
                otel_span=otel_span,
                langfuse_trace=langfuse_trace,
                langfuse_span=langfuse_span,
                model=model,
                provider=provider,
                start_time=time.time()
            )
            
            # 设置基础属性
            otel_span.set_attribute(LLMTraceAttributes.GEN_AI_SYSTEM, "gen_ai")
            otel_span.set_attribute(LLMTraceAttributes.GEN_AI_PROVIDER, provider)
            otel_span.set_attribute(LLMTraceAttributes.GEN_AI_MODEL, model)
            
            if project_id:
                otel_span.set_attribute(LLMTraceAttributes.PROJECT_ID, project_id)
            if user_id:
                otel_span.set_attribute(LLMTraceAttributes.USER_ID, user_id)
            
            try:
                yield context
                
                # 设置成功状态
                otel_span.set_status(Status(StatusCode.OK))
                
            except Exception as e:
                # 设置错误状态
                otel_span.set_status(
                    Status(StatusCode.ERROR, str(e))
                )
                otel_span.set_attribute(
                    LLMTraceAttributes.ERROR_TYPE,
                    type(e).__name__
                )
                otel_span.record_exception(e)
                
                if langfuse_span:
                    langfuse_span.end(level="ERROR")
                    langfuse_trace.update(
                        output=str(e),
                        status_message=str(e)
                    )
                
                raise
            finally:
                # 结束追踪
                duration = time.time() - context.start_time
                otel_span.set_attribute("duration_ms", duration * 1000)
                
                if langfuse_span and langfuse_span.input is not None:
                    langfuse_span.end()


class LLMTraceContext:
    """LLM 追踪上下文"""
    
    def __init__(
        self,
        otel_span,
        langfuse_trace,
        langfuse_span,
        model: str,
        provider: str,
        start_time: float
    ):
        self.otel_span = otel_span
        self.langfuse_trace = langfuse_trace
        self.langfuse_span = langfuse_span
        self.model = model
        self.provider = provider
        self.start_time = start_time
        self._input_set = False
        self._output_set = False
    
    def set_input(self, input_data: Any):
        """设置输入"""
        # 截断过长的输入
        if isinstance(input_data, str) and len(input_data) > 10000:
            input_data = input_data[:10000] + "...[truncated]"
        
        self.otel_span.set_attribute(LLMTraceAttributes.GEN_AI_PROMPT, str(input_data))
        
        if self.langfuse_span:
            self.langfuse_span.input = input_data
        
        self._input_set = True
        return self
    
    def set_output(self, output_data: Any):
        """设置输出"""
        # 截断过长的输出
        if isinstance(output_data, str) and len(output_data) > 10000:
            output_data = output_data[:10000] + "...[truncated]"
        
        self.otel_span.set_attribute(LLMTraceAttributes.GEN_AI_COMPLETION, str(output_data))
        
        if self.langfuse_span:
            self.langfuse_span.output = output_data
        
        self._output_set = True
        return self
    
    def set_token_usage(
        self,
        input_tokens: int,
        output_tokens: int,
        total_tokens: Optional[int] = None,
        cache_hit: bool = False,
        cache_creation_tokens: int = 0,
        cache_read_tokens: int = 0,
        reasoning_tokens: int = 0
    ):
        """设置 Token 使用信息"""
        # OpenTelemetry 属性
        self.otel_span.set_attribute(
            LLMTraceAttributes.GEN_AI_USAGE_PROMPT_TOKENS,
            input_tokens
        )
        self.otel_span.set_attribute(
            LLMTraceAttributes.GEN_AI_USAGE_COMPLETION_TOKENS,
            output_tokens
        )
        self.otel_span.set_attribute(
            LLMTraceAttributes.GEN_AI_USAGE_TOTAL_TOKENS,
            total_tokens or (input_tokens + output_tokens)
        )
        
        if cache_hit:
            self.otel_span.set_attribute(
                LLMTraceAttributes.GEN_AI_RESPONSE_CACHE_HIT,
                True
            )
        
        if cache_creation_tokens > 0:
            self.otel_span.set_attribute(
                LLMTraceAttributes.GEN_AI_USAGE_CACHE_CREATION_TOKENS,
                cache_creation_tokens
            )
        
        if cache_read_tokens > 0:
            self.otel_span.set_attribute(
                LLMTraceAttributes.GEN_AI_USAGE_CACHE_READ_TOKENS,
                cache_read_tokens
            )
        
        if reasoning_tokens > 0:
            self.otel_span.set_attribute(
                LLMTraceAttributes.GEN_AI_USAGE_REASONING_TOKENS,
                reasoning_tokens
            )
        
        # LangFuse metadata
        if self.langfuse_span:
            self.langfuse_span.metadata = {
                **self.langfuse_span.metadata,
                "usage": {
                    "prompt_tokens": input_tokens,
                    "completion_tokens": output_tokens,
                    "total_tokens": total_tokens or (input_tokens + output_tokens),
                    "cache_hit": cache_hit,
                    "cache_creation_tokens": cache_creation_tokens,
                    "cache_read_tokens": cache_read_tokens,
                    "reasoning_tokens": reasoning_tokens
                }
            }
        
        return self
    
    def set_rate_limit(self, remaining: int, reset_timestamp: int):
        """设置速率限制信息"""
        self.otel_span.set_attribute(
            LLMTraceAttributes.GEN_AI_RESPONSE_RATE_LIMIT_REMAINING,
            remaining
        )
        self.otel_span.set_attribute(
            LLMTraceAttributes.GEN_AI_RESPONSE_RATE_LIMIT_RESET,
            reset_timestamp
        )
        return self
    
    def add_event(self, name: str, attributes: Optional[dict] = None):
        """添加事件(用于流式输出追踪)"""
        self.otel_span.add_event(
            name=name,
            attributes=attributes or {}
        )
        return self
    
    def record_stream_token(self, token: str, is_final: bool = False):
        """
        记录流式输出的 token
        
        用于追踪 LLM 流式响应的每个 token
        """
        if not is_final:
            # 流式输出事件
            self.add_event(
                "stream_token",
                {"token": token[:10], "truncated": len(token) > 10}
            )
        return self


def trace_streaming_llm(
    tracer: LLMOtelTracer,
    name: str,
    model: str,
    provider: str,
    input_text: str,
    stream_generator,
    project_id: Optional[str] = None
):
    """
    追踪流式 LLM 调用
    
    用法:
        response = trace_streaming_llm(
            tracer=tracer,
            name="chat_stream",
            model="gpt-4o",
            provider="openai",
            input_text="Write a story",
            stream_generator=openai_client.chat.completions.create(..., stream=True)
        )
        for chunk in response:
            print(chunk)
    """
    input_tokens = 0
    output_tokens = 0
    full_output = []
    
    with tracer.trace_llm_request(
        name=name,
        model=model,
        provider=provider,
        project_id=project_id
    ) as context:
        context.set_input(input_text)
        
        # 包装流式生成器
        def tracked_generator():
            nonlocal output_tokens
            
            for chunk in stream_generator:
                # 从 chunk 中提取内容
                content = chunk.choices[0].delta.content if hasattr(chunk, 'choices') else str(chunk)
                
                if content:
                    full_output.append(content)
                    output_tokens += estimate_token_count(content)
                    
                    # 记录每个 token
                    context.record_stream_token(content)
                
                yield chunk
            
            # 更新最终的 output
            context.set_output("".join(full_output))
            
            # 从最后一个 chunk 获取 usage 信息
            if hasattr(chunk, 'usage') and chunk.usage:
                context.set_token_usage(
                    input_tokens=chunk.usage.prompt_tokens or input_tokens,
                    output_tokens=chunk.usage.completion_tokens or output_tokens,
                    total_tokens=chunk.usage.total_tokens
                )
            else:
                context.set_token_usage(
                    input_tokens=input_tokens,
                    output_tokens=output_tokens
                )
        
        return tracked_generator()


def estimate_token_count(text: str) -> int:
    """
    估算 token 数量
    
    粗略估算:中文约 1.5-2 tokens/字,英文约 0.25 tokens/字符
    精确计算需要使用 tiktoken 等库
    """
    # 简单估算
    chinese_chars = sum(1 for c in text if '\u4e00' <= c <= '\u9fff')
    other_chars = len(text) - chinese_chars
    
    return int(chinese_chars * 1.5 + other_chars * 0.25)

4.2 FastAPI 集成

# api/llm_gateway.py
"""
FastAPI LLM 网关 - 集成追踪
"""

import os
from typing import Optional, List, Union
from fastapi import FastAPI, Request, HTTPException, Depends
from fastapi.responses import StreamingResponse
from pydantic import BaseModel
import openai

from tracing.llm_tracing import LLMOtelTracer, trace_streaming_llm


# 初始化追踪器
tracer = LLMOtelTracer(
    service_name="llm-gateway",
    otlp_endpoint=os.getenv("OTLP_ENDPOINT", "http://localhost:4317"),
    langfuse_config={
        "public_key": os.getenv("LANGFUSE_PUBLIC_KEY"),
        "secret_key": os.getenv("LANGFUSE_SECRET_KEY"),
        "host": os.getenv("LANGFUSE_HOST", "http://localhost:3000")
    },
    console_export=os.getenv("ENV") != "production"
)

# OpenAI 客户端
openai_client = openai.OpenAI(
    api_key=os.getenv("OPENAI_API_KEY")
)

app = FastAPI(title="LLM Gateway with Tracing")


class ChatRequest(BaseModel):
    """聊天请求"""
    model: str = "gpt-4o"
    messages: List[dict]
    temperature: float = 0.7
    max_tokens: Optional[int] = None
    stream: bool = False
    project_id: Optional[str] = None
    user_id: Optional[str] = None


class ChatResponse(BaseModel):
    """聊天响应"""
    id: str
    model: str
    content: str
    input_tokens: int
    output_tokens: int
    total_tokens: int
    latency_ms: float


@app.post("/v1/chat/completions")
async def chat_completions(
    request: ChatRequest,
    http_request: Request
):
    """
    Chat Completions API - 带完整追踪
    """
    # 获取项目/用户 ID
    project_id = request.project_id or http_request.headers.get("X-Project-ID")
    user_id = request.user_id or http_request.headers.get("X-User-ID")
    
    # 合并消息为单个 prompt(用于追踪)
    prompt_text = "\n".join([
        f"{msg.get('role', 'user')}: {msg.get('content', '')}"
        for msg in request.messages
    ])
    
    # 非流式响应
    if not request.stream:
        with tracer.trace_llm_request(
            name="chat.completion",
            model=request.model,
            provider="openai",
            project_id=project_id,
            user_id=user_id,
            metadata={
                "endpoint": "/v1/chat/completions",
                "temperature": request.temperature,
                "max_tokens": request.max_tokens
            }
        ) as context:
            context.set_input(prompt_text)
            
            try:
                # 调用 OpenAI API
                response = openai_client.chat.completions.create(
                    model=request.model,
                    messages=request.messages,
                    temperature=request.temperature,
                    max_tokens=request.max_tokens
                )
                
                # 提取响应内容
                content = response.choices[0].message.content
                context.set_output(content)
                
                # 设置 Token 使用
                if response.usage:
                    context.set_token_usage(
                        input_tokens=response.usage.prompt_tokens,
                        output_tokens=response.usage.completion_tokens,
                        total_tokens=response.usage.total_tokens
                    )
                    
                    # 速率限制
                    if hasattr(response, 'x_ratelimit_remaining'):
                        context.set_rate_limit(
                            remaining=response.x_ratelimit_remaining,
                            reset_timestamp=int(response.x_ratelimit_reset)
                        )
                
                return ChatResponse(
                    id=response.id,
                    model=response.model,
                    content=content,
                    input_tokens=response.usage.prompt_tokens,
                    output_tokens=response.usage.completion_tokens,
                    total_tokens=response.usage.total_tokens,
                    latency_ms=response.response_ms if hasattr(response, 'response_ms') else 0
                )
                
            except openai.RateLimitError as e:
                raise HTTPException(status_code=429, detail="Rate limit exceeded")
            except openai.APIError as e:
                raise HTTPException(status_code=500, detail=str(e))
    
    # 流式响应
    else:
        # 使用追踪生成器
        raw_stream = openai_client.chat.completions.create(
            model=request.model,
            messages=request.messages,
            temperature=request.temperature,
            max_tokens=request.max_tokens,
            stream=True
        )
        
        tracked_stream = trace_streaming_llm(
            tracer=tracer,
            name="chat.completion.stream",
            model=request.model,
            provider="openai",
            input_text=prompt_text,
            stream_generator=raw_stream,
            project_id=project_id
        )
        
        return StreamingResponse(
            (chunk.model_dump_json() for chunk in tracked_stream),
            media_type="application/x-ndjson"
        )


@app.get("/health")
async def health():
    """健康检查"""
    return {"status": "healthy", "tracer": "active"}


@app.get("/v1/traces")
async def list_traces(
    project_id: Optional[str] = None,
    limit: int = 100
):
    """
    列出最近的追踪记录
    通过 Langfuse API
    """
    if not tracer.langfuse:
        raise HTTPException(status_code=503, detail="Langfuse not configured")
    
    try:
        traces = tracer.langfuse.fetch_traces(
            limit=limit,
            tag="llm"
        )
        return {"traces": traces}
    except Exception as e:
        raise HTTPException(status_code=500, detail=str(e))


# 中间件:添加请求追踪
@app.middleware("http")
async def add_trace_context(request: Request, call_next):
    """中间件:为每个请求添加追踪上下文"""
    
    # 从 header 获取 trace context(如果来自上游)
    traceparent = request.headers.get("traceparent")
    if traceparent:
        from opentelemetry.propagate import extract
        context = extract(traceparent)
        # 使用 context 开始追踪
    
    response = await call_next(request)
    
    # 添加 trace ID 到响应头
    current_span = trace.get_current_span()
    if current_span:
        span_context = current_span.get_span_context()
        if span_context.is_valid:
            trace_id = format(span_context.trace_id, '032x')
            response.headers["X-Trace-ID"] = trace_id
    
    return response

4.3 LangChain 集成

# integrations/langchain_tracing.py
"""
LangChain + OpenTelemetry + LangFuse 集成
"""

from typing import Any, Dict, List, Optional
from datetime import datetime

from langchain.callbacks.base import BaseCallbackHandler
from langchain.schema import AgentAction, AgentFinish, LLMResult
from langchain_openai import ChatOpenAI

from opentelemetry import trace
from opentelemetry.trace import SpanKind, Status, StatusCode
from opentelemetry.propagate import inject, extract

from langfuse.callback import CallbackHandler as LangfuseCallback

from tracing.llm_tracing import LLMOtelTracer
from otel_llm_semconv import LLMTraceAttributes


class CombinedCallbackHandler(BaseCallbackHandler):
    """
    组合回调处理器
    
    同时支持:
    1. OpenTelemetry 追踪
    2. LangFuse 追踪
    3. 自定义指标记录
    """
    
    def __init__(
        self,
        tracer: LLMOtelTracer,
        langfuse_handler: Optional[LangfuseCallback] = None,
        project_name: str = "langchain-app"
    ):
        self.tracer = tracer
        self.langfuse_handler = langfuse_handler
        self.project_name = project_name
        
        # 追踪上下文栈
        self._trace_stack: list = []
    
    def on_llm_start(
        self,
        serialized: Dict[str, Any],
        prompts: List[str],
        **kwargs: Any
    ) -> None:
        """LLM 调用开始"""
        # 从 kwargs 获取模型信息
        model_name = serialized.get("name", "unknown")
        
        # 获取 parent span(如果有)
        parent_span = trace.get_current_span()
        
        # 创建新 span
        span_name = f"llm.{model_name}"
        
        with self.tracer.tracer.start_span(
            span_name,
            kind=SpanKind.CLIENT,
            parent=parent_span
        ) as span:
            span.set_attribute(LLMTraceAttributes.GEN_AI_SYSTEM, "gen_ai")
            span.set_attribute(LLMTraceAttributes.GEN_AI_MODEL, model_name)
            
            # 设置输入
            for i, prompt in enumerate(prompts):
                span.set_attribute(f"llm.prompt.{i}", prompt[:1000])
            
            # 记录到栈
            self._trace_stack.append({
                "span": span,
                "start_time": datetime.now(),
                "model": model_name,
                "prompts": prompts
            })
        
        # Langfuse
        if self.langfuse_handler:
            self.langfuse_handler.handle_llm_start(
                serialized, prompts, **kwargs
            )
    
    def on_llm_end(
        self,
        response: LLMResult,
        **kwargs: Any
    ) -> None:
        """LLM 调用结束"""
        if not self._trace_stack:
            return
        
        trace_info = self._trace_stack.pop()
        span = trace_info["span"]
        
        try:
            # 处理输出
            if response.generations:
                for i, generation_list in enumerate(response.generations):
                    for j, generation in enumerate(generation_list):
                        text = generation.text
                        span.set_attribute(
                            f"llm.completion.{i}.{j}",
                            text[:1000]  # 截断
                        )
            
            # 处理 Token 使用
            if response.llm_output and "token_usage" in response.llm_output:
                usage = response.llm_output["token_usage"]
                span.set_attribute(
                    LLMTraceAttributes.GEN_AI_USAGE_PROMPT_TOKENS,
                    usage.get("prompt_tokens", 0)
                )
                span.set_attribute(
                    LLMTraceAttributes.GEN_AI_USAGE_COMPLETION_TOKENS,
                    usage.get("completion_tokens", 0)
                )
                span.set_attribute(
                    LLMTraceAttributes.GEN_AI_USAGE_TOTAL_TOKENS,
                    usage.get("total_tokens", 0)
                )
            
            span.set_status(Status(StatusCode.OK))
            
        except Exception as e:
            span.set_status(Status(StatusCode.ERROR, str(e)))
            span.record_exception(e)
        finally:
            span.end()
        
        # Langfuse
        if self.langfuse_handler:
            self.langfuse_handler.handle_llm_end(response, **kwargs)
    
    def on_llm_error(
        self,
        error: Exception,
        **kwargs: Any
    ) -> None:
        """LLM 调用错误"""
        if not self._trace_stack:
            return
        
        trace_info = self._trace_stack[-1]
        span = trace_info["span"]
        
        span.set_status(Status(StatusCode.ERROR, str(error)))
        span.record_exception(error)
        
        self._trace_stack.pop()
        span.end()
        
        # Langfuse
        if self.langfuse_handler:
            self.langfuse_handler.handle_llm_error(error, **kwargs)
    
    def on_chain_start(
        self,
        serialized: Dict[str, Any],
        inputs: Dict[str, Any],
        **kwargs: Any
    ) -> None:
        """Chain 调用开始"""
        chain_name = serialized.get("name", serialized.get("id", ["unknown"])[-1])
        
        span = self.tracer.tracer.start_span(
            f"chain.{chain_name}",
            kind=SpanKind.INTERNAL
        )
        
        span.set_attribute("chain.name", chain_name)
        
        # 注入 trace context 到 baggage
        inject()
        
        self._trace_stack.append({
            "span": span,
            "chain": chain_name
        })
        
        # Langfuse
        if self.langfuse_handler:
            self.langfuse_handler.handle_chain_start(
                serialized, inputs, **kwargs
            )
    
    def on_chain_end(
        self,
        outputs: Dict[str, Any],
        **kwargs: Any
    ) -> None:
        """Chain 调用结束"""
        if not self._trace_stack:
            return
        
        trace_info = self._trace_stack.pop()
        span = trace_info["span"]
        
        span.set_attribute("chain.output_keys", list(outputs.keys()))
        span.set_status(Status(StatusCode.OK))
        span.end()
        
        # Langfuse
        if self.langfuse_handler:
            self.langfuse_handler.handle_chain_end(outputs, **kwargs)
    
    def on_tool_start(
        self,
        serialized: Dict[str, Any],
        input_str: str,
        **kwargs: Any
    ) -> None:
        """Tool 调用开始"""
        tool_name = serialized.get("name", "unknown")
        
        span = self.tracer.tracer.start_span(
            f"tool.{tool_name}",
            kind=SpanKind.INTERNAL
        )
        
        span.set_attribute("tool.name", tool_name)
        span.set_attribute("tool.input", input_str[:500])
        
        self._trace_stack.append({"span": span, "tool": tool_name})
    
    def on_tool_end(
        self,
        output: str,
        **kwargs: Any
    ) -> None:
        """Tool 调用结束"""
        if not self._trace_stack:
            return
        
        trace_info = self._trace_stack.pop()
        span = trace_info["span"]
        
        span.set_attribute("tool.output", output[:500])
        span.set_status(Status(StatusCode.OK))
        span.end()


# 使用示例
def create_traced_llm(
    model_name: str = "gpt-4o",
    tracing_config: dict = None
):
    """
    创建带追踪的 LLM 实例
    
    用法:
        llm = create_traced_llm("gpt-4o")
        response = llm.invoke("Hello")
    """
    tracer = tracing_config.get("tracer") if tracing_config else None
    
    callbacks = []
    
    if tracer:
        callbacks.append(
            CombinedCallbackHandler(
                tracer=tracer,
                langfuse_handler=LangfuseCallback(
                    user_id=tracing_config.get("user_id"),
                    project_name=tracing_config.get("project_name", "langchain-app")
                )
            )
        )
    
    return ChatOpenAI(
        model=model_name,
        callbacks=callbacks
    )

五、Trace 数据排障案例

5.1 案例一:Token 消耗异常

问题描述
用户报告某天的 Token 消耗异常增长,日均从 100 万增长到 300 万。

排查步骤

  1. 定位时间点

LangFuse Dashboard → Traces → 筛选日期范围
  1. 查看 Token 分布

// 查询请求分布
SELECT 
  DATE(created_at) as date,
  COUNT(*) as request_count,
  AVG(usage->>'total_tokens') as avg_tokens,
  SUM((usage->>'total_tokens')::int) as total_tokens
FROM observations
WHERE type = 'GENERATION'
  AND created_at BETWEEN '2024-01-01' AND '2024-01-07'
GROUP BY DATE(created_at)
ORDER BY date
  1. 定位异常请求

# 找出 Token 消耗最高的请求
traces = langfuse.fetch_traces(
    limit=100,
    filter=[
        ("usage.total_tokens", ">=", 50000)  # 超过 50k tokens
    ],
    order_by=[("timestamp", "desc")]
)

for trace in traces:
    print(f"""
    Trace ID: {trace.id}
    Model: {trace.observations[0].model}
    Input Tokens: {trace.observations[0].usage.prompt_tokens}
    Output Tokens: {trace.observations[0].usage.completion_tokens}
    """)
  1. 根因分析

  • 发现是某次 Prompt 更新后,未清理的 Few-shot 示例积累

  • 每次请求都携带了过多历史对话

5.2 案例二:响应延迟过高

问题描述
用户反馈某些请求响应时间超过 30 秒。

排查步骤

  1. 查看延迟分布

// 使用 OpenTelemetry 查询
spans = otel_client.query_spans(
    service_name="llm-gateway",
    span_name="llm.chat.completion",
    start_time=start_date,
    end_time=end_date,
    limit=100
)

# 计算 P99 延迟
latencies = [s.duration_ms for s in spans]
latencies.sort()
p99 = latencies[int(len(latencies) * 0.99)]
  1. 分析慢请求

Trace ID: abc123
├── Preprocess: 50ms
├── LLM Call: 28500ms ← 主要耗时
│   ├── TTFT (Time to First Token): 2100ms
│   ├── Token Generation: 26000ms
│   └── Model: gpt-4o
└── Postprocess: 100ms

// 结论:模型生成速度 = 523 tokens / 26s ≈ 20 tokens/s
// 低于正常水平,可能是请求量大导致排队
  1. 关联指标

  • 发现同时段请求量激增 5 倍

  • 确认是 Rate Limit 导致排队等待

5.3 案例三:缓存命中率下降

问题描述
缓存命中率从 60% 下降到 20%。

排查步骤

  1. 查看缓存命中趋势

# 缓存命中率
rate(llm_tokens_cached_total[1h]) / 
(rate(llm_tokens_input_total[1h]) + 0.001)
  1. 分析未命中原因

// 查询未命中的请求
traces = langfuse.fetch_traces(
    filter=[
        ("cache_hit", "=", False),
        ("metadata.cache_key", "exists", True)
    ]
)

for trace in traces:
    print(f"""
    Cache Key: {trace.metadata.cache_key}
    Prompt Hash: {trace.metadata.prompt_hash}
    """)
  1. 根因

  • 发现模型从 gpt-3.5-turbo-1106 升级到 gpt-3.5-turbo-0125

  • 不同模型版本的 Tokenizer 不同,导致 Hash 不匹配


六、Docker Compose 完整部署

# docker-compose.otel.yml
# OpenTelemetry + LangFuse 完整部署

version: '3.8'

services:
  # LangFuse
  langfuse:
    image: langfuse/langfuse:latest
    container_name: langfuse
    restart: unless-stopped
    ports:
      - "3000:3000"
      - "3001:3001"
    environment:
      DATABASE_URL: postgresql://langfuse:langfuse_secret@postgres:5432/langfuse
      REDIS_URL: redis://redis:6379
      NEXTAUTH_SECRET: your-secret-key
      NEXTAUTH_URL: http://localhost:3000
      SALT: your-salt-key
    depends_on:
      postgres:
        condition: service_healthy
      redis:
        condition: service_started
    networks:
      - observability

  # PostgreSQL
  postgres:
    image: postgres:15-alpine
    container_name: otel_postgres
    restart: unless-stopped
    environment:
      POSTGRES_DB: langfuse
      POSTGRES_USER: langfuse
      POSTGRES_PASSWORD: langfuse_secret
    volumes:
      - postgres_data:/var/lib/postgresql/data
    networks:
      - observability
    healthcheck:
      test: ["CMD-SHELL", "pg_isready -U langfuse"]
      interval: 10s
      timeout: 5s
      retries: 5

  # Redis
  redis:
    image: redis:7-alpine
    container_name: otel_redis
    restart: unless-stopped
    volumes:
      - redis_data:/data
    networks:
      - observability

  # OpenTelemetry Collector
  otel-collector:
    image: otel/opentelemetry-collector-contrib:latest
    container_name: otel_collector
    restart: unless-stopped
    command: ["--config=/etc/otel-collector-config.yaml"]
    volumes:
      - ./config/otel-collector-config.yaml:/etc/otel-collector-config.yaml
    ports:
      - "4317:4317"      # OTLP gRPC
      - "4318:4318"      # OTLP HTTP
      - "8888:8888"      # Prometheus metrics
      - "8889:8889"      # PrometheusExporter metrics
    networks:
      - observability
    depends_on:
      - jaeger
      - prometheus

  # Jaeger (Trace Backend)
  jaeger:
    image: jaegertracing/all-in-one:latest
    container_name: jaeger
    restart: unless-stopped
    environment:
      COLLECTOR_OTLP_ENABLED: true
      SPAN_STORAGE_TYPE: badger
      BADGER_EPHEMERAL: true
    ports:
      - "16686:16686"    # UI
      - "14268:14268"    # Jaeger Collector
    networks:
      - observability

  # Prometheus
  prometheus:
    image: prom/prometheus:latest
    container_name: prometheus
    restart: unless-stopped
    command:
      - '--config.file=/etc/prometheus/prometheus.yml'
      - '--storage.tsdb.path=/prometheus'
      - '--web.enable-lifecycle'
    volumes:
      - ./config/prometheus.yml:/etc/prometheus/prometheus.yml
      - prometheus_data:/prometheus
    ports:
      - "9090:9090"
    networks:
      - observability

  # Grafana
  grafana:
    image: grafana/grafana:latest
    container_name: grafana
    restart: unless-stopped
    environment:
      GF_SECURITY_ADMIN_PASSWORD: admin
      GF_USERS_ALLOW_SIGN_UP: false
    volumes:
      - grafana_data:/var/lib/grafana
      - ./config/grafana/provisioning:/etc/grafana/provisioning
    ports:
      - "3002:3000"
    depends_on:
      - prometheus
    networks:
      - observability

volumes:
  postgres_data:
  redis_data:
  prometheus_data:
  grafana_data:

networks:
  observability:
    driver: bridge
# config/otel-collector-config.yaml
# OpenTelemetry Collector 配置

receivers:
  otlp:
    protocols:
      grpc:
        endpoint: 0.0.0.0:4317
      http:
        endpoint: 0.0.0.0:4318

processors:
  batch:
    timeout: 1s
    send_batch_size: 1024

  memory_limiter:
    check_interval: 1s
    limit_mib: 512

  # 过滤敏感信息
  transform:
    trace_statements:
      - context: span
        statements:
          - replace_pattern(attributes["gen_ai.prompt"], "password=[^&]*", "password=***")

exporters:
  # 发送到 Jaeger
  jaeger:
    endpoint: jaeger:14250
    tls:
      insecure: true

  # 发送到 Prometheus
  prometheus:
    endpoint: "0.0.0.0:8889"
    namespace: llm
    const_labels:
      service: llm-gateway

  # 发送到 LangFuse(通过 OTLP)
  # 注意:LangFuse 也支持 OTLP 接收
  otlp/langfuse:
    endpoint: langfuse:3001
    tls:
      insecure: true

service:
  pipelines:
    traces:
      receivers: [otlp]
      processors: [memory_limiter, batch]
      exporters: [jaeger]
      
    metrics:
      receivers: [otlp]
      processors: [memory_limiter, batch]
      exporters: [prometheus]

七、总结

7.1 核心要点回顾

  1. LLM 追踪的特殊性

  • Token 分布是核心指标

  • 流式响应需要特殊处理

  • 缓存命中是独特的维度

  1. OpenTelemetry 集成

  • 遵循 LLM 语义约定

  • 支持多种导出器组合

  • 可与现有 APM 工具集成

  1. LangFuse 优势

  • 开箱即用的 LLM 追踪 Dashboard

  • 支持 Prompt 版本管理

  • 内置 Eval 功能

  1. 最佳实践

  • 生产环境使用尾部采样

  • 控制存储成本(数据保留策略)

  • 敏感信息脱敏

7.2 扩展方向

  • 自动化告警:Token 消耗异常、延迟超标时自动告警

  • A/B 测试追踪:对比不同 Prompt/模型的指标

  • 成本归因:将成本精确归因到用户/项目/功能

  • 回归测试:每次模型更新跑回归测试,对比指标变化

7.3 注意事项

  1. 数据隐私:Prompt/Response 可能包含敏感信息,需要脱敏

  2. 存储成本:高流量场景下追踪数据量很大,需要规划存储

  3. 采样策略:生产环境建议使用尾部采样,只保留错误/慢请求

  4. 性能影响:追踪本身有开销,流式场景需要注意


本文详细介绍了 OpenTelemetry + LangFuse 的完整集成方案,代码可直接用于生产环境。建议结合 LangFuse 官方文档和 OpenTelemetry 规范进行调整。

0
  1. 支付宝打赏

    qrcode alipay
  2. 微信打赏

    qrcode weixin

评论区