上篇你给单个 Agent 装上了 Middleware——脱敏、限流、摘要、人审都能挂钩子。可当任务要「研究 + 写作 + 审稿」多种人设时,一个 create_agent 再怎么加中间件,文风仍容易打架。这篇讲 Multi Agent:外层 Orchestrator 编排,内层每个 Worker 用 LangChain create_agent 实现。。
小明用一个 Agent 写完整篇文章
学完 LangChain 支线前半段,小明自信满满:
agent = create_agent(
model,
tools=[search, ...],
system_prompt="你既是研究员又是写手又是编辑……",
middleware=[...],
)
他丢进去:「写一篇关于人工智能的科普文」。出来的稿子:开头像论文摘要,中间突然鸡汤,结尾又变说明书。Middleware 管得了 PII,管不了「一人分饰三角」。
小明去找老张:”那是不是又要退回原生 OpenAI chat.completions 手写多 Agent?”
老张说:”不用退。组织方式用 Orchestrator + Task 协议;每一环 Worker 继续用你熟悉的 create_agent——只是 tools=[]、换专属 system_prompt 和温度。原生写作助手是同一套流水线;本篇把它迁到 LangChain 技术栈。”
什么是 Multi Agent(LangChain 版)
用户输入
↓
[Orchestrator] 纯 Python 调度:派活 / 传 context / 重试 / 降级
↓
[Researcher] create_agent(tools=[], system=研究员, t=0.3) → research_output
↓
[Writer] create_agent(tools=[], system=写手, t=0.8) → draft
↓
[Editor] create_agent(tools=[], system=编辑, t=0.4) → 终稿
| 维度 |
单 create_agent |
多 Agent(本课) |
| 人设 |
一个 system_prompt |
每个 Worker 一个 create_agent |
| 大脑 |
一个 model |
各 Worker 可不同 temperature |
| 交接 |
埋在 messages |
Task.context 显式传 |
| 实现 |
LangChain |
仍是 LangChain + 外层编排 |
目录:
courseware-langchain/multi-agent-writing/
├── main.py
├── llm.py # ChatOpenAI 工厂
├── protocol.py # Task / TaskResult
└── agents/
├── orchestrator.py
├── researcher.py
├── writer.py
├── editor.py
└── worker_util.py # invoke → 取最终文本
完整实现:从工厂到流水线
1. llm.py:LangChain 模型工厂
"""多 Agent 写作助手:模型工厂(LangChain ChatOpenAI)。"""
from __future__ import annotations
import os
from dotenv import load_dotenv
from langchain_openai import ChatOpenAI
load_dotenv()
def get_chat_model(*, temperature: float = 0.7) -> ChatOpenAI:
return ChatOpenAI(
model=os.environ.get("OPENAI_MODEL", "qwen-plus"),
api_key=os.environ["OPENAI_API_KEY"],
base_url=os.environ.get("OPENAI_BASE_URL"),
temperature=temperature,
)
“换端改 .env,不改 Worker。”
2. protocol.py:Task / TaskResult(编排层,与框架无关)
"""Agent 间通信协议:Task / TaskResult。"""
from __future__ import annotations
import time
import uuid
from dataclasses import dataclass, field
from enum import Enum
from typing import Any
class TaskStatus(Enum):
DONE = "done"
FAILED = "failed"
@dataclass
class Task:
task_id: str
agent_name: str
instruction: str
context: dict[str, Any] = field(default_factory=dict)
created_at: float = field(default_factory=time.time)
@classmethod
def new(cls, agent_name: str, instruction: str, **context: Any) -> "Task":
return cls(
task_id=str(uuid.uuid4())[:8],
agent_name=agent_name,
instruction=instruction,
context=dict(context),
)
@dataclass
class TaskResult:
task_id: str
agent_name: str
status: TaskStatus
output: str
duration_ms: int = 0
metadata: dict[str, Any] = field(default_factory=dict)
@classmethod
def success(cls, task: Task, output: str, **metadata: Any) -> "TaskResult":
ms = int((time.time() - task.created_at) * 1000)
return cls(
task_id=task.task_id,
agent_name=task.agent_name,
status=TaskStatus.DONE,
output=output,
duration_ms=ms,
metadata=dict(metadata),
)
@classmethod
def failure(cls, task: Task, reason: str) -> "TaskResult":
ms = int((time.time() - task.created_at) * 1000)
return cls(
task_id=task.task_id,
agent_name=task.agent_name,
status=TaskStatus.FAILED,
output=reason,
duration_ms=ms,
)
@property
def ok(self) -> bool:
return self.status == TaskStatus.DONE
“Worker 不互相 import。只收 Task、回 TaskResult。research_output / draft 放进 context。”
3. worker_util.py:统一 invoke 取正文
"""Worker 共用:用 create_agent 跑一轮并取出最终文本。"""
from __future__ import annotations
from typing import Any
def invoke_agent_text(agent: Any, user_content: str) -> str:
result = agent.invoke({"messages": [{"role": "user", "content": user_content}]})
last = result["messages"][-1]
return getattr(last, "content", str(last)) or ""
“和快速入门篇同一调用姿势:invoke({\"messages\": [...]})。”
4. Researcher:create_agent + 低温
"""研究员 Worker:LangChain create_agent(无工具,专属 system_prompt)。"""
from __future__ import annotations
from langchain.agents import create_agent
from agents.worker_util import invoke_agent_text
from llm import get_chat_model
from protocol import Task, TaskResult
_SYSTEM = """\
你是一位经验丰富的内容研究员。
用户会给你一个写作主题,你的任务是:
1. 分析该主题的核心角度和背景
2. 整理出 5-8 个读者最需要了解的关键要点
3. 给出写作方向建议:受众定位、行文语气、重点突出哪些方面
直接输出分析内容,自然流畅地表达,不需要 JSON 格式。
"""
class ResearcherAgent:
def __init__(self) -> None:
self._agent = create_agent(
get_chat_model(temperature=0.3),
tools=[],
system_prompt=_SYSTEM,
)
def run(self, task: Task) -> TaskResult:
prompt = (
f"写作主题:{task.instruction}\n\n"
"请帮我分析这个主题,整理关键要点和写作建议。"
)
try:
output = invoke_agent_text(self._agent, prompt)
return TaskResult.success(task, output=output)
except Exception as e:
return TaskResult.failure(task, reason=str(e))
5. Writer:吃 research_output,高温创作
"""写作 Worker:基于研究报告出初稿。"""
from __future__ import annotations
from langchain.agents import create_agent
from agents.worker_util import invoke_agent_text
from llm import get_chat_model
from protocol import Task, TaskResult
_SYSTEM = """\
你是一位经验丰富的内容创作者,擅长将复杂知识写成清晰易懂的文章。
写作原则:
- 开篇用一个具体场景或问题抓住读者
- 逻辑清晰,每段只围绕一个核心观点展开
- 用类比和例子解释抽象概念,让普通读者也能看懂
- 结尾给读者留下明确的思考方向或行动建议
- 语言自然流畅,避免「首先、其次、最后」等套话
直接输出文章正文,不要加任何说明性前缀。
"""
class WriterAgent:
def __init__(self) -> None:
self._agent = create_agent(
get_chat_model(temperature=0.8),
tools=[],
system_prompt=_SYSTEM,
)
def run(self, task: Task) -> TaskResult:
research = task.context.get("research_output", "(无研究资料)")
prompt = f"""\
写作主题:{task.instruction}
以下是研究员整理的参考资料,请基于这些内容写出完整文章:
{research}
直接输出文章正文。
"""
try:
draft = invoke_agent_text(self._agent, prompt)
return TaskResult.success(task, output=draft)
except Exception as e:
return TaskResult.failure(task, reason=str(e))
6. Editor:一条建议 +【终稿】
"""编辑 Worker:一条改进建议 + 终稿。"""
from __future__ import annotations
from langchain.agents import create_agent
from agents.worker_util import invoke_agent_text
from llm import get_chat_model
from protocol import Task, TaskResult
_SYSTEM = """\
你是一位严格但高效的文章编辑。
你的工作分两步:
1. 给出【一条】最重要的改进建议(一句话,直接说问题在哪、怎么改)
2. 紧接着输出改进后的完整文章全文
输出格式:
【改进建议】<一句话建议>
【终稿】
<改进后的完整文章>
不要列多条建议,聚焦最关键的一条,然后在终稿中把它改好。
"""
class EditorAgent:
def __init__(self) -> None:
self._agent = create_agent(
get_chat_model(temperature=0.4),
tools=[],
system_prompt=_SYSTEM,
)
def run(self, task: Task) -> TaskResult:
draft = task.context.get("draft", "")
prompt = f"""\
文章主题:{task.instruction}
以下是初稿,请审核并改进:
{draft}
"""
try:
response = invoke_agent_text(self._agent, prompt)
final = _extract_final(response)
return TaskResult.success(task, output=final, full_response=response)
except Exception as e:
return TaskResult.failure(task, reason=str(e))
def _extract_final(editor_response: str) -> str:
marker = "【终稿】"
if marker in editor_response:
return editor_response[editor_response.index(marker) + len(marker):].strip()
return editor_response.strip()
7. Orchestrator:派活、重试、降级
"""调度员:Researcher → Writer → Editor。"""
from __future__ import annotations
import time
from agents.editor import EditorAgent
from agents.researcher import ResearcherAgent
from agents.writer import WriterAgent
from protocol import Task, TaskResult, TaskStatus
class OrchestratorAgent:
MAX_RETRIES = 2
def __init__(self) -> None:
self._researcher = ResearcherAgent()
self._writer = WriterAgent()
self._editor = EditorAgent()
def run(self, user_request: str) -> str:
print(f"\n[Orchestrator] 收到任务:{user_request}")
print("[Orchestrator] → 交给 Researcher...")
research_result = self._run_with_retry(
agent=self._researcher,
agent_name="researcher",
instruction=user_request,
)
if research_result.status == TaskStatus.FAILED:
return f"研究阶段失败:{research_result.output}"
print(f"[Orchestrator] Researcher 完成({research_result.duration_ms}ms)")
print("[Orchestrator] → 交给 Writer...")
write_result = self._run_with_retry(
agent=self._writer,
agent_name="writer",
instruction=user_request,
research_output=research_result.output,
)
if write_result.status == TaskStatus.FAILED:
return f"写作阶段失败:{write_result.output}"
print(f"[Orchestrator] Writer 完成({write_result.duration_ms}ms)")
print("[Orchestrator] → 交给 Editor...")
edit_result = self._run_with_retry(
agent=self._editor,
agent_name="editor",
instruction=user_request,
draft=write_result.output,
)
if edit_result.status == TaskStatus.FAILED:
print("[Orchestrator] Editor 失败,保留初稿。")
return write_result.output
print(f"[Orchestrator] Editor 完成({edit_result.duration_ms}ms)")
print("[Orchestrator] 全部流程完成。\n")
return edit_result.output
def _run_with_retry(
self,
agent: ResearcherAgent | WriterAgent | EditorAgent,
agent_name: str,
instruction: str,
**context_kwargs: str,
) -> TaskResult:
last_result: TaskResult | None = None
for attempt in range(1, self.MAX_RETRIES + 2):
task = Task.new(
agent_name=agent_name,
instruction=instruction,
**context_kwargs,
)
result = agent.run(task)
if result.status == TaskStatus.DONE:
return result
last_result = result
if attempt <= self.MAX_RETRIES:
print(f" [{agent_name}] 第{attempt}次失败:{result.output},1s 后重试...")
time.sleep(1)
else:
print(f" [{agent_name}] 已达最大重试次数,放弃。")
return last_result
8. main.py:入口
"""多 Agent 写作助手(LangChain create_agent 版)"""
import os
import sys
from agents.orchestrator import OrchestratorAgent
def main() -> None:
request = input("请输入写作需求(例如:写一篇关于 AI 的科普文章):\n> ").strip()
if not request:
print("需求不能为空。")
sys.exit(1)
article = OrchestratorAgent().run(request)
print("=" * 60)
print(article)
print("=" * 60)
if __name__ == "__main__":
main()
小明复述:”三个 Worker 都是 create_agent;调度还是 Orchestrator;上下文走 Task.context——LangChain 零件 + 多角色组织。”
“对。以后某一环要查资料,给那个 Worker 的 tools=[...] 挂上即可,不必拆掉流水线。”
总结
老张说:”第十一篇升级版:多 Agent 组织方法 + LangChain 实现零件。”
“三个核心理解:
-
Orchestrator + Worker —— 调度不写字;Worker 互不认识
-
Worker = create_agent ——
tools=[] + 专属 system + 不同 temperature
-
Task.context —— 显式传 research/draft;重试与编辑降级保证韧性”
LangChain 支线进度:
… → Middleware
↓
多 Agent 协作(本篇:LangChain 完整实现)
↓
Agentic RAG 全链路 …
小明说:”流水线用上 create_agent 了。下一篇知识库问答——Agentic RAG。”
“对。”
多 Agent 不是抛弃 LangChain,而是用多个 create_agent 扮演不同角色,再用协议把产物接起来。