返回博客

可观测性:OpenTelemetry 统一遥测框架

coomia-dip 基于 OpenTelemetry 构建统一的可观测性框架,覆盖 Traces(分布式追踪)、Metrics(指标监控)和 Logs(结构化日志)三大支柱。通过 gRPC 拦截器自动注入追踪上下文、自定义指标收集器捕获业务指标、结构化日志关联 Trace ID 实现全链路可观测。后端集成 Jaeger(追踪)、Prometheus + Grafana(指标)和 Loki(日志)。本文从遥测架构、自动instrumentation、自定义指标、日志关联到告警体系,完整解析可观测性方案。

Coomia发布于 2025年10月4日7 分钟阅读
分享本文Twitter / X

系列:S6 平台工程 · 第 20 篇 | 难度:高级 | 阅读时间:18 分钟

可观测性:OpenTelemetry 统一遥测框架

#TL;DR

coomia-dip 基于 OpenTelemetry 构建统一的可观测性框架,覆盖 Traces(分布式追踪)、Metrics(指标监控)和 Logs(结构化日志)三大支柱。通过 gRPC 拦截器自动注入追踪上下文、自定义指标收集器捕获业务指标、结构化日志关联 Trace ID 实现全链路可观测。后端集成 Jaeger(追踪)、Prometheus + Grafana(指标)和 Loki(日志)。本文从遥测架构、自动instrumentation、自定义指标、日志关联到告警体系,完整解析可观测性方案。

#1. 可观测性三大支柱

#1.1 统一遥测模型

Code
┌────────────────────────────────────────────────┐
│            OpenTelemetry SDK                    │
│  ┌──────────┐  ┌──────────┐  ┌──────────────┐ │
│  │  Traces   │  │ Metrics  │  │    Logs      │ │
│  │ (追踪)    │  │ (指标)   │  │   (日志)     │ │
│  └─────┬────┘  └─────┬────┘  └──────┬───────┘ │
│        │             │              │          │
│  ┌─────▼─────────────▼──────────────▼───────┐  │
│  │         OTLP Exporter                    │  │
│  │    (统一导出协议)                          │  │
│  └─────────────────┬────────────────────────┘  │
└────────────────────┼───────────────────────────┘
                     │
        ┌────────────┼────────────┐
        │            │            │
   ┌────▼───┐  ┌────▼────┐  ┌───▼────┐
   │ Jaeger │  │Prometheus│  │  Loki  │
   │ (追踪) │  │ (指标)   │  │ (日志) │
   └────────┘  └─────────┘  └────────┘
        │            │            │
   ┌────▼────────────▼────────────▼───┐
   │          Grafana Dashboard        │
   │      (统一可视化仪表盘)            │
   └──────────────────────────────────┘

#1.2 对标 Palantir Foundry

能力Palantir Foundrycoomia-dip
遥测框架内部OpenTelemetry
分布式追踪内部Jaeger
指标监控内部Prometheus + Grafana
日志管理内部Loki + Grafana
自定义指标有限完整支持

#2. 分布式追踪

#2.1 gRPC 自动追踪

Python
class OTelGrpcInterceptor(grpc.aio.UnaryUnaryClientInterceptor):
    """OpenTelemetry gRPC 追踪拦截器"""

    def __init__(self):
        self._tracer = trace.get_tracer("coomia-dip-sdk")

    async def intercept_unary_unary(self, continuation, client_call_details, request):
        method = client_call_details.method
        service, method_name = self._parse_method(method)

        with self._tracer.start_as_current_span(
            f"grpc.{service}/{method_name}",
            kind=trace.SpanKind.CLIENT,
            attributes={
                "rpc.system": "grpc",
                "rpc.service": service,
                "rpc.method": method_name,
                "onto.object_type": self._extract_object_type(request),
            },
        ) as span:
            # 注入 trace context 到 gRPC metadata
            metadata = list(client_call_details.metadata or [])
            inject(metadata, setter=GrpcMetadataSetter())
            new_details = client_call_details._replace(metadata=metadata)

            try:
                response = await continuation(new_details, request)
                span.set_status(StatusCode.OK)
                span.set_attribute("rpc.response_size", response.ByteSize())
                return response
            except grpc.aio.AioRpcError as e:
                span.set_status(StatusCode.ERROR, str(e))
                span.set_attribute("rpc.grpc.status_code", e.code().value[0])
                raise


class OTelGrpcServerInterceptor(grpc.aio.ServerInterceptor):
    """服务端追踪拦截器"""

    async def intercept_service(self, continuation, handler_call_details):
        # 从 gRPC metadata 提取 trace context
        context = extract(handler_call_details.invocation_metadata, getter=GrpcMetadataGetter())

        with self._tracer.start_as_current_span(
            f"grpc.server.{handler_call_details.method}",
            kind=trace.SpanKind.SERVER,
            context=context,
        ) as span:
            try:
                response = await continuation(handler_call_details)
                span.set_status(StatusCode.OK)
                return response
            except Exception as e:
                span.set_status(StatusCode.ERROR, str(e))
                span.record_exception(e)
                raise

#2.2 跨 Layer 追踪

Code
Client SDK → API Gateway → Ontology Service → Data Service → Iceberg
    │            │               │                │            │
    ├── span ────┤               │                │            │
    │            ├── span ───────┤                │            │
    │            │               ├── span ────────┤            │
    │            │               │                ├── span ────┤
    │            │               │                │            │
    └────────────┴───────────────┴────────────────┴────────────┘
                        TraceID: abc-123

#2.3 自定义 Span 属性

Python
# Ontology 操作的自定义 Span 属性
ONTO_SPAN_ATTRIBUTES = {
    "onto.object_type": "Employee",
    "onto.operation": "get",
    "onto.object_id": "emp-001",
    "onto.classification_level": "B2",
    "onto.masking_applied": True,
    "onto.result_count": 1,
    "onto.query_type": "get_by_id",
}

#3. 指标监控

#3.1 平台指标

Python
class OntoPlatformMetrics:
    """coomia-dip 平台指标"""

    def __init__(self):
        self._meter = metrics.get_meter("coomia-dip")

        # 请求指标
        self.request_counter = self._meter.create_counter(
            "onto.requests.total",
            description="Total number of requests",
            unit="1",
        )

        self.request_duration = self._meter.create_histogram(
            "onto.request.duration",
            description="Request duration in milliseconds",
            unit="ms",
        )

        self.active_requests = self._meter.create_up_down_counter(
            "onto.requests.active",
            description="Number of active requests",
        )

        # 对象操作指标
        self.object_operations = self._meter.create_counter(
            "onto.objects.operations.total",
            description="Total object operations",
            unit="1",
        )

        # 权限评估指标
        self.auth_evaluations = self._meter.create_counter(
            "onto.auth.evaluations.total",
            description="Total authorization evaluations",
        )

        self.auth_evaluation_duration = self._meter.create_histogram(
            "onto.auth.evaluation.duration",
            description="Authorization evaluation duration",
            unit="ms",
        )

        # 脱敏指标
        self.masking_operations = self._meter.create_counter(
            "onto.masking.operations.total",
            description="Total masking operations",
        )

        self.masking_duration = self._meter.create_histogram(
            "onto.masking.duration",
            description="Masking operation duration",
            unit="ms",
        )

        # 连接池指标
        self.connection_pool_size = self._meter.create_observable_gauge(
            "onto.connections.pool_size",
            callbacks=[self._observe_pool_size],
            description="Connection pool size",
        )

    def record_request(
        self,
        method: str,
        object_type: str,
        status: str,
        duration_ms: float,
    ):
        attributes = {
            "method": method,
            "object_type": object_type,
            "status": status,
        }
        self.request_counter.add(1, attributes)
        self.request_duration.record(duration_ms, attributes)

#3.2 Prometheus 集成

Python
class PrometheusExporterConfig:
    """Prometheus 导出器配置"""

    @staticmethod
    def setup():
        exporter = PrometheusMetricReader()
        provider = MeterProvider(metric_readers=[exporter])
        metrics.set_meter_provider(provider)

        # 启动 Prometheus HTTP 端点
        start_http_server(port=9464)

#4. 结构化日志

#4.1 日志与 Trace 关联

Python
class OTelLogHandler(logging.Handler):
    """OpenTelemetry 日志处理器 - 自动注入 Trace 上下文"""

    def emit(self, record: logging.LogRecord) -> None:
        # 获取当前 Span 的 Trace 上下文
        span = trace.get_current_span()
        if span.is_recording():
            ctx = span.get_span_context()
            record.trace_id = format(ctx.trace_id, "032x")
            record.span_id = format(ctx.span_id, "016x")
            record.trace_flags = ctx.trace_flags
        else:
            record.trace_id = "0" * 32
            record.span_id = "0" * 16
            record.trace_flags = 0

        # 结构化日志格式
        log_entry = {
            "timestamp": record.created,
            "level": record.levelname,
            "message": record.getMessage(),
            "logger": record.name,
            "trace_id": record.trace_id,
            "span_id": record.span_id,
            "service": "coomia-dip",
            "attributes": getattr(record, "attributes", {}),
        }

        self._export(log_entry)


# 使用示例
logger = logging.getLogger("onto.ontology")
logger.info(
    "Object accessed",
    extra={"attributes": {
        "object_type": "Employee",
        "object_id": "emp-001",
        "user_id": "user-123",
    }},
)
# 输出: {"trace_id": "abc123...", "message": "Object accessed", "attributes": {...}}

#4.2 Loki 集成

YAML
# Loki 日志标签配置
loki:
  labels:
    - service
    - level
    - trace_id
  pipeline_stages:
    - json:
        expressions:
          level: level
          trace_id: trace_id
          service: service
    - labels:
        level:
        service:

#5. 告警体系

#5.1 告警规则

YAML
# Prometheus 告警规则
groups:
  - name: coomia-dip-alerts
    rules:
      # 高延迟告警
      - alert: HighRequestLatency
        expr: histogram_quantile(0.99, onto_request_duration_bucket) > 1000
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "P99 latency exceeds 1s"

      # 错误率告警
      - alert: HighErrorRate
        expr: rate(onto_requests_total{status="error"}[5m]) / rate(onto_requests_total[5m]) > 0.05
        for: 2m
        labels:
          severity: critical
        annotations:
          summary: "Error rate exceeds 5%"

      # 权限拒绝告警
      - alert: HighAuthDenialRate
        expr: rate(onto_auth_evaluations_total{result="denied"}[5m]) > 100
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "High authorization denial rate"

      # 服务不可用
      - alert: ServiceDown
        expr: up{job=~"onto-.*"} == 0
        for: 1m
        labels:
          severity: critical
        annotations:
          summary: "Service {{ $labels.instance }} is down"

#6. 测试策略

Python
class TestObservability:
    def test_trace_propagation(self):
        """验证追踪上下文在 gRPC 调用链中传播"""
        with tracer.start_as_current_span("test-root") as root_span:
            response = await client.objects.get("Employee", "emp-001")

            # 验证子 Span 与根 Span 共享 Trace ID
            spans = exporter.get_finished_spans()
            trace_ids = set(s.context.trace_id for s in spans)
            assert len(trace_ids) == 1  # 所有 Span 共享一个 Trace ID

    def test_metrics_recorded(self):
        """验证指标被正确记录"""
        await client.objects.get("Employee", "emp-001")

        # 验证请求计数器递增
        metric_data = reader.get_metrics_data()
        request_metric = find_metric(metric_data, "onto.requests.total")
        assert request_metric.data_points[0].value >= 1

    def test_log_trace_correlation(self):
        """验证日志与追踪关联"""
        with tracer.start_as_current_span("test") as span:
            trace_id = format(span.get_span_context().trace_id, "032x")
            logger.info("test message")

            log_entry = log_exporter.get_last_entry()
            assert log_entry["trace_id"] == trace_id

#7. 生产最佳实践

#7.1 采样策略

场景采样率说明
正常请求1%减少存储开销
错误请求100%总是采集
慢请求 (>1s)100%总是采集
高分类数据访问100%合规要求

#7.2 数据保留

遥测类型保留时间存储
Traces7 天Jaeger + Elasticsearch
Metrics90 天Prometheus TSDB
Logs30 天Loki + 对象存储

#8. 总结

coomia-dip 基于 OpenTelemetry 的可观测性方案实现了三大遥测支柱的统一采集和关联分析。关键设计亮点:

  1. 统一框架:OpenTelemetry SDK 统一 Traces、Metrics、Logs
  2. 自动追踪:gRPC 拦截器自动注入追踪上下文
  3. 业务指标:自定义 Ontology 操作指标
  4. 日志关联:结构化日志自动关联 Trace ID
  5. 智能告警:基于 SLA 的多级告警规则

下一篇将探讨 coomia-dip 的 17 个仪表盘组件设计。