可观测性三支柱

支柱关注什么工具
日志(Logs)每次请求的详细信息结构化 JSON 日志
指标(Metrics)聚合统计数据Prometheus + Grafana
追踪(Traces)请求全链路耗时OpenTelemetry

1. 结构化日志

import logging
import json
import time
from datetime import datetime

# 配置 JSON 格式日志
logger = logging.getLogger('deepseek-api')
handler = logging.StreamHandler()
handler.setFormatter(logging.Formatter('%(message)s'))
logger.addHandler(handler)
logger.setLevel(logging.INFO)

class APILogger:
    @staticmethod
    def log_request(model, messages, request_id):
        log_entry = {
            'timestamp': datetime.now().isoformat(),
            'type': 'api_request',
            'request_id': request_id,
            'model': model,
            'input_tokens': sum(len(m.get('content', '')) for m in messages) // 4,  # 估算
            'stream': False
        }
        logger.info(json.dumps(log_entry, ensure_ascii=False))

    @staticmethod
    def log_response(request_id, usage, latency_ms, success=True):
        log_entry = {
            'timestamp': datetime.now().isoformat(),
            'type': 'api_response',
            'request_id': request_id,
            'prompt_tokens': usage.prompt_tokens if usage else 0,
            'completion_tokens': usage.completion_tokens if usage else 0,
            'total_tokens': usage.total_tokens if usage else 0,
            'latency_ms': round(latency_ms, 2),
            'success': success
        }
        logger.info(json.dumps(log_entry, ensure_ascii=False))

    @staticmethod
    def log_error(request_id, error_type, error_msg):
        log_entry = {
            'timestamp': datetime.now().isoformat(),
            'type': 'api_error',
            'request_id': request_id,
            'error_type': error_type,
            'error': str(error_msg)[:500]
        }
        logger.error(json.dumps(log_entry, ensure_ascii=False))

2. 请求追踪装饰器

import functools
import uuid

api_logger = APILogger()

def trace_api_call(model="deepseek-v4-flash"):
    """装饰器:自动记录每次 API 调用的日志和耗时"""
    def decorator(func):
        @functools.wraps(func)
        def wrapper(*args, **kwargs):
            request_id = str(uuid.uuid4())[:8]
            start = time.time()

            try:
                api_logger.log_request(model, kwargs.get('messages', []), request_id)
                result = func(*args, **kwargs)
                latency = (time.time() - start) * 1000
                usage = result.usage if hasattr(result, 'usage') else None
                api_logger.log_response(request_id, usage, latency)
                return result
            except Exception as e:
                latency = (time.time() - start) * 1000
                api_logger.log_error(request_id, type(e).__name__, str(e))
                api_logger.log_response(request_id, None, latency, success=False)
                raise

        return wrapper
    return decorator

# 使用
@trace_api_call(model="deepseek-v4-flash")
def call_deepseek(messages):
    return client.chat.completions.create(
        model='deepseek-v4-flash',
        messages=messages
    )

3. 实时指标收集

from collections import defaultdict
import threading

class MetricsCollector:
    def __init__(self):
        self.lock = threading.Lock()
        self.reset()

    def reset(self):
        self.total_requests = 0
        self.total_errors = 0
        self.total_tokens = 0
        self.total_latency_ms = 0
        self.latency_samples = []

    def record(self, tokens, latency_ms, is_error=False):
        with self.lock:
            self.total_requests += 1
            self.total_tokens += tokens
            self.total_latency_ms += latency_ms
            self.latency_samples.append(latency_ms)
            if is_error:
                self.total_errors += 1

            # 只保留最近 1000 个样本
            if len(self.latency_samples) > 1000:
                self.latency_samples = self.latency_samples[-1000:]

    def get_stats(self):
        with self.lock:
            samples = sorted(self.latency_samples)
            n = len(samples)
            return {
                'total_requests': self.total_requests,
                'error_rate': self.total_errors / max(self.total_requests, 1) * 100,
                'total_tokens': self.total_tokens,
                'avg_latency_ms': self.total_latency_ms / max(self.total_requests, 1),
                'p50_latency_ms': samples[n//2] if n > 0 else 0,
                'p95_latency_ms': samples[int(n*0.95)] if n > 0 else 0,
                'p99_latency_ms': samples[int(n*0.99)] if n > 0 else 0,
            }

metrics = MetricsCollector()

4. 成本追踪

class CostTracker:
    PRICES = {
        'deepseek-v4-flash': {
            'input_cache_hit': 0.02,
            'input_cache_miss': 1.00,
            'output': 2.00
        },
        'deepseek-v4-pro': {
            'input_cache_hit': 0.025,
            'input_cache_miss': 3.00,
            'output': 6.00
        }
    }

    def __init__(self):
        self.daily_costs = defaultdict(float)
        self.total_cost = 0.0

    def track(self, model, prompt_tokens, completion_tokens, cache_hit=True):
        prices = self.PRICES.get(model)
        if not prices:
            return

        input_price = prices['input_cache_hit'] if cache_hit else prices['input_cache_miss']
        cost = (prompt_tokens / 1_000_000) * input_price
        cost += (completion_tokens / 1_000_000) * prices['output']

        today = datetime.now().strftime('%Y-%m-%d')
        self.daily_costs[today] += cost
        self.total_cost += cost

        # 日消费超过阈值告警
        if self.daily_costs[today] > 100:  # 100元/天
            print(f"警告:今日 API 消费已超 ¥100(当前 ¥{self.daily_costs[today]:.2f})")

    def get_report(self):
        return {
            'today': round(self.daily_costs[datetime.now().strftime('%Y-%m-%d')], 2),
            'this_month': round(sum(v for k, v in self.daily_costs.items() if k.startswith(datetime.now().strftime('%Y-%m'))), 2),
            'total': round(self.total_cost, 2)
        }

cost_tracker = CostTracker()

5. 健康检查端点

from flask import Flask, jsonify

app = Flask(__name__)

@app.route('/health')
def health():
    stats = metrics.get_stats()
    costs = cost_tracker.get_report()

    return jsonify({
        'status': 'healthy' if stats['error_rate'] < 5 else 'degraded',
        'uptime_seconds': time.time() - start_time,
        'metrics': stats,
        'costs': costs
    })

@app.route('/metrics')
def prometheus_metrics():
    """Prometheus 格式指标端点"""
    stats = metrics.get_stats()
    costs = cost_tracker.get_report()
    return f"""# HELP api_requests_total Total API requests
# TYPE api_requests_total counter
api_requests_total {stats['total_requests']}
# HELP api_errors_total Total API errors
# TYPE api_errors_total counter
api_errors_total {stats['error_rate']}
# HELP api_latency_ms API latency in ms
# TYPE api_latency_ms gauge
api_latency_ms{{quantile="p50"}} {stats['p50_latency_ms']}
api_latency_ms{{quantile="p95"}} {stats['p95_latency_ms']}
api_latency_ms{{quantile="p99"}} {stats['p99_latency_ms']}
# HELP api_cost_yuan API cost in yuan
# TYPE api_cost_yuan gauge
api_cost_yuan{{period="today"}} {costs['today']}
api_cost_yuan{{period="month"}} {costs['this_month']}
api_cost_yuan{{period="total"}} {costs['total']}
""", 200, {'Content-Type': 'text/plain'}

使用总结

# 完整集成示例
@trace_api_call()
def safe_chat(messages):
    response = client.chat.completions.create(
        model='deepseek-v4-flash',
        messages=messages
    )
    # 记录指标
    metrics.record(response.usage.total_tokens, 0, is_error=False)
    # 记录成本
    cost_tracker.track('deepseek-v4-flash', response.usage.prompt_tokens, response.usage.completion_tokens)
    return response

result = safe_chat([{"role": "user", "content": "Hello"}])

# 查看统计
print(json.dumps(metrics.get_stats(), indent=2))
print(json.dumps(cost_tracker.get_report(), indent=2))