1. 引言:为什么说 Harness 比模型更持久
在 AI 应用快速迭代的今天,模型版本、算法框架、推理引擎都在不断演进。然而,真正决定一个 AI 系统能否长期稳定运行、能否被团队持续维护的,往往是那些看似不起眼的工程基础设施——数据管线、评估体系、监控告警、部署流程、回滚机制。这些统称为 Harness(工程约束框架)。
模型会过时,但 Harness 不会。一个设计良好的 Harness 可以在模型换了一代又一代之后依然稳定运转,成为团队最宝贵的工程资产。本文将从工程实践角度,深入剖析 Harness 的构成、设计原则,并给出完整的代码实战。
2. Harness 的核心构成
一个完整的 Harness 通常包含以下六大模块:
- 数据管线:数据的采集、清洗、标注、版本管理。
- 评估体系:离线评估、在线评估、回归测试。
- 监控告警:推理延迟、Token 消耗、错误率、漂移检测。
- 部署发布:灰度发布、A/B 测试、快速回滚。
- 配置管理:Prompt 版本、模型参数、特征开关。
- 可观测性:日志、追踪、指标聚合。
下面我们逐一展开,并给出可运行的代码示例。
3. 数据管线实战
数据是 AI 系统的燃料。一个健壮的数据管线需要保证数据的可追溯性和可重放性。以下是一个基于 Python 的数据版本管理示例:
import hashlib import json from datetime import datetime from pathlib import Path class DataVersionManager: """数据版本管理器:为每次数据变更生成不可变版本号""" def __init__(self, data_dir: str = "./data"): self.data_dir = Path(data_dir) self.data_dir.mkdir(exist_ok=True) self.manifest_path = self.data_dir / "manifest.json" self._load_manifest() def _load_manifest(self): if self.manifest_path.exists(): with open(self.manifest_path, "r") as f: self.manifest = json.load(f) else: self.manifest = {"versions": []} def _save_manifest(self): with open(self.manifest_path, "w") as f: json.dump(self.manifest, f, indent=2, ensure_ascii=False) def _compute_hash(self, content: str) -> str: return hashlib.sha256(content.encode("utf-8")).hexdigest()[:16] def register_version(self, dataset_name: str, content: str, description: str = ""): """注册一个新版本的数据集""" version_hash = self._compute_hash(content) version_id = f"{dataset_name}-{datetime.now().strftime('%Y%m%d%H%M%S')}-{version_hash}" # 保存数据文件 data_path = self.data_dir / f"{version_id}.json" with open(data_path, "w") as f: json.dump({"content": content, "hash": version_hash}, f, ensure_ascii=False) 更新清单 self.manifest["versions"].append({ "version_id": version_id, "dataset_name": dataset_name, "hash": version_hash, "description": description, "created_at": datetime.now().isoformat() }) self._save_manifest() return version_id def get_version(self, version_id: str) -> str: """按版本号读取数据""" data_path = self.data_dir / f"{version_id}.json" if not data_path.exists(): raise FileNotFoundError(f"版本 {version_id} 不存在") with open(data_path, "r") as f: record = json.load(f) 校验完整性 actual_hash = self._compute_hash(record["content"]) if actual_hash != record["hash"]: raise ValueError("数据完整性校验失败,文件可能被篡改") return record["content"] def list_versions(self, dataset_name: str = None): """列出所有版本""" versions = self.manifest["versions"] if dataset_name: versions = [v for v in versions if v["dataset_name"] == dataset_name] return versions 使用示例 if name == "main": manager = DataVersionManager() 注册 v1 数据 v1 = manager.register_version( "train_set", json.dumps([{"text": "你好,世界", "label": "greeting"}]), description="初始训练集" ) 注册 v2 数据(新增样本) v2 = manager.register_version( "train_set", json.dumps([ {"text": "你好,世界", "label": "greeting"}, {"text": "今天天气如何", "label": "weather"} ]), description="新增天气样本" ) print("所有版本:") for v in manager.list_versions("train_set"): print(f" {v['version_id']} - {v['description']}") 回滚读取 v1 content = manager.get_version(v1) print(f"回滚读取 v1 内容:{content}")</code></pre> 这个数据版本管理器解决了三个关键问题:数据不可变、可追溯、可回滚。每次变更都会生成新的版本号,旧版本永远不会被覆盖。 4. 评估体系实战 评估是 AI 工程的灵魂。没有评估,就无法判断模型是变好还是变坏。以下是一个完整的评估框架示例,支持离线评估和回归测试: import json import time from dataclasses import dataclass, field from typing import Callable, Dict, List, Optional @dataclass class EvalCase: """单个评估用例""" input_text: str expected_output: str metadata: Dict = field(default_factory=dict) @dataclass class EvalResult: """单个用例的评估结果""" case_id: str passed: bool score: float actual_output: str latency_ms: float error: Optional[str] = None class EvalSuite: """评估套件:管理用例集合并执行评估""" def init(self, name: str, cases: List[EvalCase]): self.name = name self.cases = cases self.results: List[EvalResult] = [] def run(self, model_fn: Callable[[str], str], scorer_fn: Callable[[str, str], float], threshold: float = 0.8) -> Dict: """ 执行评估 Args: model_fn: 模型推理函数,输入文本返回输出 scorer_fn: 评分函数,输入(期望输出, 实际输出)返回 0-1 分数 threshold: 通过阈值 """ self.results = [] start_time = time.time() for idx, case in enumerate(self.cases): case_start = time.time() try: actual = model_fn(case.input_text) score = scorer_fn(case.expected_output, actual) passed = score >= threshold error = None except Exception as e: actual = "" score = 0.0 passed = False error = str(e) latency = (time.time() - case_start) * 1000 self.results.append(EvalResult( case_id=f"{self.name}-{idx}", passed=passed, score=score, actual_output=actual, latency_ms=latency, error=error )) total_time = time.time() - start_time return self.summary(total_time) def summary(self, total_time: float) -> Dict: """生成评估摘要""" total = len(self.results) passed = sum(1 for r in self.results if r.passed) avg_score = sum(r.score for r in self.results) / total if total else 0 avg_latency = sum(r.latency_ms for r in self.results) / total if total else 0 return { "suite_name": self.name, "total_cases": total, "passed_cases": passed, "failed_cases": total - passed, "pass_rate": passed / total if total else 0, "avg_score": avg_score, "avg_latency_ms": avg_latency, "total_time_s": total_time, "failed_details": [ {"case_id": r.case_id, "score": r.score, "error": r.error} for r in self.results if not r.passed ] } def export_json(self, path: str): """导出评估结果到 JSON 文件""" with open(path, "w") as f: json.dump({ "suite_name": self.name, "results": [ { "case_id": r.case_id, "passed": r.passed, "score": r.score, "actual_output": r.actual_output, "latency_ms": r.latency_ms, "error": r.error } for r in self.results ] }, f, indent=2, ensure_ascii=False) 使用示例 if name == "main": 构造评估用例 cases = [ EvalCase("1+1等于几?", "2"), EvalCase("中国的首都是哪里?", "北京"), EvalCase("Python 中如何定义函数?", "def"), ] 模拟模型函数(实际中替换为真实模型调用) def mock_model(text: str) -> str: responses = { "1+1等于几?": "2", "中国的首都是哪里?": "北京", "Python 中如何定义函数?": "def" } return responses.get(text, "未知") 简单的精确匹配评分函数 def exact_match(expected: str, actual: str) -> float: return 1.0 if expected.strip() == actual.strip() else 0.0 suite = EvalSuite("核心能力测试", cases) summary = suite.run(mock_model, exact_match, threshold=0.8) print(f"通过率:{summary['pass_rate']:.0%}") print(f"平均分:{summary['avg_score']:.2f}") print(f"平均延迟:{summary['avg_latency_ms']:.1f}ms") 导出结果 suite.export_json("./eval_result.json")</code></pre> 这个评估框架的核心价值在于:它把评估从「临时脚本」升级为「可重复执行的工程资产」。每次模型更新后,都可以跑同一套用例,量化对比效果变化。 5. 监控告警实战 生产环境中的 AI 系统必须实时监控。以下示例展示如何构建一个轻量级的监控告警系统,覆盖延迟、错误率和输入漂移三个维度: import time import threading from collections import deque from dataclasses import dataclass, field from typing import Callable, Deque, Dict, List, Optional @dataclass class MetricPoint: """单个监控数据点""" timestamp: float value: float labels: Dict[str, str] = field(default_factory=dict) class MetricsCollector: """指标收集器:滑动窗口存储指标数据""" def init(self, window_size: int = 1000): self.window_size = window_size self._metrics: Dict[str, Deque[MetricPoint]] = {} self._lock = threading.Lock() def record(self, name: str, value: float, **labels): """记录一个指标点""" with self._lock: if name not in self._metrics: self._metrics[name] = deque(maxlen=self.window_size) self._metrics[name].append(MetricPoint( timestamp=time.time(), value=value, labels=labels )) def get_recent(self, name: str, seconds: float = 60) -> List[MetricPoint]: """获取最近 N 秒的指标""" with self._lock: if name not in self._metrics: return [] cutoff = time.time() - seconds return [p for p in self._metrics[name] if p.timestamp >= cutoff] def avg(self, name: str, seconds: float = 60) -> Optional[float]: """计算最近 N 秒的平均值""" points = self.get_recent(name, seconds) if not points: return None return sum(p.value for p in points) / len(points) class AlertRule: """告警规则定义""" def init(self, metric: str, operator: str, threshold: float, duration_seconds: int = 60, message: str = ""): self.metric = metric self.operator = operator # "gt", "lt", "gte", "lte" self.threshold = threshold self.duration_seconds = duration_seconds self.message = message def evaluate(self, value: Optional[float]) -> bool: """判断是否触发告警""" if value is None: return False if self.operator == "gt": return value > self.threshold elif self.operator == "gte": return value >= self.threshold elif self.operator == "lt": return value < self.threshold elif self.operator == "lte": return value <= self.threshold return False class Monitor: """监控告警系统主类""" def init(self, collector: MetricsCollector): self.collector = collector self.rules: List[AlertRule] = [] self.alert_handlers: List[Callable[[str], None]] = [] self._running = False self._thread = None def add_rule(self, rule: AlertRule): self.rules.append(rule) def add_alert_handler(self, handler: Callable[[str], None]): """注册告警处理函数(如发送邮件、钉钉通知)""" self.alert_handlers.append(handler) def _emit_alert(self, message: str): for handler in self.alert_handlers: try: handler(message) except Exception as e: print(f"告警处理失败:{e}") def check_rules(self): """检查所有规则""" for rule in self.rules: value = self.collector.avg(rule.metric, rule.duration_seconds) if rule.evaluate(value): msg = f"[告警] 指标 {rule.metric} 当前值 {value:.2f} " f"触发规则 {rule.operator} {rule.threshold}。{rule.message}" self._emit_alert(msg) def start(self, interval_seconds: float = 10.0): """启动后台监控线程""" if self._running: return self._running = True self._thread = threading.Thread( target=self._run_loop, args=(interval_seconds,), daemon=True ) self._thread.start() def _run_loop(self, interval_seconds: float): while self._running: self.check_rules() time.sleep(interval_seconds) def stop(self): self._running = False if self._thread: self._thread.join(timeout=2) 使用示例 if name == "main": collector = MetricsCollector() monitor = Monitor(collector) 定义告警规则 monitor.add_rule(AlertRule( metric="inference_latency_ms", operator="gt", threshold=500, duration_seconds=30, message="推理延迟过高,请检查模型服务" )) monitor.add_rule(AlertRule( metric="error_rate", operator="gt", threshold=0.05, duration_seconds=60, message="错误率超过 5%,请立即排查" )) 注册告警处理 def send_alert(message: str): print(f"[通知] {message}") 实际项目中可在此调用钉钉/企业微信/邮件 API monitor.add_alert_handler(send_alert) monitor.start(interval_seconds=5) 模拟生产流量 try: for i in range(100): 模拟正常请求 collector.record("inference_latency_ms", 100 + (i % 50)) collector.record("error_rate", 0.01) # 模拟异常请求 if i % 30 == 0: collector.record("inference_latency_ms", 800) collector.record("error_rate", 0.08) time.sleep(1) finally: monitor.stop()</code></pre> 这个监控系统虽然轻量,但已经具备生产可用的核心能力:滑动窗口指标存储、规则引擎、后台线程轮询、可扩展的告警通知。实际项目中可以替换为 Prometheus + AlertManager 等成熟方案。 6. 部署发布实战 AI 模型的发布必须支持灰度、可回滚。以下示例展示一个基于 Python 的灰度发布控制器: import random import time from dataclasses import dataclass from typing import Callable, Dict, Optional @dataclass class ModelVersion: """模型版本信息""" version_id: str model_fn: Callable[[str], str] weight: float = 0.0 # 流量权重 0-100 class TrafficRouter: """流量路由器:按权重分发请求到不同模型版本""" def init(self): self.versions: Dict[str, ModelVersion] = {} self._total_weight = 0 def register_version(self, version: ModelVersion): """注册模型版本""" self.versions[version.version_id] = version self._recalculate_weights() def _recalculate_weights(self): self._total_weight = sum(v.weight for v in self.versions.values()) def set_weight(self, version_id: str, weight: float): """动态调整版本权重(灰度发布核心)""" if version_id not in self.versions: raise KeyError(f"版本 {version_id} 不存在") self.versions[version_id].weight = weight self._recalculate_weights() def route(self, input_text: str) -> tuple: """ 根据权重路由请求 Returns: (version_id, output) """ if not self.versions: raise RuntimeError("没有注册任何模型版本") 按权重随机选择版本 roll = random.uniform(0, self._total_weight) cumulative = 0.0 selected_version = None for version in self.versions.values(): cumulative += version.weight if roll <= cumulative: selected_version = version break if selected_version is None: selected_version = list(self.versions.values())[-1] output = selected_version.model_fn(input_text) return selected_version.version_id, output def get_traffic_distribution(self) -> Dict[str, float]: """获取当前流量分布""" if self._total_weight == 0: return {} return { vid: v.weight / self._total_weight for vid, v in self.versions.items() } class DeploymentManager: """部署管理器:管理灰度发布流程""" def init(self, router: TrafficRouter): self.router = router self.deployment_history = [] def start_gray_release(self, new_version: ModelVersion, initial_weight: float = 5.0): """ 启动灰度发布 Args: new_version: 新模型版本 initial_weight: 初始流量权重(百分比) """ 注册新版本,初始权重很小 self.router.register_version(new_version) self.router.set_weight(new_version.version_id, initial_weight) self.deployment_history.append({ "action": "start_gray", "version_id": new_version.version_id, "weight": initial_weight, "timestamp": time.time() }) print(f"灰度发布启动:{new_version.version_id},初始权重 {initial_weight}%") def adjust_weight(self, version_id: str, weight: float): """调整灰度流量权重""" self.router.set_weight(version_id, weight) self.deployment_history.append({ "action": "adjust_weight", "version_id": version_id, "weight": weight, "timestamp": time.time() }) print(f"调整权重:{version_id} -> {weight}%") def rollback(self, version_id: str): """回滚到指定版本""" 将目标版本权重设为 100,其他版本设为 0 for vid in self.router.versions: self.router.set_weight(vid, 0.0) self.router.set_weight(version_id, 100.0) self.deployment_history.append({ "action": "rollback", "version_id": version_id, "timestamp": time.time() }) print(f"已回滚到版本:{version_id}") 使用示例 if name == "main": 定义两个模型版本 def model_v1(text: str) -> str: return f"[v1] 回答:{text}" def model_v2(text: str) -> str: return f"[v2] 回答:{text}(增强版)" router = TrafficRouter() manager = DeploymentManager(router) 初始部署 v1 v1 = ModelVersion("v1.0.0", model_v1, weight=100.0) router.register_version(v1) 灰度发布 v2 v2 = ModelVersion("v2.0.0", model_v2) manager.start_gray_release(v2, initial_weight=10.0) 模拟流量 for i in range(20): vid, output = router.route("你好") print(f"请求 {i+1}: 路由到 {vid} -> {output}") time.sleep(0.1) 观察流量分布 print(f"流量分布:{router.get_traffic_distribution()}") 逐步放量 manager.adjust_weight("v2.0.0", 50.0) manager.adjust_weight("v2.0.0", 100.0) 发现问题,回滚 manager.rollback("v1.0.0") print(f"回滚后流量分布:{router.get_traffic_distribution()}")</code></pre> 灰度发布的核心价值在于「风险可控」。新版本先接收小流量,观察指标后再逐步放量;一旦发现问题,可以秒级回滚到旧版本,把影响面降到最低。 7. 配置管理实战 Prompt 版本、模型参数、特征开关都属于配置。配置管理的关键是「变更可追溯、可回滚」。以下示例展示一个配置中心: import json import hashlib import threading from datetime import datetime from pathlib import Path from typing import Any, Callable, Dict, List, Optional class ConfigCenter: """配置中心:管理 Prompt 模板、模型参数等配置的版本化""" def init(self, config_dir: str = "./configs"): self.config_dir = Path(config_dir) self.config_dir.mkdir(exist_ok=True) self._lock = threading.Lock() self._configs: Dict[str, Dict] = {} self._listeners: List[Callable[[str, str], None]] = [] self._load_all() def _load_all(self): """加载所有配置文件""" for file in self.config_dir.glob("*.json"): key = file.stem with open(file, "r") as f: self._configs[key] = json.load(f) def _save(self, key: str): """保存配置到文件""" path = self.config_dir / f"{key}.json" with open(path, "w") as f: json.dump(self._configs[key], f, indent=2, ensure_ascii=False) def _compute_hash(self, config: Dict) -> str: content = json.dumps(config, sort_keys=True, ensure_ascii=False) return hashlib.sha256(content.encode("utf-8")).hexdigest()[:12] def set_config(self, key: str, config: Dict, description: str = "", author: str = ""): """ 更新配置(自动记录版本历史) """ with self._lock: old_hash = self._compute_hash(self._configs.get(key, {})) new_hash = self._compute_hash(config) if old_hash == new_hash: print(f"配置 {key} 无变化,跳过") return # 记录历史版本 history = self._configs.get(key, {}).get("_history", []) history.append({ "hash": old_hash, "config": self._configs.get(key, {}), "description": description, "author": author, "timestamp": datetime.now().isoformat() }) config["_history"] = history config["_current_hash"] = new_hash config["_updated_at"] = datetime.now().isoformat() config["_updated_by"] = author self._configs[key] = config self._save(key) 通知监听器 for listener in self._listeners: listener(key, new_hash) print(f"配置 {key} 已更新,新哈希:{new_hash}") def get_config(self, key: str) -> Dict: """获取当前配置""" with self._lock: config = self.configs.get(key, {}) 移除内部字段 return {k: v for k, v in config.items() if not k.startswith("")} def get_history(self, key: str) -> List[Dict]: """获取配置历史""" with self._lock: return self._configs.get(key, {}).get("_history", []) def rollback(self, key: str, history_index: int = -1): """回滚到历史版本""" with self._lock: history = self._configs.get(key, {}).get("_history", []) if not history: raise ValueError(f"配置 {key} 没有历史版本") target = history[history_index] target_config = target["config"] target_config["_history"] = history[:history_index] + history[history_index+1:] self._configs[key] = target_config self._save(key) print(f"配置 {key} 已回滚到 {target['timestamp']} 的版本") def watch(self, callback: Callable[[str, str], None]): """注册配置变更监听器""" self._listeners.append(callback) def export_snapshot(self, path: str): """导出全部配置快照""" with open(path, "w") as f: json.dump(self._configs, f, indent=2, ensure_ascii=False) 使用示例 if name == "main": center = ConfigCenter() 定义初始 Prompt 模板 initial_prompt = { "system": "你是一个专业的 AI 助手,请用简洁准确的语言回答问题。", "temperature": 0.7, "max_tokens": 512 } center.set_config("chat_prompt", initial_prompt, description="初始 Prompt 模板", author="alice") 更新 Prompt(优化提示词) optimized_prompt = { "system": "你是一个专业的 AI 助手。回答要求:1. 简洁准确 2. 分点列出 3. 给出示例", "temperature": 0.5, "max_tokens": 768 } center.set_config("chat_prompt", optimized_prompt, description="优化 Prompt:增加分点要求", author="bob") 查看历史 print("配置历史:") for h in center.get_history("chat_prompt"): print(f" {h['timestamp']} - {h['description']} (by {h['author']})") 回滚到初始版本 center.rollback("chat_prompt", 0) print(f"当前配置:{center.get_config('chat_prompt')}")</code></pre> 配置中心的本质是「把 Prompt 和参数当作代码来管理」。每一次修改都有记录、有作者、可回滚,这避免了「谁改了什么说不清」的混乱局面。 8. 可观测性实战 可观测性是 AI 系统的「仪表盘」。没有可观测性,生产事故就像在黑屋子里找东西。以下示例展示结构化日志和链路追踪: import json import time import uuid import threading from contextlib import contextmanager from dataclasses import dataclass, field from datetime import datetime from typing import Dict, List, Optional @dataclass class Span: """链路追踪中的单个跨度""" span_id: str parent_id: Optional[str] name: str start_time: float end_time: Optional[float] = None attributes: Dict = field(default_factory=dict) def finish(self): self.end_time = time.time() @property def duration_ms(self) -> float: if self.end_time is None: return 0.0 return (self.end_time - self.start_time) * 1000 def to_dict(self) -> Dict: return { "span_id": self.span_id, "parent_id": self.parent_id, "name": self.name, "duration_ms": self.duration_ms, "attributes": self.attributes } class Tracer: """链路追踪器:记录一次请求的完整调用链""" def init(self): self._local = threading.local() def start_span(self, name: str, **attributes) -> Span: """开启一个新的 span""" span_id = uuid.uuid4().hex[:12] parent_id = getattr(self._local, "current_span_id", None) span = Span( span_id=span_id, parent_id=parent_id, name=name, start_time=time.time(), attributes=attributes ) 保存到当前请求的 span 列表 if not hasattr(self._local, "spans"): self._local.spans = [] self._local.spans.append(span) 设置当前 span self._local.current_span_id = span_id return span def end_span(self, span: Span): """结束一个 span""" span.finish() 恢复父 span self._local.current_span_id = span.parent_id @contextmanager def span(self, name: str, **attributes): """上下文管理器方式使用 span""" span = self.start_span(name, **attributes) try: yield span finally: self.end_span(span) def get_trace(self) -> List[Dict]: """获取当前请求的完整调用链""" spans = getattr(self._local, "spans", []) return [s.to_dict() for s in spans] def clear(self): """清空当前请求的追踪数据""" self._local.spans = [] self._local.current_span_id = None class StructuredLogger: """结构化日志记录器:输出 JSON 格式日志""" def init(self, service_name: str, tracer: Optional[Tracer] = None): self.service_name = service_name self.tracer = tracer def _log(self, level: str, message: str, **context): record = { "timestamp": datetime.now().isoformat(), "level": level, "service": self.service_name, "message": message, **context } 附加链路追踪信息 if self.tracer: trace = self.tracer.get_trace() if trace: record["trace"] = trace record["trace_id"] = trace[0]["span_id"] print(json.dumps(record, ensure_ascii=False)) def info(self, message: str, **context): self._log("INFO", message, **context) def warning(self, message: str, **context): self._log("WARNING", message, **context) def error(self, message: str, **context): self._log("ERROR", message, **context) 使用示例 if name == "main": tracer = Tracer() logger = StructuredLogger("ai-gateway", tracer) def call_llm(prompt: str) -> str: """模拟 LLM 调用""" with tracer.span("llm.invoke", model="gpt-4", prompt_len=len(prompt)): time.sleep(0.05) # 模拟网络延迟 return "模拟回复" def process_request(user_input: str) -> str: """处理一次完整请求""" with tracer.span("request.process", user_id="u123"): logger.info("收到用户请求", input=user_input) 调用 LLM with tracer.span("llm.call"): response = call_llm(user_input) 后处理 with tracer.span("postprocess"): time.sleep(0.01) result = response.upper() logger.info("请求处理完成", output=result) return result 模拟一次请求 result = process_request("你好,请介绍一下自己") print(f"最终结果:{result}") 打印完整调用链 print("\n完整调用链:") for span in tracer.get_trace(): indent = " " * (1 if span["parent_id"] else 0) print(f"{indent}{span['name']} - {span['duration_ms']:.1f}ms")</code></pre> 可观测性的核心是「一次请求的全链路可视化」。当用户反馈「回答变慢了」,你可以通过追踪数据快速定位是网络问题、模型推理问题还是后处理问题。 9. 完整 Harness 集成示例 下面把以上模块整合成一个完整的 Harness 框架,展示它们如何协同工作: import time from typing import Callable, Dict, List, Optional class Harness: """ 完整 Harness 框架:整合数据、评估、监控、部署、配置、可观测性 """ def init(self, service_name: str = "ai-service"): self.service_name = service_name 初始化各模块 from config_center import ConfigCenter from monitor import MetricsCollector, Monitor, AlertRule from tracer import Tracer, StructuredLogger from eval_suite import EvalSuite, EvalCase from traffic_router import TrafficRouter, ModelVersion, DeploymentManager self.config = ConfigCenter() self.collector = MetricsCollector() self.monitor = Monitor(self.collector) self.tracer = Tracer() self.logger = StructuredLogger(service_name, self.tracer) self.router = TrafficRouter() self.deploy = DeploymentManager(self.router) 默认监控规则 self._setup_default_rules() def _setup_default_rules(self): """设置默认监控规则""" self.monitor.add_rule(AlertRule( metric="inference_latency_ms", operator="gt", threshold=1000, duration_seconds=60, message="推理延迟超过 1 秒" )) self.monitor.add_rule(AlertRule( metric="error_rate", operator="gt", threshold=0.05, duration_seconds=60, message="错误率超过 5%" )) def register_model(self, version_id: str, model_fn: Callable[[str], str], weight: float = 100.0): """注册模型版本""" self.router.register_version(ModelVersion(version_id, model_fn, weight)) self.logger.info("模型版本已注册", version_id=version_id) def invoke(self, input_text: str, user_id: str = "anonymous") -> str: """ 核心调用入口:带完整 Harness 能力的模型调用 包含:链路追踪、指标采集、路由分发、错误处理 """ start_time = time.time() with self.tracer.span("harness.invoke", user_id=user_id): try: 记录请求指标 self.collector.record("request_count", 1, user_id=user_id) # 路由到模型 with self.tracer.span("router.route"): version_id, output = self.router.route(input_text) # 记录成功指标 latency = (time.time() - start_time) * 1000 self.collector.record("inference_latency_ms", latency, version=version_id) self.collector.record("error_rate", 0.0) self.logger.info("推理成功", version_id=version_id, latency_ms=latency) return output except Exception as e: # 记录错误指标 self.collector.record("error_rate", 1.0) self.logger.error("推理失败", error=str(e)) raise def run_evaluation(self, cases: List[EvalCase], scorer_fn: Callable[[str, str], float]) -> Dict: """运行评估套件""" suite = EvalSuite(f"{self.service_name}-eval", cases) def model_fn(text: str) -> str: _, output = self.router.route(text) return output summary = suite.run(model_fn, scorer_fn) self.logger.info("评估完成", pass_rate=summary["pass_rate"], avg_score=summary["avg_score"]) return summary def start(self): """启动监控""" self.monitor.start(interval_seconds=10) self.logger.info("Harness 已启动", service=self.service_name) def stop(self): """停止监控""" self.monitor.stop() self.logger.info("Harness 已停止") 使用示例 if name == "main": 创建 Harness harness = Harness("demo-service") harness.start() 注册两个模型版本 def model_a(text: str) -> str: time.sleep(0.02) return f"[A] {text}" def model_b(text: str) -> str: time.sleep(0.05) return f"[B] {text}(增强版)" harness.register_model("v1.0.0", model_a, weight=100.0) harness.register_model("v2.0.0", model_b, weight=0.0) 灰度发布 v2 harness.deploy.start_gray_release( ModelVersion("v2.0.0", model_b), initial_weight=20.0 ) 模拟生产请求 for i in range(10): result = harness.invoke(f"测试请求 {i}") print(f"请求 {i}: {result}") time.sleep(0.5) 运行评估 from eval_suite import EvalCase cases = [ EvalCase("你好", "你好"), EvalCase("再见", "再见"), ] def exact_match(expected: str, actual: str) -> float: return 1.0 if expected in actual else 0.0 summary = harness.run_evaluation(cases, exact_match) print(f"评估通过率:{summary['pass_rate']:.0%}") harness.stop()</code></pre> 这个集成示例展示了 Harness 的核心价值:所有工程能力(追踪、监控、路由、评估)都透明地嵌入在调用链路中,业务代码只需要调用一个 invoke 方法,就能获得完整的工程保障。 10. 设计原则与最佳实践 基于以上实战,总结 Harness 设计的核心原则: 不可变性:数据、配置、模型版本一旦发布就不可修改,只能新增版本。这保证了任何时刻都可以精确复现历史状态。 可回滚性:任何变更都必须支持快速回滚。回滚不是「可选功能」,而是「必备能力」。 可观测性:每一次请求都要有完整的链路追踪和结构化日志。没有观测,就没有管理。 渐进式发布:新版本永远从小流量开始,逐步放量。永远不要一次性全量上线。 自动化评估:评估必须自动化、可重复。每次模型更新都要跑同一套回归用例。 11. 总结 模型是 AI 系统的「发动机」,Harness 是「底盘、仪表盘和安全带」。发动机可以随时更换,但底盘决定了这辆车能开多远、多稳。 在工程实践中,投入在 Harness 上的每一分精力都会在未来的无数次模型迭代中持续回报。数据版本管理、自动化评估、监控告警、灰度发布、配置中心、可观测性——这六大模块构成了 AI 工程的坚实底座。 希望本文的代码实战能帮助你构建属于自己的 Harness。记住:模型会过时,但 Harness 不会。