LangGraph 系列 10(下):Workflow 的 Edge、路由与 Studio 完整案例


系列 10(上)中,我们已经介绍了 Workflow 的 State、Reducer 和 Node:State 保存共享数据,Reducer 决定字段如何合并,Node 负责读取状态并提交局部更新。

本篇继续解决“节点按什么顺序执行”这个问题,重点介绍固定 Edge、并行 Fan-out、等待式 Fan-in、条件路由,以及图的构建、编译和流式执行。最后会把这些知识组合成一个城市旅行方案 Workflow,通过 langgraph dev 部署到本地 Agent Server,并在 LangGraph Studio 中观察完整执行过程。

本文继续使用同一组示例代码:

llm-learning/langgraph/p10_workflow_basics/

1. Edge:控制节点执行顺序

Edge 负责描述节点之间的执行关系。

1.1 START、END 与固定边

最简单的是固定执行路径:

# 入口:把初始 State 交给 normalize_input。
builder.add_edge(START, "normalize_input")
# 固定顺序:normalize_input 完成后一定执行 validate_input。
builder.add_edge("normalize_input", "validate_input")
# 出口:assemble_plan 完成后结束图。
builder.add_edge("assemble_plan", END)
  • START 是虚拟开始节点,把图输入发送给第一个业务 Node。
  • END 是虚拟结束节点,表示本次运行结束。
  • add_edge("a", "b") 表示 a 完成后固定执行 b

STARTEND 本身不会执行 Python 业务逻辑。

1.2 静态 Fan-out

同一个 Node 可以连接多个普通后继:

# 同一来源连接两个后继,形成静态 Fan-out。
# 两个 Node 会在下一个 Super-step 中执行。
builder.add_edge("prepare", "build_itinerary")
builder.add_edge("prepare", "build_packing_list")

prepare 完成后,两个节点会在下一个 Super-step 中执行。这种从一个节点分发到多个节点的结构叫 Fan-out。

1.3 使用列表来源完成 Fan-in

如果汇总节点必须等待多个来源全部完成,可以把来源节点列表传给 add_edge()

# 列表形式的来源表示:两个 Node 全部完成后才执行 summarize。
builder.add_edge(
    ["build_itinerary", "build_packing_list"],
    "summarize",
)

只有列表中的两个节点都完成后,summarize 才会执行。这种把多个分支重新汇聚到一个节点的结构叫 Fan-in。

并行分支在同一 Super-step 执行并通过 Fan-in 汇聚

1.4 Super-step 与并行更新

可以把一次图执行理解为多个 Super-step:

Super-step 1: prepare
Super-step 2: build_itinerary + build_packing_list
Super-step 3: summarize

同一个 Super-step 中的多个节点可能同时更新 State。如果它们更新同一个 Channel,该字段就必须有可以合并多个更新的 Reducer。

完整示例位于 04_edge_patterns.py

"""演示固定 Edge、静态 Fan-out 和列表来源 Fan-in。

Edge 决定 Node 之间的执行关系。固定 Edge 总是沿指定方向执行,
而同一来源的多个固定 Edge 会把工作流拆成可并行执行的分支。
"""

import operator
from typing import Annotated

from langgraph.graph import END, START, StateGraph
from typing_extensions import NotRequired, TypedDict


class EdgeState(TypedDict):
    city: str
    # 并行分支会同时写 suggestions 和 steps,因此必须声明 Reducer。
    suggestions: Annotated[list[str], operator.add]
    steps: Annotated[list[str], operator.add]
    summary: NotRequired[str]


def prepare(state: EdgeState) -> dict:
    """入口后的准备 Node:清理城市名称。"""

    return {"city": state["city"].strip(), "steps": ["prepare"]}


def build_itinerary(state: EdgeState) -> dict:
    """并行分支一:生成景点建议。"""

    return {
        "suggestions": [f'游览{state["city"]}的代表景点'],
        "steps": ["build_itinerary"],
    }


def build_packing_list(state: EdgeState) -> dict:
    """并行分支二:生成行李建议。"""

    return {
        "suggestions": ["准备舒适的步行鞋"],
        "steps": ["build_packing_list"],
    }


def summarize(state: EdgeState) -> dict:
    """Fan-in Node:读取两个分支合并后的 State。"""

    # 并行 Node 的完成顺序不应作为业务顺序,因此汇总前显式排序。
    suggestions = sorted(state["suggestions"])
    return {
        "summary": ";".join(suggestions),
        "steps": ["summarize"],
    }


builder = StateGraph(EdgeState)
builder.add_node(prepare)
builder.add_node(build_itinerary)
builder.add_node(build_packing_list)
builder.add_node(summarize)

# START 和 END 是特殊节点,分别表示图的入口和出口。
# START → prepare 是一条固定 Edge。
builder.add_edge(START, "prepare")

# 同一个来源连接两个后继节点,形成静态 Fan-out。
# 两个分支读取 prepare 完成后的同一份 State 快照,并在同一 Super-step 中执行。
builder.add_edge("prepare", "build_itinerary")
builder.add_edge("prepare", "build_packing_list")

# 列表形式的来源表示 Fan-in:两个来源全部完成后,summarize 才会执行。
builder.add_edge(["build_itinerary", "build_packing_list"], "summarize")
builder.add_edge("summarize", END)

# compile() 将构建器转换成可执行图,并检查节点和 Edge 是否有效。
graph = builder.compile()


if __name__ == "__main__":
    # invoke() 会等待并行分支和汇总节点全部结束,再返回最终 State。
    result = graph.invoke(
        {
            "city": " 杭州 ",
            "suggestions": [],
            "steps": [],
        }
    )
    print("suggestions:", result["suggestions"])
    print("steps:", result["steps"])
    print("summary:", result["summary"])

执行:

python langgraph/p10_workflow_basics/04_edge_patterns.py

结果:

suggestions: ['游览杭州的代表景点', '准备舒适的步行鞋']
steps: ['prepare', 'build_itinerary', 'build_packing_list', 'summarize']
summary: 准备舒适的步行鞋;游览杭州的代表景点

2. 路由函数与条件边

普通 Edge 的执行目标固定不变。如果下一步取决于当前 State,就需要路由函数和条件边。

2.1 直接返回节点名称

路由函数可以直接返回下一个节点的名称:

from typing import Literal

def route_by_node_name(
    state: WeatherState,
) -> Literal["hot_advice", "normal_advice"]:
    """返回注册后的节点名,Literal 列出所有合法目标。"""

    if state["temperature"] >= 30:
        return "hot_advice"
    return "normal_advice"

然后注册条件边:

# 条件路由也可以从 START 开始选择第一个业务 Node。
builder.add_conditional_edges(START, route_by_node_name)

Literal 不只是类型提示,它也让 LangGraph 和 Studio 知道路由函数可能到达哪些节点。

2.2 使用 path_map 映射业务符号

为了避免路由函数和节点名强耦合,可以让函数返回业务符号:

def route_by_symbol(state: TaskState) -> Literal["run", "stop"]:
    """返回业务符号,使路由函数不依赖具体节点名。"""

    return "run" if state["approved"] else "stop"

通过 path_map 映射到真正的目标:

# run 映射到普通 Node,stop 映射到 END。
builder.add_conditional_edges(
    START,
    route_by_symbol,
    {"run": "run_task", "stop": END},
)

这里的 stop 被映射为 END,因此可以提前结束图。

2.3 返回多个目标

路由函数也可以一次返回多个节点:

def choose_tasks(state: ParallelRouteState) -> list[str]:
    """返回多个目标,它们会在同一个 Super-step 中执行。"""

    return state["tasks"]

如果返回 ["task_a", "task_b"],两个节点会在下一个 Super-step 中执行。注册时应该明确候选目标:

# 第三个参数列出合法目标,也让 Studio 的图结构更准确。
builder.add_conditional_edges(
    START,
    choose_tasks,
    ["task_a", "task_b"],
)

2.4 路由函数不更新 State

路由函数应该保持简单:

  • 只读取 State。
  • 只判断下一步目标。
  • 不修改 State。
  • 不执行耗时 I/O。

如果一个步骤既要更新 State,又要动态跳转,可以使用 Command;如果需要为动态数据创建多份并行任务,可以使用 Send。这两种能力属于更高级的图控制方式,本文不展开。

2.5 不要混合固定后继和条件后继

不要从同一个来源 Node 同时添加普通 Edge 和条件 Edge:

# 不推荐:两个后继规则都会生效
builder.add_edge("check", "always_run")
builder.add_conditional_edges("check", route_result)

条件 Edge 不会覆盖普通 Edge。上面的 always_run 和路由选中的节点都可能执行。

完整示例位于 05_routing_functions.py

"""演示条件 Edge 的节点名路由、path_map、END 和多目标路由。

路由函数只读取 State 并返回执行目标,不应该在其中修改 State。
add_conditional_edges() 根据返回值,决定工作流下一步走向哪个 Node。
"""

import operator
from typing import Annotated, Literal

from langgraph.graph import END, START, StateGraph
from typing_extensions import NotRequired, TypedDict


class WeatherState(TypedDict):
    temperature: int
    advice: NotRequired[str]


def route_by_node_name(
    state: WeatherState,
) -> Literal["hot_advice", "normal_advice"]:
    """直接返回下一个节点名,Literal 明确列出所有可能目标。"""

    return "hot_advice" if state["temperature"] >= 30 else "normal_advice"


def hot_advice(state: WeatherState) -> dict:
    return {"advice": "注意防晒并及时补水。"}


def normal_advice(state: WeatherState) -> dict:
    return {"advice": "天气适宜,可以安排户外活动。"}


direct_builder = StateGraph(WeatherState)
direct_builder.add_node(hot_advice)
direct_builder.add_node(normal_advice)
# 条件 Edge 也可以从 START 开始,根据输入选择图的第一个 Node。
direct_builder.add_conditional_edges(START, route_by_node_name)
direct_builder.add_edge("hot_advice", END)
direct_builder.add_edge("normal_advice", END)
direct_graph = direct_builder.compile()


class TaskState(TypedDict):
    approved: bool
    result: NotRequired[str]


def route_by_symbol(state: TaskState) -> Literal["run", "stop"]:
    """返回与节点名解耦的业务符号,再由 path_map 映射目标。"""

    return "run" if state["approved"] else "stop"


def run_task(state: TaskState) -> dict:
    return {"result": "任务已执行"}


mapped_builder = StateGraph(TaskState)
mapped_builder.add_node(run_task)
# "run" 映射到普通 Node;"stop" 映射到 END,可直接结束工作流。
mapped_builder.add_conditional_edges(
    START,
    route_by_symbol,
    {"run": "run_task", "stop": END},
)
mapped_builder.add_edge("run_task", END)
mapped_graph = mapped_builder.compile()


class ParallelRouteState(TypedDict):
    tasks: list[str]
    # 多个目标 Node 可能并行写 results,operator.add 用于合并各分支结果。
    results: Annotated[list[str], operator.add]


def choose_tasks(state: ParallelRouteState) -> list[str]:
    """一次返回多个目标节点名,动态启动多个并行分支。"""

    return state["tasks"]


def task_a(state: ParallelRouteState) -> dict:
    return {"results": ["任务 A 完成"]}


def task_b(state: ParallelRouteState) -> dict:
    return {"results": ["任务 B 完成"]}


parallel_builder = StateGraph(ParallelRouteState)
parallel_builder.add_node(task_a)
parallel_builder.add_node(task_b)
# 第三个参数列出合法目标,也让图的可视化结果更准确。
parallel_builder.add_conditional_edges(
    START,
    choose_tasks,
    ["task_a", "task_b"],
)
parallel_builder.add_edge("task_a", END)
parallel_builder.add_edge("task_b", END)
parallel_graph = parallel_builder.compile()


if __name__ == "__main__":
    # invoke() 使用不同输入,即可观察三种路由方式的最终 State。
    print("节点名路由:", direct_graph.invoke({"temperature": 32}))
    print("path_map 执行:", mapped_graph.invoke({"approved": True}))
    print("path_map 提前结束:", mapped_graph.invoke({"approved": False}))
    print(
        "多目标路由:",
        parallel_graph.invoke(
            {
                "tasks": ["task_a", "task_b"],
                "results": [],
            }
        ),
    )

执行:

python langgraph/p10_workflow_basics/05_routing_functions.py

结果:

节点名路由: {'temperature': 32, 'advice': '注意防晒并及时补水。'}
path_map 执行: {'approved': True, 'result': '任务已执行'}
path_map 提前结束: {'approved': False}
多目标路由: {'tasks': ['task_a', 'task_b'], 'results': ['任务 A 完成', '任务 B 完成']}

3. 图的构建、编译与执行

StateGraph 是图的 Builder。它负责登记 State Schema、Node 和 Edge,但不能直接执行:

# StateGraph 是 Builder:这里只登记 Schema、Node 和 Edge。
builder = StateGraph(TravelState)
builder.add_node(normalize_input)
builder.add_edge(START, "normalize_input")
builder.add_edge("normalize_input", END)

完成结构定义后调用:

# compile() 检查图结构,并返回可以执行的 CompiledStateGraph。
graph = builder.compile()

compile() 会检查图结构,并返回可以执行的 CompiledStateGraph

3.1 invoke 与 ainvoke

同步图可以使用:

# invoke() 等待同步图运行完成,并返回最终结果。
result = graph.invoke(input_state)

包含异步 Node 的图使用:

# ainvoke() 是异步调用接口,应在 async 函数中使用。
result = await graph.ainvoke(input_state)

二者都会等待整张图运行结束,并返回经过 output_schema 过滤后的最终结果。

3.2 stream 与 astream

同步流式执行:

# updates 模式逐步返回各 Node 本轮提交的局部更新。
for update in graph.stream(input_state, stream_mode="updates"):
    print(update)

异步流式执行:

# astream() 是异步流式接口,需要使用 async for 消费。
async for update in graph.astream(input_state, stream_mode="updates"):
    print(update)

常用的流模式可以整理为:

模式 主要产出方 每次返回的内容 适用场景
values LangGraph 每个 Super-step 合并完成后的完整 State 快照 观察 State 的整体变化
updates LangGraph Node 提交的局部更新,格式为 {node_name: update} 观察哪个 Node 修改了哪些字段
checkpoints LangGraph Checkpointer 保存状态时产生的 checkpoint 事件 观察持久化检查点的写入过程
tasks / debug LangGraph 任务开始、结束、结果、错误和其他调试信息 分析任务调度和排查执行问题
messages LangGraph 与 LLM 集成 Node 内部 LLM 调用产生的消息或 token 片段及元数据 在页面中实时展示模型输出
custom Node 主动写入 Node 通过 StreamWriter 写入的自定义事件 输出业务进度或自定义中间结果

values 用于调试内部状态,因此不会因为 output_schema 而只保留公开输出字段。

4. 可部署到 LangGraph Studio 的完整案例

4.1 案例目标与执行流程

下面把前面的知识组合成一个城市旅行方案 Workflow。调用者提交城市、温度、天数、人数、每日人均预算和旅行偏好,Workflow 会完成:

  1. 输入规范化和校验。
  2. 根据温度选择天气策略。
  3. 并行生成行程、行李和预算区块。
  4. 等待三个分支完成并整理最终报告。

校验失败时会提前结束;校验成功时只有一个天气策略节点会执行。进入并行阶段后,三个规划节点处于同一个 Super-step,assemble_plan 等待它们全部完成。

这个完整案例把 assemble_plan 注册为 defer=True,并让三个并行节点分别连接它。defer=True 表示等其他待执行任务全部完成后再运行汇总节点。与节点列表形式的 Fan-in 相比,这种写法还会让 Agent Server 在 checkpoint 历史中把 assemble_plan 保存为独立任务,因此可以在 Studio 时间线中正常显示。节点列表形式的 Fan-in 仍然是有效的图结构,前面的 Edge 示例继续用它说明等待屏障;这里的调整是为了兼顾 Studio 的可观察性。

4.2 代码示例

06_studio_trip_workflow.py 的完整代码如下:

"""可部署到 LangGraph Studio 的城市旅行方案 Workflow。

这个案例把前面介绍的 State、Reducer、Node、Edge 和条件路由组合起来:
先校验输入并根据天气选择建议,再并行生成行程、行李和预算,最后汇总输出。

文件必须在模块顶层导出名为 graph 的已编译图,LangGraph Studio 才能通过
langgraph.json 加载它。直接运行本文件时,末尾的 main 演示也会执行。
"""

import operator
from typing import Annotated, Literal

from langgraph.graph import END, START, StateGraph
from langgraph.types import Overwrite
from typing_extensions import NotRequired, TypedDict


# =============================================================================
# 1. 定义 State Schema
# =============================================================================

class PlanSection(TypedDict):
    """一个旅行方案区块。

    order 是明确的业务顺序。并行 Node 的完成顺序不稳定,因此不能直接用
    draft_sections 的到达顺序决定最终展示顺序。
    """

    order: int
    title: str
    content: str


class TripInput(TypedDict):
    """调用者必须提供的公开输入字段。"""

    city: str
    temperature: int
    days: int
    people: int
    daily_budget: int
    preferences: list[str]


class TripState(TypedDict):
    """图内部使用的完整 State。

    每个字段都可以看作一个 Channel。Node 返回局部更新后,LangGraph 会按照
    该字段配置的 Reducer,把新值写入当前 State。
    """

    # 输入字段
    city: str
    temperature: int
    days: int
    people: int
    daily_budget: int
    preferences: list[str]

    # 节点逐步补充的内部字段
    errors: NotRequired[list[str]]
    weather_category: NotRequired[str]
    weather_advice: NotRequired[str]
    total_budget: NotRequired[int]
    # 三个并行 Node 都会提交一个方案区块,operator.add 负责合并列表。
    draft_sections: Annotated[list[PlanSection], operator.add]
    # steps 同样可能被并行写入,用于观察实际执行过哪些 Node。
    steps: Annotated[list[str], operator.add]
    status: NotRequired[str]
    plan_sections: NotRequired[list[PlanSection]]
    summary: NotRequired[str]


class TripOutput(TypedDict):
    """图执行完成后向调用者公开的字段。"""

    status: str
    city: str
    weather_category: str
    total_budget: int
    plan_sections: list[PlanSection]
    summary: str
    steps: list[str]


# =============================================================================
# 2. 定义输入处理与校验 Node
# =============================================================================

def normalize_input(state: TripState) -> dict:
    """初始化本次运行,再清理城市和偏好。

    Node 不直接修改传入的 state,只返回需要写回的局部字段。
    Studio 会把同一个 Thread 的 State 保留下来,因此这里用 Overwrite 清空
    上一次运行留下的 Reducer 字段,避免重复提交时继续追加旧结果。
    """

    # 用集合去重,再排序以获得稳定、便于测试的结果。
    preferences = {
        preference.strip()
        for preference in state["preferences"]
        if preference.strip()
    }
    return {
        "city": state["city"].strip(),
        "preferences": sorted(preferences),
        # Overwrite 会绕过 operator.add,直接开始一份新的草稿与执行记录。
        "draft_sections": Overwrite([]),
        "steps": Overwrite(["normalize_input"]),
    }


def validate_input(state: TripState) -> dict:
    """检查工作流输入是否在示例支持的范围内。"""

    # 一次收集所有校验错误,调用者可以同时看到全部问题。
    errors: list[str] = []

    if not state["city"]:
        errors.append("城市不能为空")
    if not -50 <= state["temperature"] <= 60:
        errors.append("温度必须在 -50℃ 到 60℃ 之间")
    if not 1 <= state["days"] <= 7:
        errors.append("旅行天数必须在 1 到 7 天之间")
    if state["people"] <= 0:
        errors.append("出行人数必须大于 0")
    if state["daily_budget"] <= 0:
        errors.append("每日人均预算必须大于 0")

    return {
        "errors": errors,
        "steps": ["validate_input"],
    }


# =============================================================================
# 3. 定义校验结果路由与失败分支
# =============================================================================

def route_after_validation(
    state: TripState,
) -> Literal["reject_input", "classify_weather"]:
    """根据校验结果返回下一个节点名。

    路由函数只做决策,不更新 State。Literal 明确列出可能目标,也能让
    LangGraph Studio 更准确地绘制条件 Edge。
    """

    return "reject_input" if state["errors"] else "classify_weather"


def reject_input(state: TripState) -> dict:
    """为非法输入生成结构稳定的输出,然后通过 Edge 直接到 END。"""

    return {
        "status": "rejected",
        "weather_category": "invalid",
        "total_budget": 0,
        "plan_sections": [],
        "summary": "输入校验失败:" + ";".join(state["errors"]),
        "steps": ["reject_input"],
    }


# =============================================================================
# 4. 定义天气分类 Node 与条件路由
# =============================================================================

def classify_weather(state: TripState) -> dict:
    """按照温度划分高温、低温和舒适三类天气。"""

    if state["temperature"] >= 30:
        category = "hot"
    elif state["temperature"] <= 10:
        category = "cold"
    else:
        category = "comfortable"

    return {
        "weather_category": category,
        "steps": ["classify_weather"],
    }


def route_weather(
    state: TripState,
) -> Literal["hot_plan", "cold_plan", "comfortable_plan"]:
    """把天气分类映射到对应的建议 Node。"""

    routes = {
        "hot": "hot_plan",
        "cold": "cold_plan",
        "comfortable": "comfortable_plan",
    }
    return routes[state["weather_category"]]


def hot_plan(state: TripState) -> dict:
    """生成高温天气建议。"""

    return {
        "weather_advice": "避开正午暴晒,准备防晒用品并及时补水。",
        "steps": ["hot_plan"],
    }


def cold_plan(state: TripState) -> dict:
    """生成低温天气建议。"""

    return {
        "weather_advice": "准备保暖衣物,优先安排室内外结合的活动。",
        "steps": ["cold_plan"],
    }


def comfortable_plan(state: TripState) -> dict:
    """生成舒适天气建议。"""

    return {
        "weather_advice": "天气舒适,可以适当增加户外步行活动。",
        "steps": ["comfortable_plan"],
    }


# =============================================================================
# 5. 定义并行规划 Node
# =============================================================================

def prepare_parallel(state: TripState) -> dict:
    """三个天气分支在此汇合,并记录并行规划阶段开始。"""

    return {"steps": ["prepare_parallel"]}


def build_itinerary(state: TripState) -> dict:
    """并行分支一:根据城市、天数、偏好和天气生成行程区块。"""

    preferences = "、".join(state["preferences"]) or "当地代表景点"
    content = (
        f'在{state["city"]}安排 {state["days"]} 天行程,'
        f"重点体验:{preferences}{state['weather_advice']}"
    )
    section: PlanSection = {
        "order": 1,
        "title": "行程建议",
        "content": content,
    }
    return {
        "draft_sections": [section],
        "steps": ["build_itinerary"],
    }


def build_packing_list(state: TripState) -> dict:
    """并行分支二:根据天气类型生成行李清单区块。"""

    if state["weather_category"] == "hot":
        weather_items = "防晒霜、遮阳帽和水杯"
    elif state["weather_category"] == "cold":
        weather_items = "保暖外套、围巾和暖宝宝"
    else:
        weather_items = "轻便外套和舒适的步行鞋"

    section: PlanSection = {
        "order": 2,
        "title": "行李清单",
        "content": f"建议准备:证件、常用药品、{weather_items}。",
    }
    return {
        "draft_sections": [section],
        "steps": ["build_packing_list"],
    }


def calculate_budget(state: TripState) -> dict:
    """并行分支三:计算总预算并生成预算区块。"""

    # 这里的预算规则保持简单,便于把注意力放在 Workflow 结构上。
    total_budget = state["days"] * state["people"] * state["daily_budget"]
    section: PlanSection = {
        "order": 3,
        "title": "预算估算",
        "content": (
            f'{state["people"]} 人出行 {state["days"]} 天,'
            f'按照每日人均 {state["daily_budget"]} 元计算,'
            f"总预算约为 {total_budget} 元。"
        ),
    }
    return {
        "total_budget": total_budget,
        "draft_sections": [section],
        "steps": ["calculate_budget"],
    }


def assemble_plan(state: TripState) -> dict:
    """Fan-in 汇总 Node:三个并行分支全部完成后才会执行。

    draft_sections 是 Reducer 合并的中间数据,plan_sections 才是整理后的
    对外结果。把二者分开,可以避免在 Reducer 中混入展示层业务逻辑。
    """

    # 不依赖并行分支的完成顺序,而是按业务字段 order 显式排序。
    plan_sections = sorted(
        state["draft_sections"],
        key=lambda section: section["order"],
    )
    return {
        "status": "success",
        "plan_sections": plan_sections,
        "summary": (
            f'{state["city"]} {state["days"]} 天旅行方案已生成,'
            f'天气类型为 {state["weather_category"]},'
            f'预计总预算 {state["total_budget"]} 元。'
        ),
        "steps": ["assemble_plan"],
    }


# =============================================================================
# 6. 创建图、注册 Node 和 Edge
# =============================================================================

# TripState 是内部完整状态;input_schema 和 output_schema 分别限制公开输入、输出。
builder = StateGraph(
    TripState,
    input_schema=TripInput,
    output_schema=TripOutput,
)

# 先注册所有 Node。函数名会作为默认节点名。
builder.add_node(normalize_input)
builder.add_node(validate_input)
builder.add_node(reject_input)
builder.add_node(classify_weather)
builder.add_node(hot_plan)
builder.add_node(cold_plan)
builder.add_node(comfortable_plan)
builder.add_node(prepare_parallel)
builder.add_node(build_itinerary)
builder.add_node(build_packing_list)
builder.add_node(calculate_budget)
# defer=True 让汇总节点等其他待执行任务全部结束后再运行。
# 这种注册方式也会让 Agent Server 把 assemble_plan 记录为独立任务,
# 因而可以在 LangGraph Studio 的时间线中正常显示。
builder.add_node(assemble_plan, defer=True)

# 固定 Edge:工作流从输入清理开始,再进行输入校验。
builder.add_edge(START, "normalize_input")
builder.add_edge("normalize_input", "validate_input")

# 条件 Edge:校验失败进入 reject_input,成功则进入天气分类。
builder.add_conditional_edges("validate_input", route_after_validation)
builder.add_edge("reject_input", END)

# 第二个条件 Edge:根据天气分类选择且只执行一个天气建议 Node。
builder.add_conditional_edges("classify_weather", route_weather)

# 三个互斥天气分支最终都会到达同一个并行准备节点。
builder.add_edge("hot_plan", "prepare_parallel")
builder.add_edge("cold_plan", "prepare_parallel")
builder.add_edge("comfortable_plan", "prepare_parallel")

# 静态 Fan-out:一个来源连接三个后继,它们会在同一 Super-step 中执行。
builder.add_edge("prepare_parallel", "build_itinerary")
builder.add_edge("prepare_parallel", "build_packing_list")
builder.add_edge("prepare_parallel", "calculate_budget")

# 三个并行分支分别触发延迟汇总节点。
# 因为 assemble_plan 配置了 defer=True,它会在所有待执行任务完成后只运行一次。
builder.add_edge("build_itinerary", "assemble_plan")
builder.add_edge("build_packing_list", "assemble_plan")
builder.add_edge("calculate_budget", "assemble_plan")
builder.add_edge("assemble_plan", END)

# =============================================================================
# 7. 编译并导出图
# =============================================================================

# compile() 把 StateGraph 构建器转换为可执行的 CompiledStateGraph。
# LangGraph Studio 会按照 langgraph.json 的配置导入这个顶层 graph 变量。
graph = builder.compile()

# =============================================================================
# 8. 本地 main 演示
# =============================================================================

def print_demo(city: str, temperature: int) -> None:
    """通过 invoke() 运行一次完整 Workflow,并打印最终输出。"""

    # invoke() 接收 TripInput,等待所有 Node 完成后返回 TripOutput。
    result = graph.invoke(
        {
            "city": city,
            "temperature": temperature,
            "days": 3,
            "people": 2,
            "daily_budget": 500,
            "preferences": ["美食", "自然风景"],
        }
    )
    print(result["summary"])
    print("steps:", result["steps"])
    # plan_sections 已经在 assemble_plan 中按 order 排好序。
    for section in result["plan_sections"]:
        print(f'- {section["title"]}: {section["content"]}')


if __name__ == "__main__":
    # 直接运行 Python 文件时,分别覆盖高温、低温和舒适三条路由。
    # 被 LangGraph Studio 导入时,这一代码块不会执行。
    print_demo("杭州", 32)
    print_demo("哈尔滨", 5)
    print_demo("昆明", 22)

这里把命令行演示放在 if __name__ == "__main__" 中。langgraph dev 导入模块时只会创建图,不会自动执行 print_demo()

4.3 配置 LangGraph 项目

langgraph.json 注册 Studio 需要加载的图:

{
  "$schema": "https://langgra.ph/schema.json",
  "dependencies": ["."],
  "graphs": {
    "trip_workflow": "./06_studio_trip_workflow.py:graph"
  },
  "python_version": "3.12"
}

路径冒号后面的 graph 必须是模块顶层真实存在的、已经编译的图。

pyproject.toml 定义项目和依赖:

# setuptools 负责把当前目录作为 Python 项目安装。
[build-system]
requires = ["setuptools==82.0.1"]
build-backend = "setuptools.build_meta"

# 项目元数据和运行依赖。
[project]
name = "workflow-basics"
version = "0.1.0"
description = "使用确定性 Python 节点学习 LangGraph Workflow"
requires-python = ">=3.12,<3.13"
dependencies = [
    # 基础示例与 StateGraph 需要的依赖。
    "langchain-core==1.2.22",
    "langgraph==1.1.3",
    # langgraph dev 使用的本地内存版 Agent Server。
    "langgraph-cli[inmem]==0.4.19",
    "pydantic==2.12.5",
    "typing-extensions==4.15.0",
]

# 本项目通过 langgraph.json 按文件路径加载图,不发布独立 Python 模块。
[tool.setuptools]
py-modules = []

这个案例不调用模型和外部接口,因此项目代码不需要 .env 或业务 API Key。

4.4 启动 Agent Server

必须进入包含 langgraph.json 的目录启动:

# langgraph dev 默认读取当前目录中的 langgraph.json。
cd /Users/bianhn/Documents/git/llm-learning/langgraph/p10_workflow_basics
source ../../.venv_langgraph/bin/activate
# 以可编辑模式安装项目和 langgraph-cli[inmem]。
uv pip install -e .
# 启动本地 Agent Server,并自动打开 Studio。
langgraph dev --host 127.0.0.1 --port 2024

langgraph dev 会启动用于本地开发和测试的 Agent Server,并在终端输出:

  • 本地 API 地址。
  • API 文档地址。
  • LangGraph Studio Web UI 地址。

这是本地开发部署,不是生产部署。普通 Python 代码修改通常会自动热重载;修改 langgraph.json、依赖,或者热重载失败时,应重新启动服务。

workflow 图

4.5 在 Studio 中提交输入

这张图的 State 不是 MessagesState,因此应该使用 Studio 的 Graph 模式,而不是 Chat 模式。

选择 trip_workflow,在 Input 区域切换到 Raw JSON,输入:

{
  "city": " 杭州 ",
  "temperature": 32,
  "days": 3,
  "people": 2,
  "daily_budget": 500,
  "preferences": ["美食", "自然风景", "美食"]
}

预期结果包括:

  • city 被清理为“杭州”。
  • weather_categoryhot
  • total_budget3000
  • plan_sections 按行程、行李、预算的顺序排列。
  • statussuccess

Studio 的时间线中,build_itinerarybuild_packing_listcalculate_budget 位于同一个并行阶段,assemble_plan 在三个节点全部完成后执行。

不要根据时间线中三个并行节点的显示先后推断业务顺序。最终顺序来自每个区块的 order 字段。

Graph 模式不会像聊天机器人那样额外生成一个回答气泡。这里的“最终结果”就是图结束时的结构化 State:在右侧时间线向下滚动并展开 assemble_plan,可以看到它提交的 statusplan_sectionssummary;也可以在最终 State 中查看经过 TripOutput 过滤后的完整输出。例如:

{
  "status": "success",
  "city": "杭州",
  "weather_category": "hot",
  "total_budget": 3000,
  "plan_sections": [
    {"order": 1, "title": "行程建议", "content": "..."},
    {"order": 2, "title": "行李清单", "content": "..."},
    {"order": 3, "title": "预算估算", "content": "..."}
  ],
  "summary": "杭州 3 天旅行方案已生成,天气类型为 hot,预计总预算 3000 元。",
  "steps": [
    "normalize_input",
    "validate_input",
    "classify_weather",
    "hot_plan",
    "prepare_parallel",
    "build_itinerary",
    "build_packing_list",
    "calculate_budget",
    "assemble_plan"
  ]
}

Studio 的同一个 Thread 会保留上一次运行的 State。因为 draft_sectionssteps 使用追加型 Reducer,如果不做处理,重复点击 Submit 会把旧列表再次带入新一轮执行。示例在 normalize_input 中使用 Overwrite 开始新的草稿和执行记录,因此可以在同一个 Thread 中反复测试,而不会得到重复区块。若要观察每次运行完全独立的 State,也可以在 Studio 中新建 Thread。

执行结果

4.6 直接执行

执行:

python 06_studio_trip_workflow.py

结果:

杭州 3 天旅行方案已生成,天气类型为 hot,预计总预算 3000 元。
steps: ['normalize_input', 'validate_input', 'classify_weather', 'hot_plan', 'prepare_parallel', 'build_itinerary', 'build_packing_list', 'calculate_budget', 'assemble_plan']
- 行程建议: 在杭州安排 3 天行程,重点体验:美食、自然风景。避开正午暴晒,准备防晒用品并及时补水。
- 行李清单: 建议准备:证件、常用药品、防晒霜、遮阳帽和水杯。
- 预算估算: 2 人出行 3 天,按照每日人均 500 元计算,总预算约为 3000 元。
哈尔滨 3 天旅行方案已生成,天气类型为 cold,预计总预算 3000 元。
steps: ['normalize_input', 'validate_input', 'classify_weather', 'cold_plan', 'prepare_parallel', 'build_itinerary', 'build_packing_list', 'calculate_budget', 'assemble_plan']
- 行程建议: 在哈尔滨安排 3 天行程,重点体验:美食、自然风景。准备保暖衣物,优先安排室内外结合的活动。
- 行李清单: 建议准备:证件、常用药品、保暖外套、围巾和暖宝宝。
- 预算估算: 2 人出行 3 天,按照每日人均 500 元计算,总预算约为 3000 元。
昆明 3 天旅行方案已生成,天气类型为 comfortable,预计总预算 3000 元。
steps: ['normalize_input', 'validate_input', 'classify_weather', 'comfortable_plan', 'prepare_parallel', 'build_itinerary', 'build_packing_list', 'calculate_budget', 'assemble_plan']
- 行程建议: 在昆明安排 3 天行程,重点体验:美食、自然风景。天气舒适,可以适当增加户外步行活动。
- 行李清单: 建议准备:证件、常用药品、轻便外套和舒适的步行鞋。
- 预算估算: 2 人出行 3 天,按照每日人均 500 元计算,总预算约为 3000 元。
(.venv_langgraph) (base) bianhn@BianhnMacBook-Pro p10_workflow_basics %

5. 总结

本篇围绕 Workflow 的执行路径和完整落地过程介绍了:

  1. Edge 使用 START、END 和固定连接控制节点执行顺序。
  2. 静态 Fan-out 可以让多个节点进入同一个 Super-step 并行执行。
  3. Fan-in、Reducer 和 defer=True 可以协调并行分支与最终汇总。
  4. 路由函数根据当前 State 选择条件路径,但不负责更新 State。
  5. compile() 把 Builder 转换成可以调用和流式执行的图。
  6. langgraph dev 可以把顶层编译图加载到 Agent Server,并在 Studio 中观察节点、State 和最终结果。

城市旅行方案案例把输入校验、两级条件路由、并行执行、Reducer 合并和延迟汇总组合到一张可运行的图中。系列 11 将在这些基础上加入评估器与优化器循环,让 Workflow 根据评估结果决定继续改进还是结束。


文章作者: hnbian
版权声明: 本博客所有文章除特別声明外,均采用 CC BY 4.0 许可协议。转载请注明来源 hnbian !
评论
  目录