一个分布式系统在生产环境中运行必须具备完善的可观测性——能够实时了解系统状态、诊断问题、分析性能。这一讲我们为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_factorydict) 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(leb) 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, defaultstr) 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}, methodmethod, pathpath, statusstatus, duration_msround(duration * 1000, 2), client_ipclient_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( fSlow query: {query}, queryquery, duration_msround(duration * 1000, 2), threshold_msround(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_factorydict) logs: List[dict] field(default_factorylist) children: List[Span] field(default_factorylist) 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_idtrace_id, span_idspan_id, parent_span_idparent_id, operation_nameoperation_name, start_timetime.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_idcarrier.get(trace_id, ), span_idcarrier.get(span_id, ), parent_span_idNone, operation_nameextracted, start_timetime.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(fminikv.{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_spanparent) 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(targetcollect, daemonTrue) 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() }六、管理API6.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, indent2).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(targetself.server.serve_forever, daemonTrue) self.thread.start() print(fAdmin 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(levellogging.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_count3, base_port9960, data_diros.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( methodGET if i % 2 0 else SET, pathf/api/kv/key{i}, status200, durationrandom.uniform(0.001, 0.5) ) time.sleep(0.01) # 更新存储指标 monitor.update_storage_metrics(keys1000, size1024000) # 导出指标 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_count3, base_port9970, data_diros.path.join(tmpdir, kv_data) ) client cluster.start() leader_id cluster.get_leader() leader_service cluster.nodes[leader_id] monitor MonitorService(leader_id, enable_tracingTrue) # 模拟一个完整的请求链路 print(\n 模拟请求链路...) # 开始根追踪 root_span monitor.start_trace(handle_request) # 模拟几个步骤 step1 monitor.start_trace(validate_input, carriermonitor.tracer.inject(root_span)) time.sleep(0.01) monitor.end_trace(step1, {valid: true}) step2 monitor.start_trace(process_data, carriermonitor.tracer.inject(root_span)) time.sleep(0.02) monitor.end_trace(step2, {processed: 100}) step3 monitor.start_trace(store_result, carriermonitor.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_count3, base_port9980, data_diros.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, port19980) 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(fperf_test:{i}, fvalue_{i}) client.get(fperf_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, version1.0.0, nodetest-node, cluster_size3) logger.logger.warn(High latency detected, endpoint/api/kv/get, latency_ms250, threshold_ms100) try: raise ValueError(Connection timeout) except Exception as e: logger.logger.error(Operation failed, exc_infoTrue, operationwrite, keytest-key) logger.logger.info(Request completed, methodGET, path/api/kv/test, status200, duration_ms12.5, client_ip192.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_spanspan) 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/HistogramPrometheus兼容结构化日志JSON格式支持等级和标签分布式追踪跨进程追踪上下文传递管理APIHTTP接口健康检查/配置管理监控服务统一整合指标/日志/追踪关键成果✅ 完整的可观测性三大支柱指标、日志、追踪✅ Prometheus兼容的指标导出✅ 分布式请求链路追踪✅ HTTP管理接口✅ 实时健康检查和配置热更新下一讲我们将实现性能优化与压力测试——对MiniKV进行全面性能评估和优化。开发之余的小工具推荐处理 Base64、JWT 解析、JSON 格式化、Crontab 计算、PDF 合并压缩这些碎片需求我常用一个纯前端本地工具箱zz365.top子页 PDF 大师PDF 大师 - zz365工具箱。所有计算在浏览器完成文件不上服务器关页即清。免费、无登录、无广告适合开发者当常驻标签页。