ex_fangwen.wen пре 1 месец
родитељ
комит
f49ddd7dd2
2 измењених фајлова са 626 додато и 0 уклоњено
  1. 396 0
      work_03/langgraph_practice.py
  2. 230 0
      work_04/langgraph_practice.py

+ 396 - 0
work_03/langgraph_practice.py

@@ -0,0 +1,396 @@
+from langgraph.graph import START, END, StateGraph, MessagesState
+from typing import TypedDict, Annotated
+from operator import add
+from langchain_openai import ChatOpenAI
+from dotenv import load_dotenv
+import os
+
+
+load_dotenv(override=True)
+deepseek_base_url = os.getenv("DEEPSEEK_BASE_URL")
+deepseek_base_key = os.getenv("DEEPSEEK_BASE_KEY")
+deepseek_base_name = os.getenv("DEEPSEEK_BASE_NAME")
+
+print("1. 三点===============================================================================================================================================================================================================")
+class CounterState1(TypedDict):
+    count: int
+
+def incr_node_1(state: CounterState1)-> dict:
+    count = state["count"] + 1
+    return {"count": count}
+
+def incr_node_2(state: CounterState1)-> dict:
+    count = state["count"] * 2
+    return {"count": count}
+
+build = StateGraph(CounterState1)
+build.add_node("incr_node_1", incr_node_1)
+build.add_node("incr_node_2", incr_node_2)
+
+build.add_edge(START, "incr_node_1")
+build.add_edge("incr_node_1", "incr_node_2")
+build.add_edge("incr_node_2", END)
+
+graph = build.compile()
+result = graph.invoke({"count": 1})
+print(result)
+
+
+print("2. 条件边===============================================================================================================================================================================================================")
+class CounterState2(TypedDict):
+    count: int
+
+def sub_node_1(state: CounterState2) -> dict:
+    count = state["count"] - 1
+    return {"count": count}
+
+def sub_node_2(state: CounterState2) -> dict:
+    count = state["count"] - 2
+    return {"count": count}
+
+def condition_edge(state: CounterState2) -> str:
+    if state["count"] > 5:
+        return "sub_node_1"
+    return "sub_node_2"
+
+build = StateGraph(CounterState2)
+
+build.add_node("sub_node_1", sub_node_1).add_node("sub_node_2", sub_node_2)
+
+build.add_edge(START, "sub_node_1")
+build.add_conditional_edges(
+    "sub_node_1",
+    condition_edge,
+    {"sub_node_1": "sub_node_1", "sub_node_2": "sub_node_2"}
+)
+build.add_edge("sub_node_2", END)
+
+graph = build.compile()
+result = graph.invoke({"count": 20})
+print(result)
+
+print("3. 覆盖===============================================================================================================================================================================================================")
+class OverwriteState(TypedDict):
+    arr:  list[str]
+
+def node_1(state: OverwriteState)->dict:
+    return {"arr":["a", "b"]}
+
+def node_2(state: OverwriteState)->dict:
+    return {"arr":["c", "d"]}
+
+build = StateGraph(OverwriteState)
+build.add_node("node_1", node_1).add_node("node_2", node_2)
+build.add_edge(START, "node_1")
+build.add_edge("node_1", "node_2")
+build.add_edge("node_2", END)
+
+graph = build.compile()
+result = graph.invoke({"arr":["e", "f"]})
+print(result)
+
+print("4. 追加===============================================================================================================================================================================================================")
+class AppendState(TypedDict):
+    arr: Annotated[list[str], add]
+
+def node_1(state: AppendState)->dict:
+    return {"arr":["aa", "bb"]}
+
+def node_2(state: AppendState)->dict:
+    return {"arr":["cc", "dd"]}
+
+build = StateGraph(AppendState)
+build.add_node("node_1", node_1).add_node("node_2", node_2)
+build.add_edge(START, "node_1")
+build.add_edge("node_1", "node_2")
+build.add_edge("node_2", END)
+
+graph = build.compile()
+result = graph.invoke({"arr":["ee", "ff"]})
+print(result)
+
+print("5. API===============================================================================================================================================================================================================")
+# llm = ChatOpenAI(
+#     base_url= deepseek_base_url,
+#     api_key= deepseek_base_key,
+#     model= deepseek_base_name
+# )
+
+# def chat_llm(state: MessagesState)->dict:
+#     response = llm.invoke("你是谁")
+#     return {"messages": [response]}
+
+# build = StateGraph(MessagesState)
+# build.add_node("chat", chat_llm)
+# build.add_edge(START, "chat")
+# build.add_edge("chat", END)
+
+# graph = build.compile()
+# result = graph.invoke({"messages": ["多学多问"]})
+# print(result)
+
+print("6. 自定义===============================================================================================================================================================================================================")
+
+class ReplaceableList(list):
+    def __init__(self, *args, replace: bool = False, **kwargs):
+        super().__init__(*args, **kwargs)
+        self.replace = replace
+
+def smart_merge(current: list, incom: list)->list:
+    if isinstance(incom, ReplaceableList) and getattr(incom, "replace", False):
+        return list(incom)
+    return current + incom
+
+class CostomState(TypedDict):
+    arr: Annotated[list[str], smart_merge]
+
+def append(state: CostomState)->dict:
+    return {"arr": ["b", "c"]}
+                 
+def overwrite(state: CostomState)->dict:
+    new_list = ReplaceableList(["reset_data"])
+    new_list.replace = True
+    return {"arr": new_list}
+
+build = StateGraph(CostomState)
+build.add_node("append", append)
+build.add_node("overwrite", overwrite)
+build.add_edge(START, "overwrite")
+build.add_edge("overwrite", "append")
+build.add_edge("append", END)
+
+graph = build.compile()
+result = graph.invoke({"arr": ["a"]})
+print(result)
+
+
+print("7. 子图===============================================================================================================================================================================================================")
+
+class DrinkState(TypedDict):
+    logs: Annotated[list[str], add]
+
+class MakeState(TypedDict):
+    sub_logs: Annotated[list[str], add]
+
+def add_coffee(state: MakeState) -> dict:
+    return {"sub_logs": ["添加咖啡"]}
+
+def add_sugar(state: MakeState) -> dict:
+    return {"sub_logs": ["添加粮"]}
+
+def do_make_coffee(state: MakeState) -> dict:
+    return {"sub_logs": ["咖啡制作中..."]}
+
+makeBuild = StateGraph(MakeState)
+makeBuild.add_node("add_coffee", add_coffee)
+makeBuild.add_node("add_sugar", add_sugar)
+makeBuild.add_node("do_make_coffee", do_make_coffee)
+makeBuild.add_edge(START, "add_coffee")
+makeBuild.add_edge("add_coffee", "add_sugar")
+makeBuild.add_edge("add_sugar", "do_make_coffee")
+makeBuild.add_edge("do_make_coffee", END)
+makeGraph = makeBuild.compile()
+
+def order_coffee(state: DrinkState) -> dict:
+    return {"logs": ["点咖啡"]}
+
+def make_coffee(state: DrinkState) -> dict:
+    result = makeGraph.invoke({"sub_logs": ["开始制作咖啡"]})
+    return {"logs": result["sub_logs"]}
+
+def drink_coffee(state: DrinkState) -> dict:
+    return {"logs": ["喝咖啡"]}
+
+build = StateGraph(DrinkState)
+build.add_node("order_coffee", order_coffee)
+build.add_node("make_coffee", make_coffee)
+build.add_node("drink_coffee", drink_coffee)
+build.add_edge(START, "order_coffee")
+build.add_edge("order_coffee", "make_coffee")
+build.add_edge("make_coffee", "drink_coffee")
+build.add_edge("drink_coffee", END)
+
+graph = build.compile()
+result = graph.invoke({"logs": ["喝咖啡"]})
+print(result)
+
+
+print("8. Send===============================================================================================================================================================================================================")
+"""
+图的拓扑结构:
+
+            START
+              │
+              ▼
+        [parse_query]            ← 节点:解析用户输入,提取品牌列表
+              │
+              ▼
+        route_to_research        ← 条件边路由函数:返回 N 个 Send
+              │
+              ├─ Send("research_brand", {brand: "苹果"})
+              ├─ Send("research_brand", {brand: "华为"})    ← 三个并行
+              └─ Send("research_brand", {brand: "小米"})
+              │
+              ▼  (并行执行)
+        [research_brand] × 3     ← 每个品牌独立研究
+              │
+              ▼  (全部完成后汇总)
+        [summarize]              ← 节点:汇总所有报告,输出对比
+              │
+              ▼
+             END
+"""
+from typing import Annotated, TypedDict
+from operator import add
+from langgraph.graph import StateGraph, START, END
+from langgraph.types import Send
+
+
+# ══════════════════════════════════════════════════════════════
+#  State 定义
+# ══════════════════════════════════════════════════════════════
+class CompareState(TypedDict):
+    """手机对比任务的 State"""
+    query: str                                # 用户的原始问题
+    brands: list[str]                         # 从 query 中解析出的品牌列表
+    reports: Annotated[list[str], add]        # 每个品牌的研究报告(add reducer 自动合并)
+    comparison: str                           # 最终的对比总结
+
+
+# ══════════════════════════════════════════════════════════════
+#  节点函数
+# ══════════════════════════════════════════════════════════════
+def parse_query(state: CompareState) -> dict:
+    """
+    节点 ①:解析用户输入,提取出要对比的品牌列表。
+
+    实际场景中这里会用 LLM 做实体提取,这里直接硬编码模拟。
+    """
+    query = state["query"]
+    print(f"[parse] 收到问题: {query}")
+
+    # 模拟 LLM 解析:从自然语言中提取品牌名称
+    brands = ["苹果", "华为", "小米"]
+    print(f"[parse] 解析出 {len(brands)} 个品牌: {brands}")
+
+    return {"brands": brands}
+
+
+def research_brand(state: CompareState) -> dict:
+    """
+    节点 ②:研究单个品牌(被 Send 并行调用的目标节点)。
+
+    它不知道总共有几个品牌,也不关心其他品牌在做什么——
+    只专注于自己拿到的这一个。
+
+    实际场景中这里会调搜索引擎 API 或 LLM,这里用数据模拟。
+    """
+    brand = state["brands"][0]  # Send 保证这里只有一个品牌
+
+    # 模拟:查询该品牌的旗舰机信息
+    phone_db = {
+        "苹果": "iPhone 16 Pro Max | A18 Pro 芯片 | 6.9吋 OLED | 4K 120fps 视频 | iOS 18",
+        "华为": "Mate 70 Pro+ | 麒麟 9100 | 6.8吋 OLED | 物理可变光圈 | 鸿蒙 NEXT",
+        "小米": "小米 15 Ultra | 骁龙 8 Gen4 | 6.73吋 AMOLED | 徕卡光学 | 澎湃 OS 2.0",
+    }
+
+    info = phone_db.get(brand, f"{brand}旗舰机信息暂缺")
+    report = f"【{brand}】{info}"
+
+    print(f"  [research] 并行查询 [{brand}] -> {info}")
+    return {"reports": [report]}  # add reducer 自动拼到总 reports 里
+
+
+def summarize(state: CompareState) -> dict:
+    """
+    节点 ③:所有品牌都研究完了,汇总生成对比结论。
+
+    实际场景中这里会把所有 reports 喂给 LLM,让它生成对比分析。
+    """
+    print(f"\n[summarize] 开始汇总 {len(state['reports'])} 份报告...")
+
+    # 模拟 LLM 生成的对比总结
+    comparison = (
+        "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n"
+        "         >> 旗舰手机对比总结\n"
+        "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n\n"
+    )
+    for report in state["reports"]:
+        comparison += f"  {report}\n"
+
+    comparison += (
+        "\n━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n"
+        "[总结] 三款旗舰各有千秋 --\n"
+        "  苹果 iPhone 16 Pro Max:视频拍摄王者,生态闭环体验最佳\n"
+        "  华为 Mate 70 Pro+:影像系统物理可变光圈独树一帜\n"
+        "  小米 15 Ultra:徕卡光学加持,性价比旗舰首选\n"
+        "━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━━\n"
+    )
+
+    return {"comparison": comparison}
+
+
+# ══════════════════════════════════════════════════════════════
+#  路由函数 —— 核心!返回 Send 列表而不是 dict
+# ══════════════════════════════════════════════════════════════
+def route_to_research(state: CompareState):
+    """
+    ★ Send API 的核心:条件边路由函数 ★
+
+    关键规则:
+    - 这个函数只能放在 add_conditional_edges 里,不能放进 add_node
+    - 普通节点返回 dict(更新 state)
+    - 路由函数返回 Send 列表(创建并行任务)或字符串(END / 节点名)
+
+    思想:几个品牌 → 发几个 Send → 启动几次 research_brand
+
+            return [
+                Send("research_brand", {"brands": ["苹果"]}),   ─┐
+                Send("research_brand", {"brands": ["华为"]}),   ─┤ 全部并行执行
+                Send("research_brand", {"brands": ["小米"]}),   ─┘
+            ]
+
+    每个 Send 的两个参数:
+      参数 1: 目标节点名 —— "派给谁做"
+      参数 2: 此分支专属的 state —— "这个任务要什么数据"
+    """
+    if not state["brands"]:
+        return END
+
+    return [
+        Send("research_brand", {"brands": [brand]})
+        for brand in state["brands"]
+    ]
+
+
+# ══════════════════════════════════════════════════════════════
+#  构图
+# ══════════════════════════════════════════════════════════════
+builder = StateGraph(CompareState)
+
+builder.add_node("parse_query", parse_query)
+builder.add_node("research_brand", research_brand)
+builder.add_node("summarize", summarize)
+
+# 边
+builder.add_edge(START, "parse_query")
+builder.add_conditional_edges("parse_query", route_to_research)  # 路由 → 动态 fan-out
+builder.add_edge("research_brand", "summarize")                  # 所有并行结果汇聚到 summarize
+builder.add_edge("summarize", END)
+
+graph = builder.compile()
+
+
+# ══════════════════════════════════════════════════════════════
+#  运行
+# ══════════════════════════════════════════════════════════════
+
+result = graph.invoke({
+    "query": "帮我对比一下苹果、华为和小米的旗舰机型",
+    "brands": [],
+    "reports": [],
+    "comparison": "",
+})
+print("======")
+print(result["brands"])
+print(result["comparison"])

+ 230 - 0
work_04/langgraph_practice.py

@@ -0,0 +1,230 @@
+import os
+import operator
+from typing import Annotated, List, Tuple, TypedDict, Union, Literal
+
+from langchain_community.tools.tavily_search import TavilySearchResults
+from langchain_community.chat_models import ChatTongyi
+from langchain_core.prompts import ChatPromptTemplate
+from langchain.agents import create_agent
+from pydantic import BaseModel, Field
+
+#https://tavily.com/
+#uv pip install tavily-python -i https://pypi.tuna.tsinghua.edu.cn/simple
+
+# ---- 配置区:请替换为你自己的 Key ----
+os.environ["TAVILY_API_KEY"] = "your-tavily-key-here"
+DASHSCOPE_API_KEY = "your-dashscope-key-here"
+
+# 初始化搜索工具,max_results 控制每次搜索返回的条数
+search_tool = TavilySearchResults(max_results=3)
+tools = [search_tool]
+
+# 规划阶段用推理能力强的 Qwen-Max
+planner_llm = ChatTongyi(
+    model="qwen-max",
+    temperature=0.1,          # 低温度保证规划稳定
+    api_key=DASHSCOPE_API_KEY,
+)
+
+# 执行阶段可以用更经济的模型,这里用 Qwen-Plus
+executor_llm = ChatTongyi(
+    model="qwen-plus",
+    temperature=0.01,
+    api_key=DASHSCOPE_API_KEY,
+)
+
+# 创建执行 Agent—— 它就是一个标准的 ReAct Agent
+executor_agent = create_agent(
+    model=executor_llm,
+    tools=tools,
+    system_prompt="你是一个执行力很强的助手。请根据给定的任务步骤,准确完成并给出结果。",
+)
+
+
+class PlanExecute(TypedDict):
+    """贯穿整个工作流的状态对象"""
+    input: str                                              # 用户原始输入
+    plan: List[str]                                         # 当前计划步骤列表
+    # past_steps 用 operator.add 做累加——每次节点返回新步骤时自动追加到历史
+    past_steps: Annotated[List[Tuple[str, str]], operator.add]
+    response: str                                           # 最终回复(非空时流程结束)
+
+
+class Plan(BaseModel):
+    """规划器输出的结构化计划"""
+    steps: List[str] = Field(
+        description="需要顺序执行的步骤列表,每一步必须独立且包含完整上下文"
+    )
+
+
+class Response(BaseModel):
+    """直接回复用户的内容"""
+    response: str
+
+
+class Act(BaseModel):
+    """
+    重规划器的决策输出——要么回复用户,要么给出更新后的计划。
+    用 Union 类型确保一次只走一条分支。
+    """
+    action: Union[Response, Plan] = Field(
+        description="如果任务已完成,返回 Response;如果还需要继续,返回 Plan"
+    )
+
+
+
+# 规划器的系统提示——核心是"拆得合理、不越界"
+planner_prompt = ChatPromptTemplate.from_messages([
+    (
+        "system",
+        """你是一个擅长任务分解的规划专家。对于用户提出的目标,请将其拆解为顺序执行的步骤列表。
+
+规则:
+1. 每一步必须自包含——执行者不需要上下文就能理解这步要做什么
+2. 不要添加无关步骤,最后一步的结果应该直接导向最终答案
+3. 优先使用中文描述步骤
+4. 步骤数量控制在 3~6 步,过少说明拆分不到位,过多说明混入了无关操作""",
+    ),
+    ("placeholder", "{messages}"),
+])
+
+# 将提示词模板和 LLM 拼成一条链,with_structured_output 让 LLM 直接输出 Plan 对象
+planner_chain = planner_prompt | planner_llm.with_structured_output(Plan)
+
+
+def plan_step(state: PlanExecute) -> dict:
+    """规划节点:分析用户原始输入,生成执行计划"""
+    result = planner_chain.invoke({"messages": [("user", state["input"])]})
+    return {"plan": result.steps}
+
+
+result = planner_chain.invoke({
+    "messages": [(
+        "user",
+        "帮我调研一下2025年国内主流代码助手工具(通义灵码、文心快码、CodeGeeX),"
+        "从代码补全准确率、支持语言、价格三个维度做对比,最后给出选型建议。"
+    )]
+})
+
+for i, step in enumerate(result.steps, 1):
+    print(f"  Step {i}: {step}")
+
+
+def execute_step(state: PlanExecute) -> dict:
+    """
+    从当前计划中取出第一步,交给执行 Agent 去跑。
+    执行完后把 (步骤描述, 执行结果) 追加到 past_steps。
+    """
+    plan = state["plan"]
+
+    # 把完整计划格式化,让 Agent 知道"我在做什么、后面还有什么"
+    plan_overview = "\n".join(f"  {i+1}. {step}" for i, step in enumerate(plan))
+    current_task = plan[0]  # 只取第一步
+
+    # 拼装给 Agent 的指令:先展示全貌,再指定当前任务
+    task_prompt = f"""以下是完整的执行计划:
+{plan_overview}
+
+现在请你只执行第 1 步,不要做后续步骤:
+  → {current_task}
+
+请直接给出这一步的执行结果,不要啰嗦。"""
+
+    agent_result = executor_agent.invoke({
+        "messages": [("user", task_prompt)]
+    })
+
+    # 取 Agent 最后一条消息的内容作为这一步的产出
+    step_output = agent_result["messages"][-1].content
+
+    return {
+        # 这一步走完后,计划需要弹出已执行的第一项——这个逻辑在 replan_step 里
+        "past_steps": [(current_task, step_output)],
+    }
+
+
+replanner_prompt = ChatPromptTemplate.from_template("""
+你是一个项目进度管控专家。根据以下信息,判断当前任务的状态并做出决策。
+
+## 用户原始目标
+{input}
+
+## 原始计划
+{plan}
+
+## 已完成的步骤及结果
+{past_steps}
+
+## 决策规则
+1. 如果所有步骤已完成且结果充分,直接给出最终答案(使用 Response)
+2. 如果还有步骤未执行,返回更新后的计划(使用 Plan),注意:
+   - 只保留尚未执行的步骤
+   - 如果已完成步骤的结果表明原计划的后续步骤需要调整,直接修改
+   - 不要返回已经做过的步骤
+3. 如果某个步骤执行失败了,请在计划中补充重试或替代方案
+""")
+
+replanner_chain = replanner_prompt | planner_llm.with_structured_output(Act)
+
+
+def replan_step(state: PlanExecute) -> dict:
+    """审视执行进度,决定下一步走向"""
+    output = replanner_chain.invoke(state)
+
+    if isinstance(output.action, Response):
+        # 任务完成,直接返回给用户
+        return {"response": output.action.response}
+    else:
+        # 还有步骤要执行,更新计划列表
+        return {"plan": output.action.steps}
+    
+
+
+from langgraph.graph import StateGraph, START, END
+
+# ---- 第一步:创建状态图,绑定状态类型 ----
+workflow = StateGraph(PlanExecute)
+
+# ---- 第二步:注册三个核心节点 ----
+workflow.add_node("planner", plan_step)       # 规划节点
+workflow.add_node("executor", execute_step)   # 执行节点
+workflow.add_node("replanner", replan_step)   # 重规划节点
+
+# ---- 第三步:编排边的流向 ----
+workflow.add_edge(START, "planner")           # 入口 → 规划
+workflow.add_edge("planner", "executor")      # 规划 → 执行
+workflow.add_edge("executor", "replanner")    # 执行 → 重规划
+
+# ---- 第四步:条件分支——重规划后走哪条路 ----
+def should_continue(state: PlanExecute) -> Literal["executor", "__end__"]:
+    """如果 response 字段非空,说明重规划器认为任务完成,结束流程"""
+    if state.get("response"):
+        return END
+    return "executor"
+
+workflow.add_conditional_edges("replanner", should_continue, {
+    "executor": "executor",
+    "__end__": END,
+})
+
+
+
+inputs = {
+    "input": (
+        "帮我调研2025年国内三大代码助手工具(通义灵码、文心快码、CodeGeeX),"
+        "从代码补全准确率、支持语言、价格三个维度做对比,最后给出选型建议。"
+        "请用中文输出最终结果。"
+    )
+}
+
+# stream 模式可以看到每一步的中间输出
+for event in app.stream(inputs):
+    for node_name, node_output in event.items():
+        print(f"\n{'='*50}")
+        print(f"📍 当前节点: {node_name}")
+        # 避免打印太长的内容
+        for key, value in node_output.items():
+            if isinstance(value, str) and len(value) > 200:
+                print(f"  {key}: {value[:200]}...")
+            else:
+                print(f"  {key}: {value}")