当前位置: 首页 > news >正文

第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工具箱)。所有计算在浏览器完成,文件不上服务器,关页即清。免费、无登录、无广告,适合开发者当常驻标签页。

http://www.jsqmd.com/news/1397581/

相关文章:

  • 节气率40%-60%的气保焊改造方案
  • 游戏开发核心:逻辑帧与物理帧的深度解析与实战优化
  • AI巡检、自动巡检、智能巡检,到底差在哪?很多企业一开始就搞错了
  • 短视频高价回收陷阱丛生,深圳处置闲置黄金,实测合规线下回收门店 - 日常前沿快讯
  • Kali Linux渗透测试入门:10天从零搭建实验环境到独立实战
  • 行星齿轮非线性动力学分析与工程应用
  • 从“会写代码”到“能碰硬件”:GaryCLI + GaryProbe 在工业场景中的开发作用与前景
  • 推荐一下青岛EMC代理正规公司:升级 - 品牌推广大师
  • 绝区零日常不再熬夜:一条龙自动化工具从零到挂机的实战配置指南
  • 2026年郑州能做智慧燃气安全监测管理系统的公司有哪些?
  • 西北环线旅游攻略七日游,青甘7天跟团游,纯玩小团玩转盐湖戈壁,新手出行必读 - 跟我去旅游
  • SGS见证兰州兰石1000Nm3/h高效低成本PEM电解水制氢系统72小时工业性试验
  • SCI投稿状态全解析:从ADM、AE到Under Review,读懂编辑部“暗号”
  • 显卡驱动清理不彻底?用 DDU 把三大品牌显卡驱动的残留一次扫净
  • 三星硬盘维修工具包下载|SHTV 4.0.6与2.2版软件+多语言教程
  • 免费解锁Wand专业版:三步永久移除2小时限制,附手机远程控制实战指南
  • Ubuntu U盘无法识别?从硬件到内核的完整排查与修复指南
  • 微串口调试工具 集成| Modbus/CAN/DBC/蓝牙 SPP/波形/脚本
  • # 避坑指南|奢侈品回收,哪些话术是商家的小圈套 - 朝夕热点速报
  • Ubuntu 22.04手动编译安装最新版CMake完整指南
  • AI专著撰写必备!精选工具助你一键生成20万字专著,快速完成出版目标! - AI写论文
  • 抖音批量下载工具完整教程(30分钟上手)
  • 2026 年南昌砸墙|微挖机租赁|废品拆除回收服务问答 - LYL仔仔
  • 家用车节气门高频脏污的底层逻辑,正确养护方式分享!
  • 不动产租赁管理系统核心能力拆解:一套系统应具备的6大模块 - 资讯报道
  • Mapshaper 完整实测:把 300MB 卡顿地图变成秒开页面,只差一条命令
  • 证券买卖五档行情接口开发与优化实战
  • DriverStoreExplorer 驱动清理实战:闯过三关,轻松释放C盘数GB空间
  • 印象笔记 token 七天就过期,我写了五版程序让它自己换
  • 基于OpenClaw Agent框架构建智能内容分发系统,实现多平台自动化发布