Skip to content

阶段 4:OpenTelemetry Logs 深入

一、Logs 在可观测性中的定位

1.1 为什么 Logs 是最后加入 OpenTelemetry 的?

时间线:
  2019  OpenTelemetry 成立,最先支持 Traces
  2020  Metrics 支持加入
  2021  Logs 规范开始制定(最晚)
  2023  Logs SDK 趋于稳定(Python 中仍有 _logs 下划线前缀,表明 API 可能变化)

原因:
  - Traces 和 Metrics 是全新的概念,需要统一标准
  - Logs 已经有成熟的生态(Python logging、Log4j、Serilog 等)
  - OTel 的策略不是"取代"现有日志系统,而是"桥接"它们

1.2 OTel Logs 的设计哲学

不替代,而是增强

传统日志:
  应用 → logging.info("用户登录") → 文件/stdout → ELK/Loki

OTel 增强后的日志:
  应用 → logging.info("用户登录") → OTel LoggingHandler → OTel SDK
                                          │                    │
                                     自动注入              统一导出
                                  trace_id/span_id         OTLP 协议
                                          │                    │
                                          ▼                    ▼
                                    日志自动关联 Trace      发送到后端
                                    在 Grafana 中点击      (Loki/ELK等)
                                    Trace 可直接跳转
                                    到相关日志

1.3 Traces、Metrics、Logs 三者的协作

场景:线上接口突然变慢

Step 1 → Metrics 告警
         "order-service 的 P99 延迟从 200ms 飙到 3s"

Step 2 → Traces 定位
         找到一条慢请求的 Trace:
         invocation (3.2s)
         └── call-llm (3.1s)  ← 瓶颈在这里

Step 3 → Logs 查看细节
         根据 trace_id 过滤日志:
         [trace_id=abc123] "LLM API 返回 429 Too Many Requests,正在重试..."
         [trace_id=abc123] "第 3 次重试成功,耗时 2.8s"
         
         根因找到:LLM API 限流导致重试,耗时变长

三者的本质区别

维度TracesMetricsLogs
数据模型结构化的 Span 树数值型时序数据半结构化文本记录
典型数据量中(可采样)低(聚合后)高(每条都记录)
查询方式按 trace_id 查整条链路按维度聚合(SUM/AVG/P99)全文搜索 + 过滤
回答的问题请求经过了哪些步骤?系统现在怎么样?趋势如何?具体发生了什么?
采样常见(1%-100%)不采样(全量聚合)通常全量,或按级别过滤

二、核心概念

2.1 OTel Logs 数据模型

一条 OTel LogRecord 包含以下字段:

┌──────────────────────────────────────────────────────┐
│                   LogRecord                          │
├──────────────────────────────────────────────────────┤
│ Timestamp         : 2025-03-15T10:23:45.123456Z      │  ← 日志产生时间
│ ObservedTimestamp  : 2025-03-15T10:23:45.123789Z      │  ← OTel 观察到的时间
│ SeverityNumber    : 9 (INFO)                         │  ← 日志级别(数值)
│ SeverityText      : "INFO"                           │  ← 日志级别(文本)
│ Body              : "User 12345 logged in"           │  ← 日志正文
│ Attributes        : {user.id: "12345", ip: "1.2.3.4"}│  ← 附加属性
│ Resource          : {service.name: "auth-service"}   │  ← 来源服务
│ TraceId           : abc123def456...                   │  ← 关联的 Trace ID
│ SpanId            : 789ghi...                        │  ← 关联的 Span ID
│ TraceFlags        : 01                               │  ← Trace 标志
│ InstrumentationScope: {name: "myapp.auth", ver: "1.0"}│ ← 产生日志的模块
└──────────────────────────────────────────────────────┘

2.2 日志级别映射

OTel 定义了标准化的 SeverityNumber,与各语言的日志级别对应:

OTel SeverityNumberOTel SeverityTextPython logging含义
1-4TRACE最详细的调试信息
5-8DEBUGlogging.DEBUG (10)调试信息
9-12INFOlogging.INFO (20)一般信息
13-16WARNlogging.WARNING (30)警告
17-20ERRORlogging.ERROR (40)错误
21-24FATALlogging.CRITICAL (50)致命错误

2.3 组件架构

┌──────────────────────────────────────────────────────┐
│                    你的应用                            │
│                                                       │
│  logging.info("msg")                                  │
│       │                                               │
│       ▼                                               │
│  ┌──────────────┐                                     │
│  │ Python Logger│ ← 标准 Python logging               │
│  └──────┬───────┘                                     │
│         │                                             │
│         ▼                                             │
│  ┌──────────────────┐                                 │
│  │ OTel LoggingHandler │ ← 桥接层(将 Python 日志转为 OTel 格式)│
│  └──────┬───────────┘                                 │
│         │                                             │
│         ▼                                             │
│  ┌──────────────────┐                                 │
│  │  LoggerProvider   │ ← OTel SDK 核心                │
│  │    │              │                                │
│  │    ▼              │                                │
│  │  OTel Logger      │ ← 创建 LogRecord               │
│  │    │              │                                │
│  │    ▼              │                                │
│  │  LogRecordProcessor│ ← 处理管道                     │
│  │  (Simple / Batch) │                                │
│  │    │              │                                │
│  │    ▼              │                                │
│  │  LogExporter      │ ← 导出到后端                    │
│  │  (Console/OTLP)   │                                │
│  └──────────────────┘                                 │
└──────────────────────────────────────────────────────┘

关键组件说明:

组件职责Python 类
LoggerProvider管理 Logger 的配置(Resource、Processor)opentelemetry.sdk._logs.LoggerProvider
Logger创建 LogRecord由 LoggerProvider 内部创建
LoggingHandlerPython logging → OTel 的桥接器opentelemetry.sdk._logs.LoggingHandler
LogRecordProcessor处理 LogRecord 的管道SimpleLogRecordProcessor / BatchLogRecordProcessor
LogExporter将 LogRecord 发送到后端ConsoleLogExporter / OTLPLogExporter

三、基础用法

3.1 安装依赖

bash
# 核心日志包
pip install opentelemetry-api opentelemetry-sdk

# OTLP 日志导出器
pip install opentelemetry-exporter-otlp-proto-grpc
# 或 HTTP 版本
pip install opentelemetry-exporter-otlp-proto-http

3.2 最小可运行示例

python
"""
最简 OTel Logs 示例:将 Python logging 桥接到 OTel
"""
import logging
from opentelemetry._logs import set_logger_provider
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import (
    SimpleLogRecordProcessor,
    ConsoleLogExporter,
)
from opentelemetry.sdk.resources import Resource

# 1. 创建 Resource
resource = Resource.create({
    "service.name": "my-app",
    "service.version": "1.0.0",
})

# 2. 创建 LoggerProvider
logger_provider = LoggerProvider(resource=resource)

# 3. 添加 Processor + Exporter
logger_provider.add_log_record_processor(
    SimpleLogRecordProcessor(ConsoleLogExporter())
)

# 4. 设置全局 LoggerProvider
set_logger_provider(logger_provider)

# 5. 创建 OTel LoggingHandler 并挂载到 Python logger
handler = LoggingHandler(
    level=logging.NOTSET,
    logger_provider=logger_provider,
)

# 6. 配置 Python logging
root_logger = logging.getLogger()
root_logger.addHandler(handler)
root_logger.setLevel(logging.INFO)

# 7. 正常使用 Python logging
logger = logging.getLogger("myapp.orders")
logger.info("订单创建成功")
logger.warning("库存不足", extra={"product.id": "SKU-001", "stock": 3})
logger.error("支付失败", extra={"order.id": "ORD-12345", "error.code": "INSUFFICIENT_FUNDS"})

# 8. 关闭
logger_provider.shutdown()

输出的 LogRecord(JSON 格式):

json
{
    "body": "订单创建成功",
    "severity_number": 9,
    "severity_text": "INFO",
    "attributes": {},
    "timestamp": "2025-03-15T10:23:45.123456Z",
    "trace_id": "0x00000000000000000000000000000000",
    "span_id": "0x0000000000000000",
    "resource": {
        "service.name": "my-app",
        "service.version": "1.0.0"
    },
    "instrumentation_scope": {
        "name": "myapp.orders"
    }
}

注意此时 trace_idspan_id 全是 0,因为我们还没有创建 Trace。


四、核心能力:日志与 Trace 的关联

这是 OTel Logs 最核心的价值——自动将日志与 Trace 关联

4.1 原理

当你同时配置了 TracerProvider 和 LoggerProvider 时:

  with tracer.start_as_current_span("handle-request"):
      # 此时 Context 中有当前 Span 的信息
      # ↓ OTel LoggingHandler 会自动读取 Context
      logger.info("处理请求中...")
      # ↓ 导出的 LogRecord 自动携带 trace_id 和 span_id

工作机制:
  1. tracer.start_as_current_span() 将 Span 放入 Context
  2. logging.info() 被 OTel LoggingHandler 拦截
  3. LoggingHandler 从 Context 读取当前 Span 的 trace_id 和 span_id
  4. 将 trace_id/span_id 写入 LogRecord
  5. 导出时 LogRecord 自动关联到对应的 Trace

4.2 完整示例

python
"""
demo_log_trace_correlation.py
展示日志与 Trace 的自动关联
"""
import logging
import time

from opentelemetry import trace
from opentelemetry._logs import set_logger_provider
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor, ConsoleSpanExporter
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import SimpleLogRecordProcessor, ConsoleLogExporter
from opentelemetry.sdk.resources import Resource

# ========== 统一 Resource ==========
resource = Resource.create({
    "service.name": "order-service",
    "service.version": "2.0.0",
    "deployment.environment": "staging",
})

# ========== 配置 Traces ==========
trace_provider = TracerProvider(resource=resource)
trace_provider.add_span_processor(SimpleSpanProcessor(ConsoleSpanExporter()))
trace.set_tracer_provider(trace_provider)
tracer = trace.get_tracer("order-service")

# ========== 配置 Logs ==========
logger_provider = LoggerProvider(resource=resource)
logger_provider.add_log_record_processor(
    SimpleLogRecordProcessor(ConsoleLogExporter())
)
set_logger_provider(logger_provider)

otel_handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider)
logging.getLogger().addHandler(otel_handler)
logging.getLogger().setLevel(logging.DEBUG)

logger = logging.getLogger("order-service.handler")

# ========== 业务逻辑 ==========

def create_order(user_id: str, product_id: str, quantity: int):
    with tracer.start_as_current_span("create_order") as span:
        span.set_attribute("user.id", user_id)
        span.set_attribute("product.id", product_id)

        # 这些日志会自动携带 trace_id 和 span_id
        logger.info("开始创建订单", extra={
            "user.id": user_id,
            "product.id": product_id,
            "quantity": quantity,
        })

        with tracer.start_as_current_span("check_inventory") as inv_span:
            logger.debug("检查库存: product=%s, quantity=%d", product_id, quantity)
            time.sleep(0.02)

            if quantity > 100:
                logger.warning("库存紧张: product=%s, 剩余不多", product_id)

        with tracer.start_as_current_span("process_payment") as pay_span:
            logger.info("处理支付: user=%s, amount=%.2f", user_id, quantity * 29.9)
            time.sleep(0.03)

            try:
                if product_id == "INVALID":
                    raise ValueError("无效的商品ID")
            except ValueError as e:
                logger.error("支付失败: %s", str(e), exc_info=True, extra={
                    "error.type": type(e).__name__,
                })
                span.set_status(trace.StatusCode.ERROR, str(e))
                span.record_exception(e)
                raise

        logger.info("订单创建成功: order_id=ORD-%s", "20250315001")

# 运行
create_order("user-001", "SKU-100", 5)
print("\n" + "=" * 60)
print("观察上面的输出:")
print("1. 每条 LogRecord 的 trace_id 和 span_id 都不为 0")
print("2. 同一个 Trace 下的所有日志共享同一个 trace_id")
print("3. 不同 Span 中的日志有不同的 span_id")
print("=" * 60)

trace_provider.shutdown()
logger_provider.shutdown()

4.3 在 Grafana 中的效果

配置 Loki + Tempo 后,可以实现:

┌──────────────── Grafana ─────────────────┐
│                                           │
│  Traces 视图 (Tempo):                     │
│  ┌─ create_order (85ms) ──────────────┐  │
│  │  ├─ check_inventory (20ms)         │  │
│  │  └─ process_payment (30ms)         │  │
│  └────────────────────────────────────┘  │
│            │                              │
│            │ 点击 "View Logs"             │
│            ▼                              │
│  Logs 视图 (Loki):                        │
│  10:23:45.100 [INFO]  开始创建订单         │
│  10:23:45.110 [DEBUG] 检查库存: SKU-100    │
│  10:23:45.130 [INFO]  处理支付: 149.50     │
│  10:23:45.160 [INFO]  订单创建成功          │
│                                           │
│  ← 反向:在日志中点击 trace_id 跳转到 Trace │
└───────────────────────────────────────────┘

五、进阶用法

5.1 自定义 LogRecord 属性

python
"""
通过 extra 字段添加自定义属性到 LogRecord
"""
import logging

logger = logging.getLogger("myapp")

# 方式一:通过 extra 参数(推荐)
logger.info("用户登录", extra={
    "user.id": "12345",
    "user.role": "admin",
    "login.method": "sso",
    "login.ip": "10.0.1.100",
})
# 这些 extra 字段会自动成为 LogRecord 的 Attributes

# 方式二:使用 LogRecord 工厂(全局添加字段)
old_factory = logging.getLogRecordFactory()

def custom_factory(*args, **kwargs):
    record = old_factory(*args, **kwargs)
    record.service_name = "order-service"
    record.environment = "production"
    return record

logging.setLogRecordFactory(custom_factory)

5.2 结构化日志

python
"""
demo_structured_logging.py
使用结构化日志 + OTel,实现可查询的日志数据
"""
import json
import logging
from opentelemetry._logs import set_logger_provider
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import SimpleLogRecordProcessor, ConsoleLogExporter
from opentelemetry.sdk.resources import Resource

resource = Resource.create({"service.name": "structured-log-demo"})
logger_provider = LoggerProvider(resource=resource)
logger_provider.add_log_record_processor(
    SimpleLogRecordProcessor(ConsoleLogExporter())
)
set_logger_provider(logger_provider)

handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider)
logging.getLogger().addHandler(handler)
logging.getLogger().setLevel(logging.INFO)

logger = logging.getLogger("myapp")

# 结构化日志:所有上下文信息都作为 Attributes
logger.info("HTTP 请求处理完成", extra={
    "http.method": "POST",
    "http.url": "/api/orders",
    "http.status_code": 201,
    "http.duration_ms": 145,
    "user.id": "user-001",
    "order.id": "ORD-12345",
    "order.total": 299.99,
})

# 在后端系统(Loki/Elasticsearch)中,这些 Attributes 都可以作为查询维度:
# 查询示例:
#   {service_name="structured-log-demo"} | json | http_status_code >= 400
#   {service_name="structured-log-demo"} | json | user_id="user-001"

logger_provider.shutdown()

5.3 BatchLogRecordProcessor 配置

python
"""
生产环境使用 BatchLogRecordProcessor
"""
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor, ConsoleLogExporter

# Simple:同步导出,每条日志立即发送(开发用)
# Batch:异步批量导出,不阻塞业务线程(生产用)

processor = BatchLogRecordProcessor(
    ConsoleLogExporter(),
    max_queue_size=2048,           # 队列最大容量,超过会丢弃
    schedule_delay_millis=5000,    # 每 5 秒导出一批
    max_export_batch_size=512,     # 单次最多导出 512 条
    export_timeout_millis=30000,   # 导出超时 30 秒
)

logger_provider.add_log_record_processor(processor)

5.4 多 Handler 共存

在实际项目中,通常需要同时保留控制台输出和 OTel 导出:

python
"""
demo_multi_handler.py
保留原有的控制台日志输出,同时通过 OTel 导出到后端
"""
import logging
import sys
from opentelemetry._logs import set_logger_provider
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor, ConsoleLogExporter
from opentelemetry.sdk.resources import Resource

resource = Resource.create({"service.name": "multi-handler-demo"})
logger_provider = LoggerProvider(resource=resource)
logger_provider.add_log_record_processor(
    BatchLogRecordProcessor(ConsoleLogExporter())
)
set_logger_provider(logger_provider)

# Handler 1: 传统的控制台输出(人类可读格式)
console_handler = logging.StreamHandler(sys.stdout)
console_handler.setFormatter(logging.Formatter(
    "[%(asctime)s] [%(levelname)s] [%(name)s] %(message)s"
))

# Handler 2: OTel LoggingHandler(导出到后端)
otel_handler = LoggingHandler(
    level=logging.NOTSET,
    logger_provider=logger_provider,
)

# 两个 Handler 都加上
root_logger = logging.getLogger()
root_logger.addHandler(console_handler)
root_logger.addHandler(otel_handler)
root_logger.setLevel(logging.INFO)

logger = logging.getLogger("myapp")

# 每条日志会同时:
#   1. 在控制台以人类可读的格式输出
#   2. 通过 OTel 以结构化格式导出到后端
logger.info("应用启动完成")
logger.warning("配置项 XYZ 已废弃,请迁移到新配置")

logger_provider.shutdown()

5.5 使用 OTLP 导出日志

python
"""
demo_otlp_log_export.py
通过 OTLP 将日志导出到 Collector / Loki
"""
import logging
from opentelemetry._logs import set_logger_provider
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor
from opentelemetry.sdk.resources import Resource

# OTLP gRPC 方式
from opentelemetry.exporter.otlp.proto.grpc._log_exporter import OTLPLogExporter as GrpcLogExporter

# OTLP HTTP 方式
# from opentelemetry.exporter.otlp.proto.http._log_exporter import OTLPLogExporter as HttpLogExporter

resource = Resource.create({
    "service.name": "production-app",
    "deployment.environment": "production",
})

logger_provider = LoggerProvider(resource=resource)

# gRPC 导出
grpc_exporter = GrpcLogExporter(
    endpoint="localhost:4317",    # Collector 的 gRPC 端口
    insecure=True,                # 开发环境不用 TLS
)

# HTTP 导出(备选)
# http_exporter = HttpLogExporter(
#     endpoint="http://localhost:4318/v1/logs",  # Collector 的 HTTP 端口
# )

logger_provider.add_log_record_processor(
    BatchLogRecordProcessor(grpc_exporter)
)
set_logger_provider(logger_provider)

handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider)
logging.getLogger().addHandler(handler)
logging.getLogger().setLevel(logging.INFO)

logger = logging.getLogger("production-app")
logger.info("日志将通过 OTLP 发送到 Collector")

logger_provider.shutdown()

5.6 自定义 LogRecordProcessor

python
"""
demo_custom_processor.py
自定义 LogRecordProcessor:实现日志过滤和脱敏
"""
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler, LogData
from opentelemetry.sdk._logs.export import (
    LogExporter,
    LogExportResult,
    SimpleLogRecordProcessor,
    ConsoleLogExporter,
)
from opentelemetry._logs import set_logger_provider, SeverityNumber
from opentelemetry.sdk.resources import Resource
import logging
import re
from typing import Sequence


class FilteringLogRecordProcessor(SimpleLogRecordProcessor):
    """只导出 WARNING 及以上级别的日志到后端(节省成本)"""

    def __init__(self, exporter, min_severity: int = SeverityNumber.WARN.value):
        super().__init__(exporter)
        self.min_severity = min_severity

    def emit(self, log_data: LogData) -> None:
        if log_data.log_record.severity_number.value >= self.min_severity:
            super().emit(log_data)


class SanitizingLogExporter(LogExporter):
    """对敏感信息进行脱敏后再导出"""

    SENSITIVE_PATTERNS = [
        (re.compile(r'password["\s:=]+["\']?(\S+)["\']?', re.IGNORECASE), r'password=***'),
        (re.compile(r'token["\s:=]+["\']?(\S+)["\']?', re.IGNORECASE), r'token=***'),
        (re.compile(r'\b\d{4}[-\s]?\d{4}[-\s]?\d{4}[-\s]?\d{4}\b'), r'****-****-****-****'),
    ]

    def __init__(self, delegate: LogExporter):
        self._delegate = delegate

    def export(self, batch: Sequence[LogData]) -> LogExportResult:
        for log_data in batch:
            body = log_data.log_record.body
            if isinstance(body, str):
                for pattern, replacement in self.SENSITIVE_PATTERNS:
                    body = pattern.sub(replacement, body)
                log_data.log_record.body = body
        return self._delegate.export(batch)

    def shutdown(self) -> None:
        self._delegate.shutdown()

    def force_flush(self, timeout_millis: int = 30000) -> bool:
        return self._delegate.force_flush(timeout_millis)


# 使用示例
resource = Resource.create({"service.name": "sanitized-logs"})
logger_provider = LoggerProvider(resource=resource)

# 套娃组合:过滤 + 脱敏
sanitizing_exporter = SanitizingLogExporter(ConsoleLogExporter())
filtering_processor = FilteringLogRecordProcessor(sanitizing_exporter)
logger_provider.add_log_record_processor(filtering_processor)

set_logger_provider(logger_provider)
handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider)
logging.getLogger().addHandler(handler)
logging.getLogger().setLevel(logging.DEBUG)

logger = logging.getLogger("secure-app")

logger.debug("这条 DEBUG 日志不会被导出")  # 被 FilteringProcessor 过滤掉
logger.info("这条 INFO 日志也不会被导出")   # 同上
logger.warning("用户密码重试: password=abc123")  # 会导出,且 password 被脱敏
logger.error("支付失败: token=sk-live-xxxx, 卡号 6222-0200-1234-5678")  # 卡号和 token 被脱敏

logger_provider.shutdown()

六、实战:完整的 Traces + Logs 可观测性

Demo:FastAPI 应用的日志与追踪联动

python
"""
demo_fastapi_logs.py
FastAPI 应用中同时使用 Traces 和 Logs

安装依赖:
  pip install fastapi uvicorn opentelemetry-api opentelemetry-sdk \
              opentelemetry-instrumentation-fastapi
"""
import logging
import uvicorn
from fastapi import FastAPI, HTTPException

from opentelemetry import trace
from opentelemetry._logs import set_logger_provider
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor, ConsoleSpanExporter
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import SimpleLogRecordProcessor, ConsoleLogExporter
from opentelemetry.sdk.resources import Resource
from opentelemetry.instrumentation.fastapi import FastAPIInstrumentor

# ========== 统一配置 ==========
resource = Resource.create({
    "service.name": "fastapi-demo",
    "service.version": "1.0.0",
})

# Traces
trace_provider = TracerProvider(resource=resource)
trace_provider.add_span_processor(SimpleSpanProcessor(ConsoleSpanExporter()))
trace.set_tracer_provider(trace_provider)
tracer = trace.get_tracer("fastapi-demo")

# Logs
log_provider = LoggerProvider(resource=resource)
log_provider.add_log_record_processor(
    SimpleLogRecordProcessor(ConsoleLogExporter())
)
set_logger_provider(log_provider)

otel_handler = LoggingHandler(level=logging.NOTSET, logger_provider=log_provider)
logging.getLogger().addHandler(otel_handler)
logging.getLogger().setLevel(logging.INFO)

logger = logging.getLogger("fastapi-demo")

# ========== FastAPI 应用 ==========
app = FastAPI()
FastAPIInstrumentor.instrument_app(app)

USERS_DB = {
    "1": {"name": "张三", "email": "zhangsan@example.com"},
    "2": {"name": "李四", "email": "lisi@example.com"},
}

@app.get("/users/{user_id}")
async def get_user(user_id: str):
    logger.info("查询用户", extra={"user.id": user_id})

    with tracer.start_as_current_span("db_lookup") as span:
        span.set_attribute("db.operation", "SELECT")
        span.set_attribute("db.table", "users")

        user = USERS_DB.get(user_id)
        if not user:
            logger.warning("用户不存在", extra={"user.id": user_id})
            raise HTTPException(status_code=404, detail="User not found")

        logger.info("用户查询成功", extra={
            "user.id": user_id,
            "user.name": user["name"],
        })

    return user

@app.get("/orders")
async def create_order():
    logger.info("创建订单请求")

    with tracer.start_as_current_span("create_order"):
        try:
            result = process_order()
            logger.info("订单创建成功", extra={"order.id": result})
            return {"order_id": result}
        except Exception as e:
            logger.error("订单创建失败: %s", str(e), exc_info=True)
            raise HTTPException(status_code=500, detail=str(e))

def process_order():
    with tracer.start_as_current_span("process_order"):
        logger.debug("开始处理订单逻辑")
        return "ORD-20250315-001"

if __name__ == "__main__":
    uvicorn.run(app, host="0.0.0.0", port=8000)

运行后访问 http://localhost:8000/users/1,可以看到:

  1. Span 输出:包含 FastAPI 自动检测生成的 HTTP Span + 手动创建的 db_lookup Span
  2. LogRecord 输出:每条日志都携带了对应的 trace_idspan_id

七、与现有日志系统集成的策略

7.1 参考 trpc-agent 的日志架构

在实际项目中,通常已经有一套自定义的日志系统。如 trpc-agent 项目定义了 BaseLogger 抽象层:

trpc_agent/log/
├── __init__.py           ← 统一导出
├── _base_logger.py       ← 抽象基类(定义 debug/info/warning/error/fatal 接口)
├── _default_logger.py    ← 默认实现(基于 Python logging)
└── _logger.py            ← 全局 logger 管理(get_logger/set_logger/register_logger)

集成策略

python
"""
将现有日志系统与 OTel 集成的模式
"""
import logging

# 方式一:在底层 Python logger 上挂载 OTel Handler
# 适用于:日志系统最终基于 Python logging(如 trpc-agent 的 DefaultLogger)
# 只需在 root logger 或业务 logger 上添加 OTel Handler 即可

from opentelemetry.sdk._logs import LoggingHandler

otel_handler = LoggingHandler(level=logging.NOTSET, logger_provider=logger_provider)

# 找到实际使用的 Python logger 并添加 handler
underlying_logger = logging.getLogger("trpc_agent")
underlying_logger.addHandler(otel_handler)
# 这样所有通过 trpc_agent logger 记录的日志都会被 OTel 捕获


# 方式二:创建自定义的 OTel-aware Logger
# 适用于:日志系统不基于 Python logging,或需要更精细的控制

from opentelemetry.sdk._logs import LoggerProvider
from opentelemetry._logs import get_logger_provider, SeverityNumber

class OTelAwareLogger:
    """集成了 OTel 的自定义 Logger"""

    def __init__(self, name: str):
        self.name = name
        self._otel_logger = get_logger_provider().get_logger(name)

    def info(self, message: str, **attributes):
        # 同时输出到控制台和 OTel
        print(f"[INFO] {message}")
        self._otel_logger.emit(
            self._create_log_record(message, SeverityNumber.INFO, attributes)
        )

7.2 集成决策树

你的日志系统是什么?

├── 基于 Python logging(最常见)
│   └── ✅ 直接添加 OTel LoggingHandler → 最简单

├── 基于 loguru
│   └── ✅ loguru 底层也用 logging → 可以在 root logger 上加 Handler
│       或使用 loguru 的 add() 自定义 sink

├── 基于 structlog
│   └── ✅ structlog 可以配置输出到 logging → 在 logging 上加 Handler
│       或实现自定义的 structlog processor

├── 完全自定义的日志系统
│   └── 需要手动创建 LogRecord 并通过 OTel Logger 发送

└── 不想改任何代码
    └── ✅ 使用 Collector 的 filelog receiver 从日志文件中采集
        (不推荐,会丢失结构化信息和 Trace 关联)

八、Collector 中的 Logs Pipeline 配置

8.1 Collector 配置示例

yaml
# otel-collector-config.yaml
receivers:
  otlp:
    protocols:
      grpc:
        endpoint: 0.0.0.0:4317
      http:
        endpoint: 0.0.0.0:4318

processors:
  batch:
    timeout: 5s
    send_batch_size: 1024

  # 日志级别过滤:只保留 WARN 及以上
  filter/severity:
    logs:
      log_record:
        - 'severity_number < 13'  # 过滤掉 DEBUG 和 INFO

  # 属性处理:添加/删除/修改属性
  attributes/logs:
    actions:
      - key: "cluster.name"
        value: "production-cluster"
        action: upsert
      - key: "password"
        action: delete

exporters:
  # 发送到 Grafana Loki
  loki:
    endpoint: http://loki:3100/loki/api/v1/push
    labels:
      attributes:
        service.name: "service"
        severity_text: "level"

  # 发送到 Elasticsearch
  elasticsearch:
    endpoints: ["http://elasticsearch:9200"]
    logs_index: "otel-logs"

service:
  pipelines:
    logs:
      receivers: [otlp]
      processors: [batch, attributes/logs]   # 生产环境全量收集
      exporters: [loki]

    logs/filtered:
      receivers: [otlp]
      processors: [batch, filter/severity]   # 只保留 WARN+ 到 ES
      exporters: [elasticsearch]

九、日志格式化:让 trace_id 出现在控制台输出中

在开发和排查问题时,能在控制台日志中直接看到 trace_id 非常有用:

python
"""
demo_log_format_with_trace.py
在传统的控制台日志格式中包含 trace_id 和 span_id
"""
import logging
from opentelemetry import trace
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import SimpleSpanProcessor, ConsoleSpanExporter
from opentelemetry.sdk.resources import Resource


class TraceIdFormatter(logging.Formatter):
    """自定义 Formatter,在日志中注入 trace_id 和 span_id"""

    def format(self, record):
        span = trace.get_current_span()
        ctx = span.get_span_context()

        if ctx and ctx.trace_id != 0:
            record.trace_id = format(ctx.trace_id, '032x')
            record.span_id = format(ctx.span_id, '016x')
        else:
            record.trace_id = "0" * 32
            record.span_id = "0" * 16

        return super().format(record)


# 配置 Trace
resource = Resource.create({"service.name": "trace-in-log-demo"})
provider = TracerProvider(resource=resource)
provider.add_span_processor(SimpleSpanProcessor(ConsoleSpanExporter()))
trace.set_tracer_provider(provider)
tracer = trace.get_tracer(__name__)

# 配置 logging 使用自定义 Formatter
handler = logging.StreamHandler()
handler.setFormatter(TraceIdFormatter(
    "[%(asctime)s] [%(levelname)s] [trace=%(trace_id)s span=%(span_id)s] %(message)s"
))

logger = logging.getLogger("myapp")
logger.addHandler(handler)
logger.setLevel(logging.INFO)

# 使用
logger.info("无 Trace 上下文的日志")
# 输出: [2025-03-15 10:23:45] [INFO] [trace=00000000... span=00000000...] 无 Trace 上下文的日志

with tracer.start_as_current_span("request-handler"):
    logger.info("有 Trace 上下文的日志")
    # 输出: [2025-03-15 10:23:45] [INFO] [trace=abc123def456... span=789ghi...] 有 Trace 上下文的日志

    with tracer.start_as_current_span("sub-operation"):
        logger.info("子 Span 中的日志")
        # trace_id 与父 Span 相同,但 span_id 不同

provider.shutdown()

也可以通过环境变量自动注入(需要 opentelemetry-instrumentation-logging):

bash
pip install opentelemetry-instrumentation-logging

# 设置环境变量启用日志关联
export OTEL_PYTHON_LOG_CORRELATION=true

# 之后 Python logging 会自动在 LogRecord 中添加以下字段:
#   otelTraceID
#   otelSpanID
#   otelTraceSampled
#   otelServiceName
python
# 配合使用
logging.basicConfig(
    format="%(asctime)s [%(levelname)s] trace=%(otelTraceID)s span=%(otelSpanID)s - %(message)s",
    level=logging.INFO,
)

十、常见问题 QA

Q1:为什么 OTel Logs 的包名有下划线前缀(_logs)?

A:下划线前缀表示该 API 尚未完全稳定

python
from opentelemetry._logs import set_logger_provider      # 注意 _logs
from opentelemetry.sdk._logs import LoggerProvider         # 注意 _logs

截至 2025 年底,OTel Python 的 Logs API 仍处于 "实验性" 阶段。这意味着:

  • API 可能会在未来版本中有变化
  • 功能已经可用,只是接口还没有最终冻结
  • Traces 和 Metrics 的 API 已经稳定(没有下划线)

实际影响:可以放心在生产中使用,只是升级 SDK 版本时需要关注变更日志。

Q2:OTel LoggingHandler 会不会导致日志被记录两次?

A:可能会。需要注意 Handler 的配置。

python
# ❌ 错误:root logger 已有 StreamHandler,再加 OTel Handler 不会重复
# 但如果 OTel 用的是 ConsoleLogExporter,日志就会在控制台出现两次

root_logger = logging.getLogger()
root_logger.addHandler(logging.StreamHandler())          # 第一次输出到控制台
root_logger.addHandler(LoggingHandler(...))               # OTel 捕获
# 如果 OTel 用 ConsoleLogExporter,又会输出到控制台     # 第二次

# ✅ 方案一:OTel 只用 OTLP Exporter(不用 Console)
logger_provider.add_log_record_processor(
    BatchLogRecordProcessor(OTLPLogExporter(endpoint="localhost:4317"))
)

# ✅ 方案二:去掉原有的控制台 Handler,完全由 OTel 管理
root_logger = logging.getLogger()
root_logger.handlers.clear()  # 清除原有 Handler
root_logger.addHandler(LoggingHandler(logger_provider=logger_provider))

Q3:如果 Exporter 发送失败,日志会丢失吗?

A:取决于 Processor 类型和配置:

Processor发送失败的行为
SimpleLogRecordProcessor日志丢失,不重试(同步导出)
BatchLogRecordProcessor日志在队列中,会重试,但队列满了会丢弃最早的
python
# 降低丢失风险的配置
processor = BatchLogRecordProcessor(
    exporter,
    max_queue_size=8192,           # 增大队列(默认 2048)
    schedule_delay_millis=2000,    # 更频繁地导出(默认 5000)
    export_timeout_millis=60000,   # 增大超时(默认 30000)
)

# 但即使如此,极端情况下(Exporter 持续失败 + 队列满)日志仍会丢失
# 最佳实践:同时保留本地文件日志作为兜底

Q4:日志量很大,怎么控制导出成本?

A:几种策略:

python
# 策略一:只导出特定级别以上的日志
otel_handler = LoggingHandler(
    level=logging.WARNING,  # 只导出 WARNING 及以上
    logger_provider=logger_provider,
)

# 策略二:在 Collector 中过滤(推荐,灵活度最高)
# 见上文 Collector 配置中的 filter/severity

# 策略三:自定义 Processor 按条件过滤(见 5.6 自定义 Processor 示例)

# 策略四:控制 extra 属性的数量
# 避免在日志中塞入大量属性,每个属性都会增加传输和存储成本
logger.info("请求完成", extra={
    "http.method": "GET",        # ✅ 有用的维度
    "http.status": 200,          # ✅ 有用的维度
    # "request.body": huge_json, # ❌ 避免存储大对象
    # "response.body": ...,      # ❌ 避免存储大对象
})

Q5:如何在日志中记录异常堆栈?

A:使用 exc_info=Truelogger.exception()

python
try:
    risky_operation()
except Exception as e:
    # 方式一:exc_info=True
    logger.error("操作失败: %s", str(e), exc_info=True)

    # 方式二:logger.exception()(等价于 error + exc_info=True)
    logger.exception("操作失败")

    # 方式三:同时在 Span 中记录异常
    span = trace.get_current_span()
    span.record_exception(e)
    span.set_status(trace.StatusCode.ERROR, str(e))
    logger.error("操作失败: %s", str(e), exc_info=True)

# OTel LogRecord 会包含完整的异常堆栈信息
# 在后端系统中可以按异常类型过滤和统计

Q6:OTel 的 Logger 和 Python 的 Logger 是什么关系?

A:完全是两个独立的概念,通过 LoggingHandler 桥接。

Python Logger                        OTel Logger
─────────────                        ──────────
logging.getLogger("myapp")           logger_provider.get_logger("myapp")
│                                    │
├── 由 logging 模块管理               ├── 由 LoggerProvider 管理
├── 支持 Handler + Formatter          ├── 只负责创建 LogRecord
├── 开发者直接使用 ✅                  ├── 通常不直接使用
└── 生态成熟(loguru 等兼容)          └── 通过 LoggingHandler 自动接入

桥接关系:
  Python Logger ──Handler──→ OTel LoggingHandler ──→ OTel Logger ──→ Processor ──→ Exporter

总结:你写代码时用 Python logging,OTel 在背后自动采集和导出。

Q7:loguru 可以和 OTel 一起用吗?

A:可以,有两种方式:

python
# 方式一:让 loguru 输出到 Python logging(推荐)
from loguru import logger
import logging

class InterceptHandler(logging.Handler):
    def emit(self, record):
        logger.opt(depth=6, exception=record.exc_info).log(
            record.levelname, record.getMessage()
        )

# 然后在 Python logging 上挂 OTel Handler
logging.getLogger().addHandler(otel_handler)

# 方式二:使用 loguru 的 sink 功能
from loguru import logger

def otel_sink(message):
    record = message.record
    otel_logger = logger_provider.get_logger(record["name"])
    # 手动创建并发送 OTel LogRecord
    # ...

logger.add(otel_sink, level="INFO")

Q8:Resource 是 Traces 和 Logs 共享的吗?

A:是的,Resource 用来标识"谁产生了这些遥测数据"。建议 Traces、Metrics、Logs 使用同一个 Resource,这样在后端系统中才能正确关联。

python
# ✅ 推荐:共享同一个 Resource
resource = Resource.create({
    "service.name": "my-service",
    "service.version": "1.0.0",
})

trace_provider = TracerProvider(resource=resource)       # Traces 用同一个 Resource
logger_provider = LoggerProvider(resource=resource)      # Logs 用同一个 Resource
meter_provider = MeterProvider(resource=resource)        # Metrics 用同一个 Resource

# ❌ 不推荐:每个 Provider 用不同的 Resource
# 这会导致后端系统无法将同一服务的 Traces、Logs、Metrics 关联起来

Q9:如何优雅关闭日志导出?

A

python
import atexit

logger_provider = LoggerProvider(resource=resource)
logger_provider.add_log_record_processor(
    BatchLogRecordProcessor(exporter)
)

# 方式一:使用 atexit(推荐)
atexit.register(logger_provider.shutdown)

# 方式二:在 finally 块中关闭
try:
    run_app()
finally:
    logger_provider.shutdown()  # 刷新队列中的日志

# 方式三:在 FastAPI 的 lifespan 中关闭
from contextlib import asynccontextmanager

@asynccontextmanager
async def lifespan(app):
    yield
    logger_provider.shutdown()
    trace_provider.shutdown()

app = FastAPI(lifespan=lifespan)

shutdown() 的作用:

  1. 通知 BatchProcessor 立即导出队列中剩余的日志
  2. 等待导出完成
  3. 关闭 Exporter 连接

如果不调用 shutdown(),BatchProcessor 队列中的日志可能会丢失。

Q10:生产环境的日志配置模板是怎样的?

A

python
"""
production_log_setup.py
生产环境日志 + Trace 配置模板
"""
import os
import logging
import atexit

from opentelemetry import trace
from opentelemetry._logs import set_logger_provider
from opentelemetry.sdk.trace import TracerProvider
from opentelemetry.sdk.trace.export import BatchSpanProcessor
from opentelemetry.sdk._logs import LoggerProvider, LoggingHandler
from opentelemetry.sdk._logs.export import BatchLogRecordProcessor
from opentelemetry.sdk.resources import Resource
from opentelemetry.sdk.trace.sampling import ParentBased, TraceIdRatioBased
from opentelemetry.exporter.otlp.proto.grpc.trace_exporter import OTLPSpanExporter
from opentelemetry.exporter.otlp.proto.grpc._log_exporter import OTLPLogExporter


def setup_observability():
    """一次性配置 Traces + Logs 的可观测性"""

    # 从环境变量读取配置
    service_name = os.getenv("OTEL_SERVICE_NAME", "unknown-service")
    otlp_endpoint = os.getenv("OTEL_EXPORTER_OTLP_ENDPOINT", "localhost:4317")
    sample_rate = float(os.getenv("OTEL_TRACES_SAMPLER_ARG", "0.1"))
    log_level = os.getenv("OTEL_LOG_LEVEL", "INFO").upper()

    # 统一 Resource
    resource = Resource.create({
        "service.name": service_name,
        "service.version": os.getenv("SERVICE_VERSION", "unknown"),
        "deployment.environment": os.getenv("DEPLOYMENT_ENV", "unknown"),
    })

    # ===== Traces =====
    trace_provider = TracerProvider(
        resource=resource,
        sampler=ParentBased(root=TraceIdRatioBased(sample_rate)),
    )
    trace_provider.add_span_processor(
        BatchSpanProcessor(OTLPSpanExporter(endpoint=otlp_endpoint, insecure=True))
    )
    trace.set_tracer_provider(trace_provider)

    # ===== Logs =====
    log_provider = LoggerProvider(resource=resource)
    log_provider.add_log_record_processor(
        BatchLogRecordProcessor(
            OTLPLogExporter(endpoint=otlp_endpoint, insecure=True),
            max_queue_size=4096,
            schedule_delay_millis=3000,
        )
    )
    set_logger_provider(log_provider)

    # OTel Handler(只用于导出到后端)
    otel_handler = LoggingHandler(
        level=getattr(logging, log_level),
        logger_provider=log_provider,
    )
    logging.getLogger().addHandler(otel_handler)

    # 控制台 Handler(人类可读格式,开发调试时查看)
    console_handler = logging.StreamHandler()
    console_handler.setFormatter(logging.Formatter(
        "[%(asctime)s] [%(levelname)s] [%(name)s] %(message)s"
    ))
    logging.getLogger().addHandler(console_handler)
    logging.getLogger().setLevel(getattr(logging, log_level))

    # 优雅关闭
    atexit.register(trace_provider.shutdown)
    atexit.register(log_provider.shutdown)

    return trace_provider, log_provider


# 使用
trace_provider, log_provider = setup_observability()

# 启动命令示例:
# OTEL_SERVICE_NAME=order-service \
# OTEL_EXPORTER_OTLP_ENDPOINT=collector.internal:4317 \
# OTEL_TRACES_SAMPLER_ARG=0.1 \
# DEPLOYMENT_ENV=production \
# python app.py

十一、小结与知识图谱

                        OTel Logs

          ┌─────────────────┼─────────────────┐
          │                 │                 │
       桥接策略           核心组件          最佳实践
          │                 │                 │
   ┌──────┼──────┐    ┌────┼────┐      ┌────┼────┐
   │      │      │    │    │    │      │    │    │
 Python loguru struct Logger  Log    共享  级别  脱敏
logging       log   Provider Record Resource 过滤
   │                  │    Processor
   └──── LoggingHandler    Exporter
         (桥接层)     │
                      ├── Simple(开发)
                      └── Batch(生产)

关联能力:
  Trace Context ──自动注入──→ LogRecord.trace_id / span_id

                    在 Grafana 中 Trace ↔ Log 双向跳转

本阶段核心收获

  1. 理解 OTel Logs 的设计哲学:不替代现有日志系统,而是桥接和增强
  2. 掌握 LoggerProviderLoggingHandlerLogRecordProcessorLogExporter 的配置链路
  3. 掌握日志与 Trace 自动关联的原理和实现
  4. 能够编写自定义 Processor 实现过滤和脱敏
  5. 理解 Simple vs Batch Processor 的选型
  6. 掌握生产环境的日志配置模板

下一阶段预告:阶段 5 将深入 Context 传播机制,学习 W3C TraceContext、Baggage 以及跨服务/跨线程/异步场景下的 Context 传播策略。