第9讲:监控与运维——可观测性与管理能力
一个分布式系统在生产环境中运行,必须具备完善的可观测性——能够实时了解系统状态、诊断问题、分析性能。
这一讲,我们为MiniKV添加监控、指标、日志、追踪和管理API,让它成为一个可运维的生产级系统。
一、设计思路
1.1 可观测性的三大支柱
可观测性 ├── Metrics(指标)— 数值化的系统度量 │ ├── 吞吐量:QPS、TPS │ ├── 延迟:P50/P95/P99 │ ├── 资源:CPU、内存、磁盘 │ └── 业务:数据量、请求数 │ ├── Logging(日志)— 结构化的事件记录 │ ├── 访问日志 │ ├── 错误日志 │ ├── 慢查询日志 │ └── 审计日志 │ └── Tracing(追踪)— 请求的全链路跟踪 ├── 分布式追踪 ├── 调用链分析 └── 依赖关系图1.2 监控架构
┌─────────────┐ ┌─────────────┐ ┌─────────────┐ │ MiniKV │ │ MiniKV │ │ MiniKV │ │ Node 1 │ │ Node 2 │ │ Node 3 │ ├─────────────┤ ├─────────────┤ ├─────────────┤ │ Metrics │ │ Metrics │ │ Metrics │ │ Exporter │ │ Exporter │ │ Exporter │ └──────┬──────┘ └──────┬──────┘ └──────┬──────┘ │ │ │ └──────────────────┼──────────────────┘ ▼ ┌─────────────────┐ │ Prometheus │ │ (聚合存储) │ └────────┬────────┘ │ ┌────────▼────────┐ │ Grafana │ │ (可视化面板) │ └─────────────────┘二、指标系统
2.1 核心指标定义
# minikv/monitor/metrics.py import time import threading from typing import Dict, List, Optional, Callable from collections import defaultdict from dataclasses import dataclass, field import json @dataclass class MetricValue: """指标值""" value: float timestamp: float labels: Dict[str, str] = field(default_factory=dict) @dataclass class HistogramBucket: """直方图桶""" le: float # 上限 count: int = 0 class Counter: """计数器:只增不减""" def __init__(self, name: str, help_text: str = "", labels: Dict[str, str] = None): self.name = name self.help = help_text self.labels = labels or {} self._value = 0 self.lock = threading.Lock() def inc(self, value: float = 1): with self.lock: self._value += value def get(self) -> float: with self.lock: return self._value def reset(self): with self.lock: self._value = 0 class Gauge: """仪表盘:可增可减""" def __init__(self, name: str, help_text: str = "", labels: Dict[str, str] = None): self.name = name self.help = help_text self.labels = labels or {} self._value = 0 self.lock = threading.Lock() def set(self, value: float): with self.lock: self._value = value def inc(self, value: float = 1): with self.lock: self._value += value def dec(self, value: float = 1): with self.lock: self._value -= value def get(self) -> float: with self.lock: return self._value class Histogram: """直方图:测量分布""" def __init__(self, name: str, help_text: str = "", buckets: List[float] = None, labels: Dict[str, str] = None): self.name = name self.help = help_text self.labels = labels or {} self.buckets = [HistogramBucket(le=b) for b in (buckets or [0.001, 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0])] self.total_count = 0 self.total_sum = 0.0 self.lock = threading.Lock() def observe(self, value: float): with self.lock: self.total_count += 1 self.total_sum += value for bucket in self.buckets: if value <= bucket.le: bucket.count += 1 def get_percentile(self, p: float) -> float: """获取百分位数""" with self.lock: if self.total_count == 0: return 0 target = int(self.total_count * p / 100) count = 0 for bucket in self.buckets: count += bucket.count if count >= target: return bucket.le return self.buckets[-1].le if self.buckets else 0 class MetricsRegistry: """指标注册中心""" def __init__(self): self.counters: Dict[str, Counter] = {} self.gauges: Dict[str, Gauge] = {} self.histograms: Dict[str, Histogram] = {} self.lock = threading.Lock() def counter(self, name: str, help_text: str = "", labels: Dict[str, str] = None) -> Counter: with self.lock: if name not in self.counters: self.counters[name] = Counter(name, help_text, labels) return self.counters[name] def gauge(self, name: str, help_text: str = "", labels: Dict[str, str] = None) -> Gauge: with self.lock: if name not in self.gauges: self.gauges[name] = Gauge(name, help_text, labels) return self.gauges[name] def histogram(self, name: str, help_text: str = "", buckets: List[float] = None, labels: Dict[str, str] = None) -> Histogram: with self.lock: if name not in self.histograms: self.histograms[name] = Histogram(name, help_text, buckets, labels) return self.histograms[name] def export_prometheus(self) -> str: """导出Prometheus格式""" lines = [] for counter in self.counters.values(): lines.append(f"# HELP {counter.name} {counter.help}") lines.append(f"# TYPE {counter.name} counter") labels_str = ",".join(f'{k}="{v}"' for k, v in counter.labels.items()) if labels_str: lines.append(f'{counter.name}{{{labels_str}}} {counter.get()}') else: lines.append(f'{counter.name} {counter.get()}') for gauge in self.gauges.values(): lines.append(f"# HELP {gauge.name} {gauge.help}") lines.append(f"# TYPE {gauge.name} gauge") labels_str = ",".join(f'{k}="{v}"' for k, v in gauge.labels.items()) if labels_str: lines.append(f'{gauge.name}{{{labels_str}}} {gauge.get()}') else: lines.append(f'{gauge.name} {gauge.get()}') for hist in self.histograms.values(): lines.append(f"# HELP {hist.name} {hist.help}") lines.append(f"# TYPE {hist.name} histogram") with hist.lock: for bucket in hist.buckets: lines.append(f'{hist.name}_bucket{{le="{bucket.le}"}} {bucket.count}') lines.append(f'{hist.name}_count {hist.total_count}') lines.append(f'{hist.name}_sum {hist.total_sum}') return "\n".join(lines) def export_json(self) -> dict: """导出JSON格式""" return { 'counters': {k: v.get() for k, v in self.counters.items()}, 'gauges': {k: v.get() for k, v in self.gauges.items()}, 'histograms': { k: { 'p50': v.get_percentile(50), 'p95': v.get_percentile(95), 'p99': v.get_percentile(99), 'count': v.total_count, 'sum': v.total_sum } for k, v in self.histograms.items() } }三、日志系统
3.1 结构化日志
# minikv/monitor/logging.py import logging import json import time import traceback from typing import Dict, Any, Optional from datetime import datetime from enum import Enum, auto class LogLevel(Enum): DEBUG = auto() INFO = auto() WARN = auto() ERROR = auto() FATAL = auto() class StructuredLogger: """ 结构化日志记录器 输出JSON格式的日志,便于日志收集和分析 """ def __init__(self, name: str, output_file: str = None): self.name = name self.logger = logging.getLogger(name) # 设置格式 formatter = logging.Formatter( '%(message)s' # 我们自己格式化 ) # 控制台输出 console_handler = logging.StreamHandler() console_handler.setFormatter(formatter) self.logger.addHandler(console_handler) # 文件输出 if output_file: file_handler = logging.FileHandler(output_file) file_handler.setFormatter(formatter) self.logger.addHandler(file_handler) self.logger.setLevel(logging.INFO) def _log(self, level: LogLevel, message: str, **kwargs): """记录结构化日志""" log_entry = { 'timestamp': datetime.utcnow().isoformat() + 'Z', 'level': level.name, 'logger': self.name, 'message': message, **kwargs } log_line = json.dumps(log_entry, default=str) level_map = { LogLevel.DEBUG: logging.DEBUG, LogLevel.INFO: logging.INFO, LogLevel.WARN: logging.WARNING, LogLevel.ERROR: logging.ERROR, LogLevel.FATAL: logging.CRITICAL } self.logger.log(level_map[level], log_line) def info(self, message: str, **kwargs): self._log(LogLevel.INFO, message, **kwargs) def warn(self, message: str, **kwargs): self._log(LogLevel.WARN, message, **kwargs) def error(self, message: str, exc_info: bool = False, **kwargs): if exc_info: kwargs['exception'] = traceback.format_exc() self._log(LogLevel.ERROR, message, **kwargs) def debug(self, message: str, **kwargs): self._log(LogLevel.DEBUG, message, **kwargs) def fatal(self, message: str, **kwargs): self._log(LogLevel.FATAL, message, **kwargs) class AccessLogger: """访问日志记录器""" def __init__(self, logger: StructuredLogger): self.logger = logger def log_request(self, method: str, path: str, status: int, duration: float, client_ip: str, **kwargs): self.logger.info( f"{method} {path} {status}", method=method, path=path, status=status, duration_ms=round(duration * 1000, 2), client_ip=client_ip, **kwargs ) class SlowQueryLogger: """慢查询日志记录器""" def __init__(self, threshold: float = 0.1): self.threshold = threshold self.logger = StructuredLogger("slow_query") def check(self, query: str, duration: float, **kwargs): if duration > self.threshold: self.logger.warn( f"Slow query: {query}", query=query, duration_ms=round(duration * 1000, 2), threshold_ms=round(self.threshold * 1000, 2), **kwargs )四、分布式追踪
4.1 追踪实现
# minikv/monitor/tracing.py import time import uuid import threading from typing import Dict, List, Optional from dataclasses import dataclass, field import json @dataclass class Span: """追踪跨度""" trace_id: str span_id: str parent_span_id: Optional[str] operation_name: str start_time: float end_time: float = 0 tags: Dict[str, str] = field(default_factory=dict) logs: List[dict] = field(default_factory=list) children: List['Span'] = field(default_factory=list) def finish(self): self.end_time = time.time() def duration(self) -> float: return (self.end_time - self.start_time) if self.end_time else 0 def log_event(self, event: str, **fields): self.logs.append({ 'timestamp': time.time(), 'event': event, 'fields': fields }) class Tracer: """ 分布式追踪器 支持跨进程的追踪上下文传递 """ def __init__(self, service_name: str): self.service_name = service_name self.spans: Dict[str, List[Span]] = {} self.lock = threading.Lock() def start_span(self, operation_name: str, parent_span: Span = None, trace_id: str = None) -> Span: """开始一个追踪跨度""" span_id = str(uuid.uuid4()) if parent_span: trace_id = parent_span.trace_id parent_id = parent_span.span_id else: trace_id = trace_id or str(uuid.uuid4()) parent_id = None span = Span( trace_id=trace_id, span_id=span_id, parent_span_id=parent_id, operation_name=operation_name, start_time=time.time(), tags={ 'service': self.service_name, 'span.kind': 'server' } ) with self.lock: if trace_id not in self.spans: self.spans[trace_id] = [] self.spans[trace_id].append(span) return span def inject(self, span: Span) -> dict: """将追踪上下文注入到载体(如HTTP头)""" return { 'trace_id': span.trace_id, 'span_id': span.span_id, 'service': self.service_name } def extract(self, carrier: dict) -> Optional[Span]: """从载体中提取追踪上下文""" if not carrier: return None return Span( trace_id=carrier.get('trace_id', ''), span_id=carrier.get('span_id', ''), parent_span_id=None, operation_name='extracted', start_time=time.time() ) def get_trace(self, trace_id: str) -> List[Span]: """获取完整的追踪链路""" with self.lock: return self.spans.get(trace_id, []).copy() def export_trace(self, trace_id: str) -> Optional[dict]: """导出追踪数据""" spans = self.get_trace(trace_id) if not spans: return None return { 'trace_id': trace_id, 'spans': [ { 'span_id': s.span_id, 'parent_span_id': s.parent_span_id, 'operation': s.operation_name, 'duration_ms': round(s.duration() * 1000, 2), 'tags': s.tags, 'logs': s.logs } for s in spans ] } class TraceContext: """追踪上下文(线程局部)""" _context = threading.local() @classmethod def get_current_span(cls) -> Optional[Span]: return getattr(cls._context, 'current_span', None) @classmethod def set_current_span(cls, span: Span): cls._context.current_span = span @classmethod def clear(cls): cls._context.current_span = None五、监控服务
5.1 集成监控
# minikv/monitor/service.py import time import threading import json from typing import Dict, Any, Optional from .metrics import MetricsRegistry, Counter, Gauge, Histogram from .logging import StructuredLogger, AccessLogger, SlowQueryLogger from .tracing import Tracer, TraceContext class MonitorService: """ 监控服务 整合指标、日志、追踪的统一入口 """ def __init__(self, node_id: str, enable_tracing: bool = True): self.node_id = node_id # 指标 self.metrics = MetricsRegistry() # 日志 self.logger = StructuredLogger(f"minikv.{node_id}") self.access_logger = AccessLogger(self.logger) self.slow_query_logger = SlowQueryLogger() # 追踪 self.tracer = Tracer(node_id) if enable_tracing else None # 预定义指标 self._init_default_metrics() # 后台采集 self._start_collector() def _init_default_metrics(self): """初始化默认指标""" # 请求相关 self.requests_total = self.metrics.counter( 'minikv_requests_total', 'Total requests', {'node': self.node_id} ) self.requests_active = self.metrics.gauge( 'minikv_requests_active', 'Active requests', {'node': self.node_id} ) self.request_duration = self.metrics.histogram( 'minikv_request_duration_seconds', 'Request duration', [0.001, 0.005, 0.01, 0.05, 0.1, 0.5, 1.0] ) # 存储相关 self.storage_keys = self.metrics.gauge( 'minikv_storage_keys', 'Number of stored keys', {'node': self.node_id} ) self.storage_size = self.metrics.gauge( 'minikv_storage_size_bytes', 'Storage size in bytes', {'node': self.node_id} ) # Raft相关 self.raft_leader_changes = self.metrics.counter( 'minikv_raft_leader_changes', 'Leader changes', {'node': self.node_id} ) self.raft_log_size = self.metrics.gauge( 'minikv_raft_log_size', 'Raft log size', {'node': self.node_id} ) # 系统资源 self.memory_usage = self.metrics.gauge( 'minikv_memory_bytes', 'Memory usage', {'node': self.node_id} ) def record_request(self, method: str, path: str, status: int, duration: float, client_ip: str = "unknown"): """记录请求""" self.requests_total.inc() self.request_duration.observe(duration) self.access_logger.log_request( method, path, status, duration, client_ip ) self.slow_query_logger.check(path, duration) def start_trace(self, operation: str, carrier: dict = None) -> Optional[object]: """开始追踪""" if not self.tracer: return None parent = None if carrier: parent = self.tracer.extract(carrier) span = self.tracer.start_span(operation, parent_span=parent) TraceContext.set_current_span(span) return span def end_trace(self, span, tags: dict = None): """结束追踪""" if span: if tags: span.tags.update(tags) span.finish() TraceContext.clear() def update_storage_metrics(self, keys: int, size: int): """更新存储指标""" self.storage_keys.set(keys) self.storage_size.set(size) def update_raft_metrics(self, log_size: int, is_leader: bool): """更新Raft指标""" self.raft_log_size.set(log_size) def _start_collector(self): """启动后台采集""" def collect(): while True: time.sleep(15) try: # 采集系统资源 import psutil process = psutil.Process() self.memory_usage.set(process.memory_info().rss) except ImportError: pass thread = threading.Thread(target=collect, daemon=True) thread.start() def get_metrics_prometheus(self) -> str: """获取Prometheus格式指标""" return self.metrics.export_prometheus() def get_metrics_json(self) -> dict: """获取JSON格式指标""" return self.metrics.export_json() def get_health_status(self) -> dict: """获取健康状态""" return { 'node_id': self.node_id, 'status': 'healthy', 'uptime': time.time() - self._start_time if hasattr(self, '_start_time') else 0, 'metrics': self.get_metrics_json() }六、管理API
6.1 HTTP管理接口
# minikv/monitor/admin_api.py import json import threading from http.server import HTTPServer, BaseHTTPRequestHandler from typing import Dict, Any, Optional class AdminAPIHandler(BaseHTTPRequestHandler): """管理API处理器""" server_instance = None # 由AdminServer设置 def do_GET(self): path = self.path.rstrip('/') if path == '/health': self._json_response(self.server_instance.get_health()) elif path == '/metrics': self._text_response(self.server_instance.get_metrics_prometheus()) elif path == '/metrics/json': self._json_response(self.server_instance.get_metrics_json()) elif path == '/status': self._json_response(self.server_instance.get_cluster_status()) elif path.startswith('/trace/'): trace_id = path.split('/')[-1] self._json_response(self.server_instance.get_trace(trace_id)) elif path == '/config': self._json_response(self.server_instance.get_config()) else: self._json_response({'error': 'Not found'}, 404) def do_POST(self): path = self.path.rstrip('/') content_length = int(self.headers.get('Content-Length', 0)) body = self.rfile.read(content_length) if content_length else b'{}' data = json.loads(body) if body else {} if path == '/config/update': self._json_response(self.server_instance.update_config(data)) elif path == '/debug/set_level': level = data.get('level', 'INFO') self._json_response(self.server_instance.set_log_level(level)) else: self._json_response({'error': 'Not found'}, 404) def _json_response(self, data: dict, status: int = 200): self.send_response(status) self.send_header('Content-Type', 'application/json') self.end_headers() self.wfile.write(json.dumps(data, indent=2).encode()) def _text_response(self, text: str, status: int = 200): self.send_response(status) self.send_header('Content-Type', 'text/plain') self.end_headers() self.wfile.write(text.encode()) def log_message(self, format, *args): pass # 抑制默认日志 class AdminServer: """ 管理服务器 提供HTTP管理接口 """ def __init__(self, monitor_service, kv_service, host: str = '127.0.0.1', port: int = 8080): self.monitor = monitor_service self.kv_service = kv_service self.host = host self.port = port self.server = None self.thread = None # 配置 self.config = { 'max_key_size': 1024, 'max_value_size': 1048576, 'slow_query_threshold': 0.1, 'log_level': 'INFO' } def start(self): """启动管理服务器""" AdminAPIHandler.server_instance = self self.server = HTTPServer((self.host, self.port), AdminAPIHandler) self.thread = threading.Thread(target=self.server.serve_forever, daemon=True) self.thread.start() print(f"Admin API running on http://{self.host}:{self.port}") def stop(self): if self.server: self.server.shutdown() def get_health(self) -> dict: """获取健康状态""" return self.monitor.get_health_status() def get_metrics_prometheus(self) -> str: """获取Prometheus指标""" return self.monitor.get_metrics_prometheus() def get_metrics_json(self) -> dict: """获取JSON指标""" return self.monitor.get_metrics_json() def get_cluster_status(self) -> dict: """获取集群状态""" if self.kv_service: raft_state = self.kv_service.raft.get_state() return { 'node_id': self.kv_service.node_id, 'role': raft_state['role'], 'term': raft_state['current_term'], 'log_size': raft_state['log_size'], 'commit_index': raft_state['commit_index'], 'data_size': len(self.kv_service.state_machine.data) } return {'error': 'KV service not available'} def get_trace(self, trace_id: str) -> dict: """获取追踪信息""" if self.monitor.tracer: trace = self.monitor.tracer.export_trace(trace_id) return trace or {'error': 'Trace not found'} return {'error': 'Tracing disabled'} def get_config(self) -> dict: """获取配置""" return self.config def update_config(self, updates: dict) -> dict: """更新配置""" self.config.update(updates) return {'success': True, 'config': self.config} def set_log_level(self, level: str) -> dict: """设置日志级别""" import logging level_map = { 'DEBUG': logging.DEBUG, 'INFO': logging.INFO, 'WARN': logging.WARNING, 'ERROR': logging.ERROR } log_level = level_map.get(level.upper(), logging.INFO) logging.getLogger().setLevel(log_level) self.config['log_level'] = level return {'success': True, 'level': level}七、完整演示
# examples/monitoring_demo.py import time import logging import sys import os import tempfile import threading import random logging.basicConfig(level=logging.INFO) sys.path.insert(0, '..') from minikv.kv.cluster import MiniKVCluster from minikv.monitor.service import MonitorService from minikv.monitor.admin_api import AdminServer def demo_metrics(): """演示指标收集""" print("=" * 90) print("📊 指标系统演示") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: cluster = MiniKVCluster( node_count=3, base_port=9960, data_dir=os.path.join(tmpdir, 'kv_data') ) client = cluster.start() leader_id = cluster.get_leader() leader_service = cluster.nodes[leader_id] # 创建监控服务 monitor = MonitorService(leader_id) # 模拟请求 print("\n📝 模拟请求...") for i in range(100): monitor.record_request( method='GET' if i % 2 == 0 else 'SET', path=f'/api/kv/key{i}', status=200, duration=random.uniform(0.001, 0.5) ) time.sleep(0.01) # 更新存储指标 monitor.update_storage_metrics(keys=1000, size=1024000) # 导出指标 print("\n📈 Prometheus 格式指标:") prometheus_output = monitor.get_metrics_prometheus() for line in prometheus_output.split('\n')[:15]: print(f" {line}") print(f" ... ({len(prometheus_output.split(chr(10)))} lines total)") print("\n📊 JSON 格式指标:") json_metrics = monitor.get_metrics_json() print(f" 请求总数: {json_metrics['counters']['minikv_requests_total']}") print(f" 请求延迟 P50: {json_metrics['histograms']['minikv_request_duration_seconds']['p50']}s") print(f" 请求延迟 P95: {json_metrics['histograms']['minikv_request_duration_seconds']['p95']}s") print(f" 请求延迟 P99: {json_metrics['histograms']['minikv_request_duration_seconds']['p99']}s") cluster.stop() def demo_tracing(): """演示分布式追踪""" print("\n" + "=" * 90) print("🔍 分布式追踪演示") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: cluster = MiniKVCluster( node_count=3, base_port=9970, data_dir=os.path.join(tmpdir, 'kv_data') ) client = cluster.start() leader_id = cluster.get_leader() leader_service = cluster.nodes[leader_id] monitor = MonitorService(leader_id, enable_tracing=True) # 模拟一个完整的请求链路 print("\n🔄 模拟请求链路...") # 开始根追踪 root_span = monitor.start_trace("handle_request") # 模拟几个步骤 step1 = monitor.start_trace("validate_input", carrier=monitor.tracer.inject(root_span)) time.sleep(0.01) monitor.end_trace(step1, {'valid': 'true'}) step2 = monitor.start_trace("process_data", carrier=monitor.tracer.inject(root_span)) time.sleep(0.02) monitor.end_trace(step2, {'processed': '100'}) step3 = monitor.start_trace("store_result", carrier=monitor.tracer.inject(root_span)) time.sleep(0.015) monitor.end_trace(step3, {'stored': 'true'}) monitor.end_trace(root_span, {'status': 'success'}) # 导出追踪 trace_data = monitor.tracer.export_trace(root_span.trace_id) print(f"\n📋 追踪链路:") print(f" Trace ID: {trace_data['trace_id']}") for span in trace_data['spans']: indent = " " if span['parent_span_id'] else "" print(f" {indent}▶ {span['operation']}: " f"{span['duration_ms']}ms") cluster.stop() def demo_admin_api(): """演示管理API""" print("\n" + "=" * 90) print("🌐 管理API演示") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: cluster = MiniKVCluster( node_count=3, base_port=9980, data_dir=os.path.join(tmpdir, 'kv_data') ) client = cluster.start() leader_id = cluster.get_leader() leader_service = cluster.nodes[leader_id] # 启动管理服务器 monitor = MonitorService(leader_id) admin = AdminServer(monitor, leader_service, port=19980) admin.start() print(f"\n📡 管理API已启动:") print(f" http://localhost:19980/health") print(f" http://localhost:19980/metrics") print(f" http://localhost:19980/status") print(f" http://localhost:19980/config") # 模拟一些请求 print("\n📝 模拟请求...") for i in range(50): client.set(f"perf_test:{i}", f"value_{i}") client.get(f"perf_test:{i}") # 获取状态 print("\n📊 集群状态:") status = admin.get_cluster_status() for key, value in status.items(): print(f" {key}: {value}") # 获取配置 print("\n⚙️ 当前配置:") config = admin.get_config() for key, value in config.items(): print(f" {key}: {value}") # 更新配置 print("\n🔄 更新配置...") result = admin.update_config({'slow_query_threshold': 0.05}) print(f" 结果: {result}") admin.stop() cluster.stop() def demo_structured_logging(): """演示结构化日志""" print("\n" + "=" * 90) print("📝 结构化日志演示") print("=" * 90) with tempfile.TemporaryDirectory() as tmpdir: log_file = os.path.join(tmpdir, 'minikv.log') logger = MonitorService("test-node") print("\n📋 日志输出:") logger.logger.info("System started", version="1.0.0", node="test-node", cluster_size=3) logger.logger.warn("High latency detected", endpoint="/api/kv/get", latency_ms=250, threshold_ms=100) try: raise ValueError("Connection timeout") except Exception as e: logger.logger.error("Operation failed", exc_info=True, operation="write", key="test-key") logger.logger.info("Request completed", method="GET", path="/api/kv/test", status=200, duration_ms=12.5, client_ip="192.168.1.100") print("\n (日志已输出到控制台,格式为JSON)") if __name__ == "__main__": demo_metrics() demo_tracing() demo_admin_api() demo_structured_logging()八、测试
# tests/test_monitor.py import unittest import time import threading from minikv.monitor.metrics import Counter, Gauge, Histogram, MetricsRegistry from minikv.monitor.tracing import Tracer, Span class TestMetrics(unittest.TestCase): """指标测试""" def test_counter(self): c = Counter('test_counter', 'Test counter') self.assertEqual(c.get(), 0) c.inc() self.assertEqual(c.get(), 1) c.inc(5) self.assertEqual(c.get(), 6) def test_gauge(self): g = Gauge('test_gauge', 'Test gauge') g.set(100) self.assertEqual(g.get(), 100) g.inc(10) self.assertEqual(g.get(), 110) g.dec(20) self.assertEqual(g.get(), 90) def test_histogram(self): h = Histogram('test_hist', 'Test histogram', buckets=[0.1, 0.5, 1.0]) for v in [0.05, 0.2, 0.3, 0.8, 1.5]: h.observe(v) self.assertEqual(h.total_count, 5) p50 = h.get_percentile(50) p95 = h.get_percentile(95) self.assertGreater(p95, p50) class TestTracing(unittest.TestCase): """追踪测试""" def test_basic_trace(self): tracer = Tracer("test-service") span = tracer.start_span("root_op") time.sleep(0.01) child = tracer.start_span("child_op", parent_span=span) time.sleep(0.005) child.finish() span.finish() trace = tracer.export_trace(span.trace_id) self.assertIsNotNone(trace) self.assertEqual(len(trace['spans']), 2) def test_trace_context(self): tracer = Tracer("test-service") span = tracer.start_span("test") carrier = tracer.inject(span) extracted = tracer.extract(carrier) self.assertEqual(extracted.trace_id, span.trace_id) self.assertEqual(extracted.span_id, span.span_id) if __name__ == "__main__": unittest.main()九、总结
这一讲我们为MiniKV构建了完整的监控运维体系:
组件 | 功能 |
|---|---|
指标系统 | Counter/Gauge/Histogram,Prometheus兼容 |
结构化日志 | JSON格式,支持等级和标签 |
分布式追踪 | 跨进程追踪,上下文传递 |
管理API | HTTP接口,健康检查/配置管理 |
监控服务 | 统一整合指标/日志/追踪 |
关键成果:
✅ 完整的可观测性三大支柱(指标、日志、追踪)
✅ Prometheus兼容的指标导出
✅ 分布式请求链路追踪
✅ HTTP管理接口
✅ 实时健康检查和配置热更新
下一讲:我们将实现性能优化与压力测试——对MiniKV进行全面性能评估和优化。
🧰开发之余的小工具推荐
处理 Base64、JWT 解析、JSON 格式化、Crontab 计算、PDF 合并压缩这些碎片需求,我常用一个纯前端本地工具箱:zz365.top(子页 PDF 大师:PDF 大师 - zz365工具箱)。所有计算在浏览器完成,文件不上服务器,关页即清。免费、无登录、无广告,适合开发者当常驻标签页。
