函数式 API **函数式 API** 允许您将 LangGraph 的关键功能 (持久化, 内存, human-in-the-loop 和 流式处理) 添加到您的应用程序中,只需对现有代码进行最少的更改。
它旨在将这些功能集成到可能使用标准语言原语进行分支和流程控制的现有代码中,例如 if 语句, for 循环和函数调用。与许多需要将代码重构为明确管道或 DAG 的数据编排框架不同,函数式 API 允许您整合这些功能而无需强制执行严格的执行模型。
函数式 API 使用两个关键构建块:
- * **
@entrypoint**:将函数标记为工作流的起点,封装逻辑并管理执行流程,包括处理长时间运行的任务和中断。 - * **
@task**:表示离散的工作单元,例如 API 调用或数据处理步骤,可以在入口点内异步执行。任务返回类似 future 的对象,可以同步等待或解析。
这为构建具有状态管理和流式处理的workflows提供了一个最小化的抽象。
函数式 API 与图 API
对于喜欢更声明式方法的用户,LangGraph 的 图 API 允许您使用图范式定义工作流。两个 API 共享相同的底层运行时,因此您可以在同一应用程序中一起使用它们。
以下是一些关键区别:
- * **控制流**:函数式 API 不需要考虑图结构。您可以使用标准 Python 构造来定义工作流。这通常会减少您需要编写的代码量。
- * **短期内存**: **GraphAPI** 需要声明 **状态** 并可能需要定义 **reducer** 来管理图状态的更新。
@entrypoint和@tasks不需要显式状态管理,因为它们的状态仅限于函数内部,不会跨函数共享。 - * **检查点**:两个 API 都生成和使用检查点。在 **图 API** 中,每个 superstep后都会生成一个新的检查点。在 **函数式 API**当任务执行时,它们的结果会被保存到与给定入口点关联的现有检查点,而不是创建新的检查点。
- * **可视化**:Graph API 使得将工作流可视化为图形变得容易,这对于调试、理解工作流和与他人分享很有用。Functional API 不支持可视化,因为图形是在运行时动态生成的。
示例
下面我们演示一个简单的应用程序,它会写一篇论文并 中断 以请求人工审查。
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.func import entrypoint, task
from langgraph.types import interrupt
@task
def write_essay(topic: str) -> str:
"""Write an essay about the given topic."""
time.sleep(1) # A placeholder for a long-running task.
return f"An essay about topic: {topic}"
@entrypoint(checkpointer=InMemorySaver())
def workflow(topic: str) -> dict:
"""A simple workflow that writes an essay and asks for a review."""
essay = write_essay("cat").result()
is_approved = interrupt({
# Any json-serializable payload provided to interrupt as argument.
# It will be surfaced on the client side as an Interrupt when streaming data
# from the workflow.
"essay": essay, # The essay we want reviewed.
# We can add any additional information that we need.
# For example, introduce a key called "action" with some instructions.
"action": "Please approve/reject the essay",
})
return {
"essay": essay, # The essay that was generated
"is_approved": is_approved, # Response from HIL
}
Detailed Explanation
此工作流将撰写一篇关于"猫"主题的论文,然后暂停以获得人工审查。工作流可以无限期中断,直到提供审查。
当工作流恢复时,它从最开始执行,但由于 writeEssay 任务的结果已经被保存,任务结果将从检查点加载,而不是重新计算。
论文已写好,可以进行审查。一旦提供了审查,我们可以恢复工作流:
工作流已完成,审查已添加到论文中。
入口点
@entrypoint装饰器可用于从函数创建工作流。它封装了工作流逻辑并管理执行流程,包括处理 _长时间运行的任务_ 和 中断.
定义
An **入口点** 通过使用 @entrypoint decorator.
函数来定义 **必须接受单个位置参数**,作为工作流输入。如果您需要传递多个数据,请使用字典作为第一个参数的输入类型。
使用 entrypoint 装饰函数会产生一个Pregel][Pregel.stream]实例,它有助于管理工作流的执行(例如,处理流式传输、恢复和检查点)。
您通常需要传递一个 **检查点器** 到 @entrypoint 装饰器以启用持久化并使用 **human-in-the-loop**.
Sync
from langgraph.func import entrypoint
@entrypoint(checkpointer=checkpointer)
def my_workflow(some_input: dict) -> int:
# some logic that may involve long-running tasks like API calls,
# and may be interrupted for human-in-the-loop.
...
return result
Async
from langgraph.func import entrypoint
@entrypoint(checkpointer=checkpointer)
async def my_workflow(some_input: dict) -> int:
# some logic that may involve long-running tasks like API calls,
# and may be interrupted for human-in-the-loop
...
return result
可注入参数
在声明 entrypoint时,您可以请求访问在运行时自动注入的附加参数。这些参数包括:
| 参数 | 描述 |
|---|---|
| **previous** | 访问与之前的 checkpoint 相关的状态。请参阅 short-term-memory. |
| **存储** | [BaseStore][langgraph.store.base.BaseStore] 的一个实例。适用于 长期记忆. |
| **writer** | 用于在使用 Async Python < 3.11 时访问 StreamWriter。请参阅 使用功能 API 进行流式传输了解更多详情. |
| **config** | 用于访问运行时配置。请参阅 RunnableConfig 了解更多信息。 |
Requesting Injectable Parameters
from langchain_core.runnables import RunnableConfig
from langgraph.func import entrypoint
from langgraph.store.base import BaseStore
from langgraph.store.memory import InMemoryStore
from langgraph.checkpoint.memory import InMemorySaver
from langgraph.types import StreamWriter
in_memory_checkpointer = InMemorySaver(...)
in_memory_store = InMemoryStore(...) # An instance of InMemoryStore for long-term memory
@entrypoint(
checkpointer=in_memory_checkpointer, # Specify the checkpointer
store=in_memory_store # Specify the store
)
def my_workflow(
some_input: dict, # The input (e.g., passed via `invoke`)
*,
previous: Any = None, # For short-term memory
store: BaseStore, # For long-term memory
writer: StreamWriter, # For streaming custom data
config: RunnableConfig # For accessing the configuration passed to the entrypoint
) -> ...:
执行
使用 @entrypoint 会产生一个 Pregel 对象,可以使用 invoke, ainvoke, stream和 astream methods.
Invoke
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
my_workflow.invoke(some_input, config) # Wait for the result synchronously
Async Invoke
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
await my_workflow.ainvoke(some_input, config) # Await result asynchronously
Stream
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
stream = my_workflow.stream_events(some_input, config, version="v3")
for message in stream.messages:
for token in message.text:
print(token, end="", flush=True)
Async Stream
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
stream = await my_workflow.astream_events(some_input, config, version="v3")
async for message in stream.messages:
async for token in message.text:
print(token, end="", flush=True)
恢复
在 interrupt 后恢复执行可以通过传递一个 **resume** 值到 Command 原语。
Invoke
from langgraph.types import Command
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
my_workflow.invoke(Command(resume=some_resume_value), config)
Async Invoke
from langgraph.types import Command
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
await my_workflow.ainvoke(Command(resume=some_resume_value), config)
Stream
from langgraph.types import Command
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
stream = my_workflow.stream_events(Command(resume=some_resume_value), config, version="v3")
for message in stream.messages:
for token in message.text:
print(token, end="", flush=True)
Async Stream
from langgraph.types import Command
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
stream = await my_workflow.astream_events(Command(resume=some_resume_value), config, version="v3")
async for message in stream.messages:
async for token in message.text:
print(token, end="", flush=True)
错误后恢复
若要在错误后恢复,请运行 entrypoint 并提供 None 和相同的 **线程 ID** (config)。
这假定底层 **错误** 已得到解决,执行可以成功继续。
Invoke
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
my_workflow.invoke(None, config)
Async Invoke
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
await my_workflow.ainvoke(None, config)
Stream
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
stream = my_workflow.stream_events(None, config, version="v3")
for message in stream.messages:
for token in message.text:
print(token, end="", flush=True)
Async Stream
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
stream = await my_workflow.astream_events(None, config, version="v3")
async for message in stream.messages:
async for token in message.text:
print(token, end="", flush=True)
短期记忆
当 entrypoint 使用 checkpointer定义时,它会在同一 **线程 ID** in 检查点.
这允许使用 previous parameter.
默认情况下, previous 参数是前一次调用的返回值。
@entrypoint(checkpointer=checkpointer)
def my_workflow(number: int, *, previous: Any = None) -> int:
previous = previous or 0
return number + previous
config = {
"configurable": {
"thread_id": "some_thread_id"
}
}
my_workflow.invoke(1, config) # 1 (previous was None)
my_workflow.invoke(2, config) # 3 (previous was 1 from the previous invocation)
entrypoint.final
entrypoint.final是一个特殊的原语,可以从入口点返回,并允许 **解耦** 被保存的值 **在检查点中** 与 **入口点的返回值**.
第一个值是入口点的返回值,第二个值是将被保存在检查点中的值。类型注解是 entrypoint.final[return_type, save_type].
@entrypoint(checkpointer=checkpointer)
def my_workflow(number: int, *, previous: Any = None) -> entrypoint.final[int, int]:
previous = previous or 0
# This will return the previous value to the caller, saving
# 2 * number to the checkpoint, which will be used in the next invocation
# for the `previous` parameter.
return entrypoint.final(value=previous, save=2 * number)
config = {
"configurable": {
"thread_id": "1"
}
}
my_workflow.invoke(3, config) # 0 (previous was None)
my_workflow.invoke(1, config) # 6 (previous was 3 * 2 from the previous invocation)
任务
A **任务** 表示一个离散的工作单元,例如API调用或数据处理步骤。它有两个关键特性:
- * **异步执行**:任务被设计为异步执行,允许多个操作并发运行而不会阻塞。
- * **检查点保存**:任务结果被保存到检查点,从而能够从最后保存的状态恢复工作流。(请参阅 持久化 了解更多详情)。
定义
任务使用 @task 装饰器定义,该装饰器包装一个常规Python函数。
from langgraph.func import task
@task()
def slow_computation(input_value):
# Simulate a long-running operation
...
return result
执行
任务 只能从 **入口点**、另一个 **任务**, or a 状态图节点.
任务 _不能_ 直接从主应用程序代码调用。
当您调用 **任务**时,它会立即返回 _立即_ 一个 future 对象。future 是一个占位符,代表稍后可用的结果。
要获取 **任务**的结果,您可以同步等待(使用 result())或异步等待(使用 await).
Synchronous Invocation
@entrypoint(checkpointer=checkpointer)
def my_workflow(some_input: int) -> int:
future = slow_computation(some_input)
return future.result() # Wait for the result synchronously
Asynchronous Invocation
@entrypoint(checkpointer=checkpointer)
async def my_workflow(some_input: int) -> int:
return await slow_computation(some_input) # Await result asynchronously
何时使用任务
任务 在以下场景中很有用:
- * **检查点**:当您需要将长时间运行操作的结果保存到检查点时,这样您就不需要在恢复工作流时重新计算。
- * **Human-in-the-loop**:如果您正在构建需要人工干预的工作流,您必须使用 **任务** 来封装任何随机性(例如 API 调用),以确保工作流能够正确恢复。请参阅 确定性 部分了解更多详情。
- * **并行执行**: For I/O-bound tasks, **任务** 支持并行执行,允许多个操作同时运行而不会阻塞(例如调用多个 API)。
- * **可观测性**:将操作封装在 **任务** 中提供了一种跟踪工作流进度和监控各个操作执行的方法,使用 LangSmith.
- * **可重试工作**:当工作需要重试以处理失败或不一致时, **任务** 提供了一种封装和管理重试逻辑的方法。
序列化
LangGraph 中的序列化有两个关键方面:
entrypoint输入和输出必须是 JSON 可序列化的。task输出必须是 JSON 可序列化的。
这些要求是启用检查点和工作流恢复的必要条件。请使用 Python 原生类型(如字典、列表、字符串、数字和布尔值)来确保输入和输出可序列化。
序列化确保工作流状态(如任务结果和中间值)可以被可靠地保存和恢复。这对于实现人机交互、容错和并行执行至关重要。
如果提供不可序列化的输入或输出,在配置了检查点的工作流中会导致运行时错误。
确定性
恢复工作流运行时,代码会 **NOT** 从执行停止的 **同一行代码** 处恢复执行。执行会返回到检查点边界,然后工作流 **重放** 直到再次到达暂停点。
对于函数式 API,重放从 **入口点** 的开头开始,而 LangGraph 会从检查点恢复已完成的 **任务** 和 **子图** 结果,而不是重新计算它们。这会保留暂停期间记录的步骤顺序,包括长时间运行或非确定性的 **任务** outputs.
要使用 **human-in-the-loop**等功能,必须将非确定性工作(如随机值)和副作用(如文件写入或 API 调用)放在 **任务**.
工作流的不同运行可能产生不同的结果,但恢复 **特定** 线程应该重放相同的持久化 **任务** 和 **子图** results.
为确保工作流具有确定性并可以一致地重放,请遵循以下准则:
- * **避免重复工作**: In an **入口点**中,如果要链接多个副作用(如日志记录、文件写入或网络调用),请为每个操作创建单独的 **任务** ,这样恢复时会从检查点还原其输出,而不会再次运行。
- * **封装非确定性操作**:将可能在不同尝试之间变化的值(如随机数或系统时钟读取)保留在 **任务**内部,这样重放才能与检查点记录对齐。
- * **使用幂等操作**:关于部分任务失败和重试,请参阅 幂等性.
幂等性
幂等性确保多次运行同一操作会产生相同的结果。这有助于防止因失败而重新运行步骤时产生重复的 API 调用和冗余处理。请始终将 API 调用放在 **任务** 用于检查点的函数,并设计它们在重新执行时具有幂等性。 这对于导致数据写入的操作尤为重要。 当工作流恢复时,LangGraph 会重放已完成的 **任务** 结果(从检查点)。一个 **任务** 已启动但未完成的任务可能会在恢复时再次运行,因此设计副作用时应具有幂等性。使用幂等性密钥或验证现有结果以避免意外重复。
常见陷阱
处理副作用
将副作用(例如写入文件、发送电子邮件)封装在任务中,以确保恢复工作流时不会多次执行。
Incorrect
在此示例中,副作用(写入文件)直接包含在工作流中,因此在恢复工作流时将第二次执行它。
@entrypoint(checkpointer=checkpointer)
def my_workflow(inputs: dict) -> int:
# This code will be executed a second time when resuming the workflow.
# Which is likely not what you want.
with open("output.txt", "w") as f: # [!code highlight]
f.write("Side effect executed") # [!code highlight]
value = interrupt("question")
return value
Correct
在此示例中,副作用被封装在任务中,确保恢复时的一致执行。
from langgraph.func import task
@task # [!code highlight]
def write_to_file(): # [!code highlight]
with open("output.txt", "w") as f:
f.write("Side effect executed")
@entrypoint(checkpointer=checkpointer)
def my_workflow(inputs: dict) -> int:
# The side effect is now encapsulated in a task.
write_to_file().result()
value = interrupt("question")
return value
非确定性控制流
可能每次给出不同结果的操作(如获取当前时间或随机数)应被封装在任务中,以确保恢复时返回相同的结果。
- * 在任务中:获取随机数 (5) → 中断 → 恢复 → (再次返回 5)→ ...
- * 不在任务中:获取随机数 (5) → 中断 → 恢复 → 获取新的随机数 (7) → ...
这在使用时尤为重要 **human-in-the-loop** workflows with multiple interrupt calls. LangGraph keeps a list of resume values for each task/entrypoint. When an interrupt is encountered, it's matched with the corresponding resume value. This matching is strictly **index-based**,因此恢复值的顺序应与中断的顺序相匹配。
如果恢复时未保持执行顺序,一个 interrupt 调用可能会与错误的 resume 值匹配,从而导致错误的结果。
请阅读关于 确定性 的章节以获取更多详细信息。
Incorrect
在此示例中,工作流使用当前时间来确定要执行的任务。这是非确定性的,因为工作流的结果取决于其执行的时间。
from langgraph.func import entrypoint
@entrypoint(checkpointer=checkpointer)
def my_workflow(inputs: dict) -> int:
t0 = inputs["t0"]
t1 = time.time() # [!code highlight]
delta_t = t1 - t0
if delta_t > 1:
result = slow_task(1).result()
value = interrupt("question")
else:
result = slow_task(2).result()
value = interrupt("question")
return {
"result": result,
"value": value
}
Correct
在此示例中,工作流使用输入 t0 来确定要执行的任务。这是确定性的,因为工作流的结果仅取决于输入。
from langgraph.func import task
@task # [!code highlight]
def get_time() -> float: # [!code highlight]
return time.time()
@entrypoint(checkpointer=checkpointer)
def my_workflow(inputs: dict) -> int:
t0 = inputs["t0"]
t1 = get_time().result() # [!code highlight]
delta_t = t1 - t0
if delta_t > 1:
result = slow_task(1).result()
value = interrupt("question")
else:
result = slow_task(2).result()
value = interrupt("question")
return {
"result": result,
"value": value
}