从零构建大模型管控系统:Harness设计模式与Python实战
1. 项目概述:从“Harness”的迷雾到清晰的实现路径
最近在技术社区里,“Harness”这个词的热度有点高,但讨论起来总感觉隔着一层纱。有人把它和AI智能体(Agent)混为一谈,有人觉得它是个神秘的框架,还有人直接问“大模型Harness是什么意思”。作为一个喜欢动手把概念落地的开发者,我决定不再空谈,而是用代码来“复盘”一次Harness的构建过程。这不仅仅是为了搞清楚Harness到底是什么,更是想探索在当下LLM(大语言模型)能力爆发的时代,我们如何设计一个既灵活又可靠的系统,来“驾驭”这些强大的、但有时又不太可控的模型能力。
简单来说,我理解的“Harness”,其核心目标就是对LLM或其他复杂AI组件进行安全、可控、可观测的封装与编排。你可以把它想象成给一匹充满力量但野性未驯的骏马(LLM)套上缰绳(Harness)、安上马鞍(工具接口)、配上地图(工作流),并让骑手(开发者或用户)能够清晰地知道马跑到了哪里、状态如何。它不同于一个追求完全自主行动的Agent,Harness更强调外部控制、流程的可预测性以及对每一步操作的“刹车”能力。这次实战,我就用Python,从零开始搭建了这样一个系统的核心骨架,过程中踩了不少坑,也总结出一些架构设计上的心得。
2. 核心架构设计:不是框架,而是一种设计模式
在动手写第一行代码之前,最关键的是厘清思路。我浏览了网络上很多关于“Harness vs Agent”的讨论,发现混淆点在于粒度。Agent通常指一个具备目标、能自主规划、调用工具去完成任务的智能体,它内部可能就运用了Harness的思想来管理其LLM核心。而Harness本身,更像是一种设计模式或中间件层,它的关注点不在于“自主性”,而在于“管控”。
2.1 设计目标与原则
我的Harness实现主要围绕以下几个目标展开:
- 隔离与稳定性:用户或上游系统不应直接、裸调LLM API。任何模型服务的抖动、超时、输出格式异常,都应在Harness层被消化和处理,向上提供稳定的接口。
- 可观测性:每一次对LLM的调用,其输入(Prompt)、输出(Response)、耗时、消耗的Token数、甚至中间步骤的思考过程,都必须被清晰记录。这是后期优化、审计和成本核算的基础。
- 流程编排:很多任务不是一次LLM调用就能解决的。可能需要链式调用(根据第一次输出构造第二次输入),或并行调用多个模型进行结果比对(民主投票模式)。Harness需要提供简洁的方式来定义这种工作流。
- 安全与合规:对输入输出进行过滤和检查,防止Prompt注入攻击,对输出内容进行合规性校验,避免产生有害内容。这层检查必须在返回给用户前完成。
- 灵活性:不应绑定任何特定的LLM提供商(如OpenAI、Anthropic、国内大模型)或特定框架(如LangChain)。它应该是一个轻量的粘合层。
基于这些目标,我放弃了直接使用功能庞大但有时显得笨重的全功能框架,决定采用一种“微内核+插件”的轻量级架构。
2.2 核心模块划分
最终,我将系统划分为五个核心模块,它们之间的依赖关系清晰,职责单一:
- 核心执行引擎 (Engine):这是大脑。它定义工作流的执行逻辑,是顺序执行、并行执行还是条件分支。它不关心具体做什么,只关心流程。
- 任务单元 (Task):这是骨骼和肌肉。一个Task封装了一次对LLM的完整调用,包括:准备Prompt、调用模型、解析输出、处理错误。它是可配置、可复用的基本单位。
- 上下文管理器 (Context):这是共享记忆。在整个工作流执行过程中,各个Task需要传递和共享数据。Context就是一个全局的、类型安全的键值存储,负责管理这些状态。
- 可观测性套件 (Observability):这是神经系统和仪表盘。它集成日志记录、指标收集(如延迟、Token消耗)和分布式追踪。每一个Task的执行痕迹都会被它捕获。
- 安全与验证层 (Guardrails):这是免疫系统。在Task的输入前和输出后,执行预定义的安全规则检查,例如检查输入中是否包含敏感词,输出格式是否符合JSON Schema等。
这个架构的好处是,每个模块都可以独立开发和测试。例如,我可以先实现一个最简单的顺序执行引擎,然后逐步替换为更复杂的DAG(有向无环图)引擎,而其他模块几乎不需要改动。
3. 关键代码实现与核心细节解析
理论说得再多,不如一行代码。我选择Python作为实现语言,因为它生态丰富,异步支持好,非常适合做这种胶水层开发。下面我拆解几个最核心的实现片段。
3.1 定义任务单元:从松散到严谨
最初,我简单地把一个Task定义为一个异步函数。但这很快带来了问题:配置如何传递?错误怎么统一处理?元数据(如本次调用的用途描述)如何附加?
重构后,我采用了基于Pydantic的数据类来定义Task。Pydantic提供了强大的数据验证和序列化能力,这让Task的配置管理变得非常清晰。
from pydantic import BaseModel, Field from typing import Any, Dict, Optional, Callable, Awaitable from enum import Enum class TaskStatus(Enum): PENDING = “pending” RUNNING = “running” SUCCESS = “success” FAILED = “failed” CANCELLED = “cancelled” class Task(BaseModel): “”“Harness系统中的基本执行单元。”“” id: str = Field(default_factory=lambda: str(uuid.uuid4())[:8]) name: str = Field(..., description=“任务名称,用于日志和追踪”) # 核心执行函数 execute_fn: Callable[[Dict[str, Any]], Awaitable[Any]] = Field(..., exclude=True) # 输入参数,将从上下文中解析 input_template: Optional[str] = None # 输出结果的上下文存储键 output_key: Optional[str] = None # 重试配置 retry_config: RetryConfig = Field(default_factory=RetryConfig) # 超时时间(秒) timeout: Optional[int] = 30 # 状态 status: TaskStatus = TaskStatus.PENDING result: Any = None error: Optional[str] = None started_at: Optional[datetime] = None finished_at: Optional[datetime] = None class Config: arbitrary_types_allowed = True async def run(self, context: “Context”) -> Any: “”“执行任务,并更新上下文。”“” self.status = TaskStatus.RUNNING self.started_at = datetime.now() try: # 1. 准备输入:从上下文渲染模板,或直接使用上下文 execution_input = self._prepare_input(context) # 2. 执行核心逻辑 self.result = await asyncio.wait_for( self.execute_fn(execution_input), timeout=self.timeout ) # 3. 处理输出:存储到上下文 if self.output_key: context.set(self.output_key, self.result) self.status = TaskStatus.SUCCESS return self.result except asyncio.TimeoutError: self.error = f“Task ‘{self.name}’ timed out after {self.timeout}s” self.status = TaskStatus.FAILED raise except Exception as e: self.error = str(e) self.status = TaskStatus.FAILED # 这里可以加入重试逻辑 raise finally: self.finished_at = datetime.now() # 触发可观测性事件 context.emit(“task_finished”, self)这个设计的精妙之处在于:
- 强类型与自描述:所有配置字段都有类型提示和描述,IDE支持好,减少了运行时错误。
- 状态管理内置:Task对象自己管理生命周期状态,外部只需调用
run方法。 - 与上下文解耦又关联:Task通过
context参数获取输入和存储输出,而不是直接操作全局变量,这使得Task可以被安全地复用和组合。 - 可观测性钩子:在
finally块中触发事件,让可观测性模块可以无侵入地收集数据。
3.2 构建上下文:共享状态的安全通道
上下文是连接各个Task的桥梁。它必须线程/协程安全,并且要能处理复杂对象。我实现了一个简单的Context类,内部使用字典存储,但提供了基于“键路径”的访问方式(如context.get(“user.profile.age”)),并集成了简单的事件发射器,用于可观测性。
class Context: def __init__(self, initial_data: Optional[Dict] = None): self._store = {} self._event_handlers = defaultdict(list) if initial_data: self._store.update(initial_data) def set(self, key: str, value: Any): “”“设置值,支持点分隔的路径,如 ‘a.b.c’。”“” keys = key.split(‘.’) d = self._store for k in keys[:-1]: d = d.setdefault(k, {}) d[keys[-1]] = value self.emit(“context_updated”, {“key”: key, “value”: value}) def get(self, key: str, default: Any = None) -> Any: “”“获取值,支持点分隔的路径。”“” try: keys = key.split(‘.’) value = self._store for k in keys: if not isinstance(value, dict): return default value = value.get(k) if value is None: return default return value except (KeyError, AttributeError): return default def emit(self, event: str, data: Any): “”“触发事件,用于可观测性。”“” for handler in self._event_handlers[event]: # 在实际项目中,这里应使用异步调度 handler(data)注意:在生产环境中,如果工作流非常复杂或需要持久化,这个内存中的上下文可能需要替换为更强大的状态管理后端,比如Redis。但作为核心模式,这个轻量实现已经足够清晰。
3.3 实现顺序执行引擎:工作流的基石
有了Task和Context,就可以组装它们了。我首先实现了一个最简单的顺序执行引擎。它接受一个Task列表,并按顺序执行。
class SequentialEngine: “”“顺序执行引擎。”“” def __init__(self): self.tasks = [] self.context = Context() def add_task(self, task: Task): self.tasks.append(task) async def run(self): “”“顺序执行所有任务。”“” execution_report = {“start_time”: datetime.now(), “tasks”: []} for task in self.tasks: task_report = {“task_id”: task.id, “name”: task.name} try: await task.run(self.context) task_report[“status”] = “success” task_report[“result”] = task.result except Exception as e: task_report[“status”] = “failed” task_report[“error”] = str(e) # 这里可以定义引擎级别的失败策略:继续或终止 # 例如:if task.critical: raise finally: task_report[“duration”] = (task.finished_at - task.started_at).total_seconds() if task.started_at and task.finished_at else None execution_report[“tasks”].append(task_report) execution_report[“end_time”] = datetime.now() return execution_report这个引擎虽然简单,但已经能处理很多线性业务流程。例如,“总结一篇长文”可以分解为:Task1(分段),Task2(并行总结各段),Task3(汇总各段总结)。
4. 实战演练:构建一个内容安全审核工作流
现在,让我们把各个部分组合起来,实现一个具体的场景:内容安全审核工作流。这个工作流模拟一个用户生成内容(UGC)平台,在发布前对文本进行多维度检查。
需求:用户提交一段文本,系统需要:
- 检查是否包含违禁词汇。
- 调用LLM判断内容的情感倾向(积极/消极/中立)和潜在风险(如仇恨、暴力)。
- 如果LLM认为有风险,则再调用一个专门的分类模型进行二次确认。
- 汇总所有结果,做出最终决策(通过、拒绝、需人工复核)。
4.1 定义具体任务
首先,我们创建四个具体的Task。为了简化,这里用模拟函数代替真实的API调用。
import asyncio import re async def check_prohibited_words(input_data: Dict) -> Dict: “”“任务1:违禁词检查。”“” text = input_data[“text”] prohibited_words = [“暴力”, “极端”] # 示例词库 found = [word for word in prohibited_words if word in text] return {“has_prohibited”: len(found) > 0, “words”: found} async def llm_sentiment_analysis(input_data: Dict) -> Dict: “”“任务2:LLM情感与风险分析。”“” text = input_data[“text”] # 模拟LLM API调用 await asyncio.sleep(0.5) # 模拟LLM返回一个结构化的结果 return { “sentiment”: “neutral”, # 模拟结果 “risk_level”: “medium” if “争议” in text else “low”, # 模拟逻辑 “reason”: “文本涉及争议性话题。” if “争议” in text else “无显著风险。” } async def dedicated_risk_classifier(input_data: Dict) -> Dict: “”“任务3:专用风险分类模型。”“” # 模拟一个更精确但更慢的分类模型 await asyncio.sleep(1) return {“is_high_risk”: False, “confidence”: 0.85} # 模拟结果 async def make_final_decision(input_data: Dict) -> Dict: “”“任务4:汇总决策。”“” context = input_data[“context”] # 注意,这里我们从input_data里拿到了context prohibited_result = context.get(“prohibited_check”) llm_result = context.get(“llm_analysis”) classifier_result = context.get(“risk_classifier”) decision = “pass” reason = [] if prohibited_result and prohibited_result[“has_prohibited”]: decision = “reject” reason.append(f“包含违禁词: {prohibited_result[‘words’]}”) elif llm_result and llm_result[“risk_level”] == “high”: decision = “reject” reason.append(f“LLM高风险判定: {llm_result[‘reason’]}”) elif llm_result and llm_result[“risk_level”] == “medium”: if classifier_result and classifier_result[“is_high_risk”]: decision = “reject” reason.append(“专用分类器确认高风险”) else: decision = “human_review” reason.append(“中风险,建议人工复核”) else: decision = “pass” reason.append(“无风险”) return {“final_decision”: decision, “reasons”: reason}4.2 组装工作流并执行
接下来,我们使用Harness的核心类来组装并执行这个工作流。
async def main(): # 1. 初始化引擎和上下文 engine = SequentialEngine() initial_context = {“text”: “这是一段包含争议话题的示例文本,但用词比较克制。”} context = Context(initial_context) # 2. 创建并添加任务 # 任务1:违禁词检查,结果存到 ‘prohibited_check’ task1 = Task( name=“prohibited_words_check”, execute_fn=check_prohibited_words, input_template=“text”, # 从上下文的 ‘text’ 键获取输入 output_key=“prohibited_check” ) # 任务2:LLM分析,结果存到 ‘llm_analysis’ task2 = Task( name=“llm_sentiment_analysis”, execute_fn=llm_sentiment_analysis, input_template=“text”, output_key=“llm_analysis” ) # 任务3:条件性任务 - 只有LLM分析结果为中风险时才执行 # 注意:这里我们需要一个更智能的引擎来处理条件逻辑。为了演示,我们先创建它,但由引擎决定是否执行。 # 在实际的DAG引擎中,我们可以定义依赖关系:task3 依赖于 task2 且 task2.output.risk_level == ‘medium’ task3 = Task( name=“dedicated_risk_classifier”, execute_fn=dedicated_risk_classifier, input_template=“text”, output_key=“risk_classifier” ) # 任务4:最终决策,它需要读取前面所有任务的结果 # 这里execute_fn需要能访问整个context,所以我们通过一个包装函数来实现 async def decision_maker(input_data: Dict): # 这个函数的input_data在Harness调用时会被传入context # 但我们希望decision_maker内部能直接使用context # 一种方法是在创建Task时,将context绑定进去(但这样不纯粹)。 # 更好的方式是:Harness在调用execute_fn时,除了渲染后的input,还把context作为额外参数传入。 # 为了简化演示,我们修改一下设计,让execute_fn接收一个包含‘data’和‘context’的字典。 pass # 具体实现见下方调整 # 让我们调整一下设计,让Task的execute_fn固定接收一个包含‘data’和‘harness_context’的字典。 # 重构Task.run方法中的 `execution_input = self._prepare_input(context)` 部分, # 使其返回 {“data”: rendered_input, “harness_context”: context}。 # 然后,我们重写decision_maker: async def decision_maker_wrapper(input_with_context: Dict): # input_with_context 是 {“data”: …, “harness_context”: context} harness_context = input_with_context[“harness_context”] return await make_final_decision({“context”: harness_context}) task4 = Task( name=“make_final_decision”, execute_fn=decision_maker_wrapper, # input_template 为空,因为决策不依赖特定输入,而是依赖整个上下文 output_key=“final_decision_report” ) # 由于我们当前的SequentialEngine不支持条件分支,我们先执行一个简化版:总是运行task1, task2, task4。 # 并将task3的逻辑以“是否需要在task4中触发”的形式处理。 # 更完善的实现需要一个支持DAG的引擎。 engine.add_task(task1) engine.add_task(task2) engine.add_task(task4) # 3. 执行引擎 report = await engine.run() # 4. 打印结果 print(“工作流执行报告:”) print(f“总耗时: {(report[‘end_time’] - report[‘start_time’]).total_seconds():.2f}秒”) for t in report[“tasks”]: print(f“ - [{t[‘status’]}] {t[‘name’]}: {t.get(‘result’, t.get(‘error’, ‘N/A’))}”) print(“\n最终上下文内容:”) import pprint pprint.pprint(engine.context._store) # 运行 if __name__ == “__main__”: asyncio.run(main())这个例子展示了Harness如何将复杂的、多步骤的LLM应用流程模块化、清晰化。每个Task职责单一,通过Context共享数据,引擎控制流程。虽然我们的顺序引擎还很简单,但已经勾勒出了核心模式。
5. 深入探讨:从顺序引擎到DAG引擎
上面的例子暴露了顺序引擎的局限性:无法处理条件分支和并行任务。在内容审核的例子中,专用分类器(task3)只有在LLM分析结果为“中风险”时才需要执行。这需要一个支持有向无环图(DAG)的引擎。
5.1 DAG引擎的设计思路
一个DAG引擎需要解决两个核心问题:
- 任务依赖定义:每个Task需要声明它依赖哪些其他Task的输出。
- 调度与执行:根据依赖关系,计算出可以并行执行的任务,并按拓扑顺序执行。
我们可以这样扩展Task和Engine:
class DAGTask(Task): “”“支持依赖关系的Task。”“” upstream_task_ids: List[str] = Field(default_factory=list) # 依赖的上游任务ID downstream_task_ids: List[str] = Field(default_factory=list) # 被哪些下游任务依赖 class DAGEngine: def __init__(self): self.tasks: Dict[str, DAGTask] = {} self.context = Context() self.graph = {} # 可以用 networkx 库,这里简化为邻接表 def add_task(self, task: DAGTask): self.tasks[task.id] = task def add_dependency(self, upstream_task_id: str, downstream_task_id: str): “”“添加依赖:downstream 依赖于 upstream。”“” if upstream_task_id in self.tasks and downstream_task_id in self.tasks: self.tasks[upstream_task_id].downstream_task_ids.append(downstream_task_id) self.tasks[downstream_task_id].upstream_task_ids.append(upstream_task_id) # 更新内部图结构... else: raise ValueError(“Task ID not found”) async def run(self): “”“执行DAG。”“” # 1. 拓扑排序,确定执行顺序 execution_order = self._topological_sort() # 2. 为每一层(可并行执行的任务)创建异步任务 for layer in execution_order: tasks_to_run = [self.tasks[task_id] for task_id in layer] # 使用 asyncio.gather 并行执行这一层的所有任务 await asyncio.gather(*[task.run(self.context) for task in tasks_to_run]) # 3. 返回报告...这样,在内容审核工作流中,我们可以建立依赖:task4依赖于task1,task2,task3。而task3依赖于task2的某个输出状态(risk_level == ‘medium’)。这需要引擎支持动态依赖或条件边,实现会更复杂,但模式是清晰的。
5.2 与现有框架的对比
你可能会问,这和LangChain或LlamaIndex的Chain、Agent有什么区别?我的体会是:
- LangChain:功能极其全面,但抽象层次高,有时“黑盒”感强,定制特定流程或深入调试时可能感觉笨重。它的很多概念(LCEL, Runnable)学习曲线不低。
- LlamaIndex:更专注于RAG(检索增强生成)领域,在该领域提供了非常优秀的抽象。但对于通用的、复杂的LLM工作流编排,可能不是最轻量的选择。
- 我们的Harness实现:更底层、更透明。它不试图提供开箱即用的100种工具和记忆模块,而是提供一套清晰的核心模式(Task, Context, Engine),让你可以像搭积木一样构建自己想要的任何流程,并且对每一环节都有完全的控制权和可见性。它更适合需要深度定制、对性能和可观测性有极高要求的场景。
6. 可观测性与生产级考量
一个玩具级的Harness和可用于生产的Harness,差距就在可观测性和健壮性上。
6.1 集成结构化日志与追踪
在生产环境中,打印print语句是远远不够的。我们需要结构化日志(如使用structlog或loggingJSON格式化),并集成分布式追踪(如OpenTelemetry)。
# 在Task.run方法中集成更完善的可观测性 async def run(self, context: “Context”, trace_id: str, span_id: str): import structlog logger = structlog.get_logger() with logger.bind(task_id=self.id, task_name=self.name, trace_id=trace_id, span_id=span_id): logger.info(“task.started”) try: # ... 执行逻辑 logger.info(“task.completed”, duration=…, token_usage=…) except Exception as e: logger.error(“task.failed”, error=str(e), exc_info=True) raise我们可以创建一个Observability模块,自动为每个Task调用注入日志器和追踪span。
6.2 实现重试与熔断机制
网络调用和第三方API失败是常态。我们的RetryConfig需要真正发挥作用。
from tenacity import retry, stop_after_attempt, wait_exponential, retry_if_exception_type class RetryConfig(BaseModel): max_attempts: int = 3 wait_multiplier: float = 1.0 # 指数退避的乘数 stop_after_delay: Optional[float] = 30.0 # 最大重试总时长 # 在Task内部,可以用tenacity装饰器包装execute_fn def _wrap_with_retry(fn, config: RetryConfig): @retry( stop=stop_after_attempt(config.max_attempts), wait=wait_exponential(multiplier=config.wait_multiplier), retry=retry_if_exception_type((TimeoutError, APIConnectionError)) # 只对特定异常重试 ) async def wrapped(*args, **kwargs): return await fn(*args, **kwargs) return wrapped6.3 安全层实现示例
安全层(Guardrails)可以作为Task执行前后的钩子(Hook)来实现。
class Guardrail: async def before_execution(self, task: Task, context: Context) -> Optional[Dict]: “”“执行前检查,返回非None则阻止执行并返回该结果。”“” # 例如:检查输入是否超长,是否包含明显恶意模式 input_text = context.get(task.input_template) if input_text and len(input_text) > 10000: return {“error”: “Input text too long”, “code”: “VALIDATION_ERROR”} return None async def after_execution(self, task: Task, context: Context, result: Any) -> Any: “”“执行后检查,可以修改或验证结果。”“” # 例如:确保LLM输出符合JSON Schema,过滤敏感信息 if task.output_key == “llm_analysis”: # 假设我们期望result是一个dict if not isinstance(result, dict): raise ValueError(“LLM output is not a dictionary”) # 对result中的某些字段进行内容过滤 if “reason” in result: result[“reason”] = filter_sensitive_content(result[“reason”]) return result然后,在Task.run方法中,分别在execute_fn调用前后插入guardrail.before_execution和guardrail.after_execution的调用。
7. 踩坑心得与最佳实践
在实现和迭代这个Harness原型的过程中,我总结了几条血泪教训:
Context的设计要早,且要稳定:Context是系统的血液流通系统,一旦定义好结构,后期修改成本很高。尽早确定是用扁平键值、嵌套对象还是支持路径查询。考虑是否需要版本兼容。
区分“业务错误”和“系统错误”:LLM返回内容不符合要求是“业务错误”(如未按指定格式输出),应该被捕获并作为Task失败的一种正常结果。网络超时、鉴权失败是“系统错误”,需要重试或熔断。在Harness层明确处理这两种错误,向上提供清晰的错误类型。
可观测性数据要结构化、可查询:不要只记录“成功了”或“失败了”。记录模型名称、Prompt模板版本、消耗的Token(Prompt+Completion)、耗时、输入输出的采样或指纹。这些数据对于成本优化、Prompt工程和性能调优至关重要。
Prompt模板管理是另一个大课题:我们的Harness处理了流程,但Prompt本身的管理(版本化、A/B测试、参数化)同样重要。可以考虑将Prompt模板也作为外部化配置,通过一个
PromptTemplate类来管理,其render方法接收Context并生成最终Prompt字符串。测试策略:Harness的每个部分都要易于测试。
- 单元测试:每个
execute_fn、每个Guardrail。 - 集成测试:测试整个工作流,可以使用Mock来模拟LLM API的返回,确保流程逻辑正确。
- 端到端测试:用真实但配额很小的API Key,对关键流程进行少量真实调用测试。
- 单元测试:每个
性能考量:如果Task之间没有依赖,一定要利用
asyncio.gather实现并行。对于IO密集型的LLM调用,并行能大幅降低总延迟。同时,考虑为耗时长的Task设置合理的超时,避免一个任务卡死整个流程。
这个用代码实现Harness的过程,本质上是一次对复杂软件系统“控制力”的追寻。它让我更深刻地理解到,在AI应用开发中,比追求“全自动”的智能更重要的,是构建“可预测、可观测、可干预”的可靠系统。这套模式不一定适合所有场景,但对于需要将LLM能力稳健、可控地集成到核心业务流程中的项目,提供一个清晰的架构起点和深刻的实践认知。
