Day 47 / 共 60 天 · 阶段 8 函数式 API 与子图

@entrypoint / @task:不画图,也能用 LangGraph

前 46 天我们都在用 StateGraph:手动 add_node、add_edge、画出一张图。今天登场的是另一套写法——函数式 API:你只写普通 Python 函数,@task 装饰一段可缓存/可重试的活儿,@entrypoint 把整个流程函数直接变成一张能跑、能存档、能中断恢复的图。它不是新引擎,而是"图引擎的函数糖衣"。看懂它,你会更懂昨天学的持久化到底给了函数式什么魔法。

📍 阶段 8 · 函数式 API 与子图(6 天)你在这里
D47 @entrypoint/@task D48 func→pregel D49 子图基础 D50 子图隔离/stream D51 CachePolicy D52 重试容错
🤔 痛点:画图太"重"了 想写个"生成文章 → 人工审阅 → 返回结果"的流程,用 StateGraph 你得:定义 State TypedDict、写节点函数、add_node、add_edge、set entry……三行业务逻辑要配十行脚手架。有没有办法:就写一个普通函数,中间该并行的并行、该断点的断点、该存档的存档,全自动?
💡 本质:函数式 API = 把"控制流"藏进函数调用里 StateGraph 是声明式(先把图的形状全画出来,再跑);函数式 API 是命令式(顺着代码往下读,调用到哪算哪)。类比:声明式像先画好地铁线路图;命令式像你直接边走边导航。二者跑的是同一列地铁(Pregel 引擎),只是"描述路线"的方式不同。@task 标记"一个站点",@entrypoint 把"你今天的整条行程"打包成一张可复用的线路。
L01

两种写法,一份引擎

函数式 API 的全部公开符号只有两个,看模块开头 func/__init__.py:56

# func/__init__.py:56
__all__ = ("task", "entrypoint")

对比一下同一个"加 1 流程"的两种写法:

# 声明式(前 46 天的写法)
from langgraph.graph import StateGraph, START, END
from typing import TypedDict
class S(TypedDict):
    n: int
g = StateGraph(S)
g.add_node("inc", lambda s: {"n": s["n"] + 1})
g.add_edge(START, "inc")
g.add_edge("inc", END)
app = g.compile()
app.invoke({"n": 1})            # {'n': 2}

# 命令式(今天的函数式 API)
from langgraph.func import entrypoint, task
@task
def inc(n: int) -> int:
    return n + 1
@entrypoint()
def app(n: int) -> int:
    return inc(n).result()       # 直接像调函数
app.invoke(1)                    # 2
@taskinc 标记成"一个可被引擎调度的任务"——它能被缓存、重试、并行。
@entrypoint()app 这个普通函数变成一张 Pregel 图,于是它也有了 .invoke / .stream
inc(n).result()调用 task 不直接返回结果,而是返回一个 Future(L03 细讲),.result() 才拿到值——这是"能并行"的关键。
💡 一句话函数式 API 没有另起炉灶。@entrypoint 最后 return 的就是一个 Pregel 对象(L04 会看到源码里那行 graph: Pregel = Pregel(...))。所以你前 46 天学的 checkpoint、interrupt、stream,函数式全能用。
L02

@task 到底做了什么:包成 _TaskFunction

task 是个可带参也可裸用的装饰器 func/__init__.py:132。核心是把你的函数塞进一个 _TaskFunction 对象 func/__init__.py:59

# func/__init__.py:59
class _TaskFunction(Generic[P, T]):
    def __init__(self, func, *, retry_policy, cache_policy=None,
                 timeout=None, name=None):
        ...
        self.func = func                    # 你原始的函数
        self.retry_policy = retry_policy     # 失败重试策略(D52)
        self.cache_policy = cache_policy     # 结果缓存策略(D51)
        self.timeout = timeout
        functools.update_wrapper(self, func) # 把 __name__/__doc__ 抄过来,伪装成原函数

task(...) 本体只是"整理参数 + 造装饰器",看 func/__init__.py:226

# func/__init__.py:226
retry_policies: Sequence[RetryPolicy] = (
    ()                                       # 没给 → 空元组(不重试)
    if retry_policy is None
    else (retry_policy,)                     # 给了单个 → 包成 1 元组
    if isinstance(retry_policy, RetryPolicy)
    else retry_policy                        # 给了列表 → 原样
)
def decorator(func):
    ...
    return _TaskFunction(func, retry_policy=retry_policies,
                         cache_policy=cache_policy, timeout=timeout_policy, name=name)
if __func_or_none__ is not None:             # 裸用 @task(没括号)
    return decorator(__func_or_none__)
return decorator                             # 带参 @task(...)(有括号)
retry_policy 统一成元组无论你给 None / 单个 / 列表,内部一律存成"策略序列"。D52 会看到执行时按顺序找第一个匹配的策略。
__func_or_none__ 分支经典装饰器技巧:裸写 @task 时第一个位置参数就是被装饰函数,直接 decorate;写 @task(retry_policy=...) 时先返回 decorator 等 Python 再喂函数。
💡 设计取舍①:为什么用"类实例"包装,而不是普通闭包函数?朴素做法是 functools.wraps 返回一个闭包函数。但 LangGraph 用了 _TaskFunction 这个,因为它需要挂载额外能力:clear_cache() / aclear_cache() 方法(func/__init__.py:96)让你能主动清掉这个 task 的缓存,还要保存 retry_policy / cache_policy 供引擎读取。闭包塞不下这些"对象方法 + 结构化属性"。代价是它不是真函数(是 callable 对象),所以源码里靠 functools.update_wrapper 把名字/文档抄过来,让它"看起来像"原函数、IDE 补全和堆栈信息不出戏。
L03

调用 @task 为什么返回 Future?

关键在 _TaskFunction.__call__ func/__init__.py:86——你以为在"调函数",其实在"提交一个任务":

# func/__init__.py:86
def __call__(self, *args, **kwargs) -> SyncAsyncFuture[T]:
    return _call_with_options(
        self.func, args, kwargs,
        retry_policy=self.retry_policy,
        cache_policy=self.cache_policy,
        timeout=self.timeout,
    )

_call_with_options 不当场执行你的函数,而是把它交给"当前正在跑的引擎",返回一个 Future pregel/_call.py:276

# pregel/_call.py:276
def _call_with_options(func, args, kwargs, *, retry_policy=None,
                       cache_policy=None, timeout=None) -> SyncAsyncFuture[T]:
    ...
    config = get_config()                    # 拿到"我现在在哪张图里跑"的上下文
    impl = config[CONF][CONFIG_KEY_CALL]     # 引擎放在 config 里的"提交任务"回调
    fut = impl(func, (args, kwargs), retry_policy=retry_policy,
               cache_policy=cache_policy, callbacks=config["callbacks"],
               timeout=timeout)
    return fut                               # 返回 Future,不是结果本身
get_config()函数式的魔法源头:task 只能在 entrypoint(或 StateGraph)内部被调用,此时线程/上下文里存着"当前图"的 config。脱离图直接调 task 会报错。
config[CONF][CONFIG_KEY_CALL]引擎在启动每个任务前,往 config 塞了一个"提交任务"的回调 impl。task 调用就是"喊引擎:帮我排一个活儿"。D48 会看到这个回调是谁塞的。
返回 fut因为是 Future,你可以先发起 N 个 task(N 个活儿并行排队),再统一 .result() 收割——这就是函数式并行的写法。

SyncAsyncFuture 本身是个"既能 .result() 又能 await"的双面 Future pregel/_call.py:253

# pregel/_call.py:253
class SyncAsyncFuture(Generic[T], concurrent.futures.Future[T]):
    def __await__(self) -> Generator[T, None, T]:
        yield cast(T, ...)
📝 并行的真实写法(来自官方 docstring, func/__init__.py:185)
@entrypoint()
def add_one(numbers: list[int]) -> list[int]:
    futures = [add_one_task(n) for n in numbers]   # 一口气发起 N 个任务
    results = [f.result() for f in futures]         # 再统一收结果
    return results
add_one.invoke([1, 2, 3])                           # [2, 3, 4]
循环里 add_one_task(n) 瞬间返回,不阻塞;三个任务被引擎并行调度,最后 .result() 才等结果。
控制流:调 task 返回 Future,先发起后收割(并行) inc(1) inc(2) inc(3) 瞬间返回 3 个 Future 引擎并行调度 3 个任务同时跑 f.result() ×3 统一收割结果 对比:若直接返回结果 → 只能串行 inc(1) 等完再 inc(2),无法并行、无法中途 checkpoint
图注:分离"发起"与"取值",才能在两者之间做并行调度、缓存命中、中断保存。
💡 本质:Future 是"命令式代码里塞进并行/断点"的接口如果 task 直接返回结果,代码就是纯串行、也无法在中间插入 checkpoint/interrupt。返回 Future 后,"发起"和"取结果"分开了,引擎得以在两者之间做调度、缓存命中、甚至中断保存——这正是为什么函数式也能用上前 46 天的所有能力。
L04

@entrypoint 把函数变成一张 Pregel 图

entrypoint 是个类(为了挂 .final),真正的转换在 __call__ func/__init__.py:516。它读你函数的签名,然后 手工组装一张只有一个节点的 Pregel 图 func/__init__.py:576

# func/__init__.py:576
graph: Pregel = Pregel(
    nodes={
        func.__name__: PregelNode(              # 你的函数 = 唯一的节点
            bound=bound,                         # 包装后的可运行体
            triggers=[START],                    # 被 START 触发
            channels=START,                      # 从 START 通道读输入
            timeout=self.timeout,
            writers=[
                ChannelWrite([
                    ChannelWriteEntry(END, mapper=_pluck_return_value),  # 返回值 → END
                    ChannelWriteEntry(PREVIOUS, mapper=_pluck_save_value),# 存档值 → PREVIOUS
                ])
            ],
        )
    },
    channels={
        START: EphemeralValue(input_type),       # 输入通道(用完即弃)
        END: LastValue(output_type, END),        # 输出通道
        PREVIOUS: LastValue(save_type, PREVIOUS),# "上一次的结果"通道
    },
    input_channels=START, output_channels=END, stream_channels=END,
    stream_mode="updates", stream_eager=True,
    checkpointer=self.checkpointer, store=self.store, cache=self.cache,
    cache_policy=self.cache_policy, retry_policy=self.retry_policy or (),
    context_schema=self.context_schema,
)
nodes={func.__name__: PregelNode(...)}整个 entrypoint 只编译成一个节点——你函数体里的 task 调用是运行时动态派生的子任务(D48),不在这张静态图里。
triggers=[START], channels=START节点被 START 触发、从 START 读输入。等于"图一开,就跑你的函数"。
writers=ChannelWrite([...])函数返回后,把返回值写到 END(给调用方),把"要存的值"写到 PREVIOUS(给下次用)。两个 mapper 决定怎么拆(L06)。
checkpointer=self.checkpointer你在 @entrypoint(checkpointer=...) 传的存档器,原样接到 Pregel 上——于是函数式也能断点续跑、时间旅行。

bound(被包装的函数体)从哪来?看 func/__init__.py:531

# func/__init__.py:526
if inspect.isgeneratorfunction(func) or inspect.isasyncgenfunction(func):
    raise NotImplementedError("Generators are not supported in the Functional API.")
bound = get_runnable_for_entrypoint(func)     # 把普通函数裹成 Runnable
数据结构:@entrypoint 组装出的单节点 Pregel 图 START Ephemeral PregelNode = 你的函数 bound (Runnable 包装) writers: 拆返回值 → END / PREVIOUS triggers=[START] END LastValue PREVIOUS 存档给下次(previous)
图注:一个 entrypoint 只有一个节点,三条通道 START/END/PREVIOUS 各司其职。
L05

三个隐形通道:START / END / PREVIOUS

StateGraph 里你自己定义 State 字段当通道;entrypoint 里通道是引擎硬编码的三条 func/__init__.py:593

通道类型作用
STARTEphemeralValue装 invoke 传进来的输入。用完即弃(Day 31 讲的短暂通道),不进档。
ENDLastValue装函数返回值,覆盖写。invoke 的返回值从这里取。
PREVIOUSLastValue装"要存给下一次的值"。下次同 thread 调用时,通过 previous 参数注入。

输入类型不是瞎猜的——entrypoint 读你函数第一个参数的类型注解 func/__init__.py:535

# func/__init__.py:535
sig = inspect.signature(func)
first_parameter_name = next(iter(sig.parameters.keys()), None)
if not first_parameter_name:
    raise ValueError("Entrypoint function must have at least one parameter")
input_type = (
    sig.parameters[first_parameter_name].annotation
    if sig.parameters[first_parameter_name].annotation is not inspect.Signature.empty
    else Any
)
必须有至少一个参数entrypoint 函数只接收一个输入参数(想传多个值就用 dict)。没有参数直接报错——因为 START 通道不知道往哪灌数据。
读第一个参数的注解当 input_typedef app(n: int) → START 通道就是 int 类型;没写注解就退化成 Any
⚠️ 边界:entrypoint 函数只能有一个"业务输入参数"初学者常写 def app(a, b) 想传两个值,但引擎只把第一个参数接到 START。config / previous / runtime 这几个是引擎注入的特殊参数(按名字识别,见 docstring func/__init__.py:276),不算业务输入。要传多个业务值,规范做法是 def app(inputs: dict) 然后 inputs["a"]
L06

previous 与 entrypoint.final:返回值 ≠ 存档值

默认情况下,函数返回什么,就存什么、下次 previous 拿到的也是它。但有时你想"返回给调用方 A,却给下次存 B"。这就是 entrypoint.final func/__init__.py:475

# func/__init__.py:475
@dataclass(**_DC_KWARGS)
class final(Generic[R, S]):
    value: R   # 返回给调用方的值(invoke 拿到这个)
    save: S    # 存进 checkpoint 的值(下次 previous 拿到这个)

拆分逻辑就是 L04 那两个 mapper func/__init__.py:546

# func/__init__.py:546
def _pluck_return_value(value):    # 写 END 用:取 .value,否则原样
    return value.value if isinstance(value, entrypoint.final) else value
def _pluck_save_value(value):      # 写 PREVIOUS 用:取 .save,否则原样
    return value.save if isinstance(value, entrypoint.final) else value
普通 return x两个 mapper 都拿到 x:END=x(返回)、PREVIOUS=x(存档)。返回值和存档值相同。
return entrypoint.final(value=A, save=B)_pluck_return_value 取 A 写 END;_pluck_save_value 取 B 写 PREVIOUS。一次返回,两个去向。
📝 官方例子(func/__init__.py:417):一个"记忆上次输入"的计数器
@entrypoint(checkpointer=InMemorySaver())
def my_workflow(number: int, *, previous: Any = None) -> entrypoint.final[int, int]:
    previous = previous or 0
    return entrypoint.final(value=previous, save=2 * number)

config = {"configurable": {"thread_id": "some_thread"}}
my_workflow.invoke(3, config)  # 返回 0   (previous 为 None → 0;存 2*3=6)
my_workflow.invoke(1, config)  # 返回 6   (previous 为上次存的 6;存 2*1=2)
第一次返回 0 但偷偷存了 6;第二次 previous 就是 6。返回值和存档值完全解耦。
💡 设计取舍②:为什么不直接让用户操作 checkpoint,而搞个 final 对象?朴素想法是"给用户一个 save_to_checkpoint(x) 函数随便存"。但那会让存档时机变得不可控、也难序列化。entrypoint.final 把"返回 vs 存档"收敛成一个纯数据的 dataclass:引擎在函数返回后一次性用两个 mapper 拆开,时机明确(就在写 END/PREVIOUS 那一刻)、内容可序列化(frozen dataclass)。用"结构化返回值"代替"随处副作用",这与 LangGraph 全局"用通道写代替直接改状态"的哲学一致。
L07

边界:不支持生成器 + 今日小结

回看 L04 那行——entrypoint 明确拒绝生成器函数 func/__init__.py:526

# func/__init__.py:526
if inspect.isgeneratorfunction(func) or inspect.isasyncgenfunction(func):
    raise NotImplementedError("Generators are not supported in the Functional API.")
⚠️ 边界:为什么不能用 yield?你可能想在 entrypoint 里 yield 流式产出。但引擎需要函数有一个明确的返回值去写 END/PREVIOUS 通道;生成器没有"返回值"这个概念(它是惰性产出流)。要流式输出,正确姿势是用 StreamWriter(Day 57),或直接 @entrypoint().stream(...) 消费。生成器和"返回值语义"冲突,所以直接在编译期拦下,给出清晰报错而非运行时诡异行为。

👶 那 @task 和 @entrypoint 里的函数,能不能像 StateGraph 那样跨节点共享一个大 State?

👨‍🏫 不能,也不需要。函数式的"状态"就是普通 Python 局部变量——你在 entrypoint 函数体里定义的变量,天然被后面的代码读到。跨"运行"的记忆才用 previous。所以函数式适合"流程清晰、状态简单"的场景;需要复杂多字段 reducer、并行写同字段时,StateGraph 更顺手。二者可混用(task 也能在 StateGraph 节点里调)。

🧠 今天你应该能回答

  • 函数式 API 和 StateGraph 的关系?(同一个 Pregel 引擎的两种描述方式,命令式 vs 声明式)
  • @task 把函数包成了什么?(_TaskFunction 对象,带 retry/cache 策略和 clear_cache 方法)
  • 为什么调用 task 返回 Future 而不是结果?(分离"发起"与"取值",好并行 + 好插入 checkpoint/interrupt)
  • @entrypoint 把函数变成了什么?(一个只含单节点的 Pregel 图,return 的就是 Pregel 对象)
  • entrypoint 的三个隐形通道?(START Ephemeral 输入 / END LastValue 输出 / PREVIOUS LastValue 跨运行记忆)
  • entrypoint.final 解决什么?(返回值和存档值解耦:value 给调用方,save 给下次的 previous)
  • 为什么不支持生成器?(生成器没有明确返回值,无法写 END/PREVIOUS 通道)

✋ 10 分钟动手

# 1. 读三段核心源码
sed -n '59,106p'  libs/langgraph/langgraph/func/__init__.py   # _TaskFunction
sed -n '516,620p' libs/langgraph/langgraph/func/__init__.py   # entrypoint.__call__ 组装 Pregel

# 2. 亲手验证"返回值≠存档值"
python - <<'PY'
from typing import Any
from langgraph.func import entrypoint
from langgraph.checkpoint.memory import InMemorySaver
@entrypoint(checkpointer=InMemorySaver())
def wf(number: int, *, previous: Any = None) -> entrypoint.final[int, int]:
    previous = previous or 0
    return entrypoint.final(value=previous, save=2 * number)
cfg = {"configurable": {"thread_id": "t"}}
print(wf.invoke(3, cfg))   # 0
print(wf.invoke(1, cfg))   # 6
PY
明天预告 · Day 48:今天留了个悬念——config[CONF][CONFIG_KEY_CALL] 那个"提交任务"的回调是谁塞的?@task 调用又是如何变成引擎里一个真正的 PregelExecutableTask 节点的?明天钻进 pregel/_call.py_algo.py,看 func 和 pregel 的接缝。
← Day 46 replay 与幂等 Day 48 · func 与 pregel 的关系 →