第 7 章:持久化、流式与观测
checkpointer、thread_id、LangSmith 与事件流。
学习目标
理解 LangGraph 如何通过 checkpointer 和 thread_id 保存线程状态,以及如何在运行时用 astream_events 和 get_state() 观察图执行。
它是什么
checkpointer 负责把图状态按线程保存下来;thread_id 决定“哪几次调用属于同一条会话”。astream_events 提供运行时事件流,适合看节点、模型和状态更新过程;get_state() 则用于在某个时刻直接读取该线程的状态快照。LangSmith 不是另一套执行机制,而是把这些运行轨迹发到外部观测平台里。
当前项目怎么用
当前仓库的 deep_researcher 本身没有在源码里直接编译 checkpointer:
deep_researcher = deep_researcher_builder.compile()
这意味着项目核心图默认是“无持久化”的,适合本地直接跑。但 LangGraph 官方文档说明:
- 若要持久化线程记忆,需要在
compile(checkpointer=...)时提供 saver。 - 同一
thread_id的多次调用会共享状态。 get_state(config, subgraphs=True)可查看子图状态。astream_events(..., version="v3")可以流式看到运行中的事件。
当前依赖里 InMemorySaver 可直接用于学习和测试;它只保存在内存里,进程重启后数据就没了。
LangSmith 放在哪
LangSmith 是外部观测平台,不改变图本身的状态机制。只要你的模型调用链已带 LangChain/LangGraph instrumentation,并且配置了 LANGSMITH_TRACING=true 和有效 API key,同一运行就会自动出现在 LangSmith 里。这个章节默认不主动发 traces 到外部服务,只验证本地 astream_events 和 get_state()。
最小真实 Agent
示例文件:07_persistence_streaming_observability.py。
它做三件事:
- 用
InMemorySaver()编译一个只有一个模型节点的图。 - 用同一个
thread_id连续调用两次:- 第一次告诉模型“我叫老李”。
- 第二次问“我刚才说我叫什么?”。
- 在第一次调用时用
astream_events(..., version="v3")统计事件数量;第二次后用get_state(config)查看线程里累计消息数。
运行
uv run python docs/langgraph-learning/examples/07_persistence_streaming_observability.py
预期现象:
- 打印第一轮
事件数。 - 第二轮回答能回忆出名字。
持久化消息数大于 2,说明第一轮内容已保存在同一线程里。
常见误区
以为 thread_id 只是日志标签。 不是。没有相同 thread_id,同一个 checkpointer 也不会把两次调用接成同一线程。
以为 InMemorySaver 就算长期记忆。 不是。它只适合本地学习、测试和单进程调试,进程一重启全没了。
以为开了 LangSmith 才能看运行过程。 不对。astream_events 和 get_state() 本地就能用;LangSmith 是把这些轨迹发到外部平台方便检索、比较和评估。
本次真实验证
已使用默认 openai:gpt-5.5 运行一次。结果为:
事件数: 20
第二轮答复: 你刚才说你叫老李。
持久化消息数: 4
运行时出现 v3 streaming protocol on Pregel is experimental 告警,这是当前版本对 astream_events(..., version="v3") 的 beta 提示,不影响本次结果。
相关资源
查看示例代码:docs/langgraph-learning/examples/07_persistence_streaming_observability.py
"""Chapter 7: persist one thread, stream events, inspect saved state.""" import asyncio from dotenv import load_dotenv from langchain.chat_models import init_chat_model from langchain_core.messages import HumanMessage from langgraph.checkpoint.memory import InMemorySaver from langgraph.graph import END, START, MessagesState, StateGraph from langgraph.runtime import Runtime from open_deep_research.configuration import Configuration load_dotenv() configurable_model = init_chat_model( configurable_fields=("model", "max_tokens", "api_key"), ) async def reply(state: MessagesState, runtime: Runtime[Configuration]): settings = runtime.context model = configurable_model.with_config( { "configurable": { "model": settings.research_model, "max_tokens": 120, }, "tags": ["langsmith:nostream"], } ) response = await model.ainvoke(state["messages"]) return {"messages": [response]} memory = InMemorySaver() graph = ( StateGraph(MessagesState, context_schema=Configuration) .add_node("reply", reply) .add_edge(START, "reply") .add_edge("reply", END) .compile(checkpointer=memory) ) async def count_events(input_value, config, context): count = 0 stream = await graph.astream_events( input_value, config=config, context=context, version="v3", ) async for _event in stream: count += 1 return count async def main(): config = {"configurable": {"thread_id": "learning-thread-1"}} context = Configuration.from_env() first_input = { "messages": [HumanMessage(content="你好,我叫老李。请记住我的名字。")] } event_count = await count_events(first_input, config, context) second_result = await graph.ainvoke( { "messages": [HumanMessage(content="我刚才说我叫什么?只用一句中文回答。")] }, config=config, context=context, ) snapshot = graph.get_state(config) print(f"事件数: {event_count}") print("第二轮答复: " + str(second_result["messages"][-1].content)) print(f"持久化消息数: {len(snapshot.values['messages'])}") if __name__ == "__main__": asyncio.run(main())