Day1

第1步:LangChain是什么

1.1W2的痛点
你在W2写了⼀个完整的research_pipeline.py ――能⽤,但有问题。
W2当前的写法

import os, json 
from openai import OpenAI 
from dotenv import load_dotenv 
from pydantic import BaseModel, Field 
load_dotenv() 
client = OpenAI(api_key=os.getenv("DEEPSEEK_API_KEY"), base_url="https://api.deepseek.com") 
class ExperimentRecord(BaseModel): experiment_id: str 
started_at: str 
# ... 其他字段 
# 每次调⽤ LLM 都重复这些: 
# 1. 写 messages 列表(system + user) 
# 2. 调 client.chat.completions.create(...) # 3. 拿 response.choices[0].message.content # 4. json.loads + Pydantic 校验 
# 5. 失败时重试 / 报错 
resp = client.chat.completions.create( 
model="deepseek-chat", 
messages=[ 
{"role": "system", "content": "你是⼀个实验⽇志解析助⼿..."}, {"role": "user", "content": f"实验⽇志:\n{log_text}"}, 
], 
response_format={"type": "json_object"}, 
) 
json_str = resp.choices[0].message.content 
record = ExperimentRecord.model_validate_json(json_str) 

1.2LangChain是什么
LangChain 是一款开源、模块化的大语言模型 (LLM) 应用开发框架LangChain。

⚠️ 它本身不是大模型,不能直接聊天;它是胶水 / 流水线工具,用来把 DeepSeek、GPT‑4、Claude 这类大模型,和文件、数据库、自定义代码、搜索引擎串联起来,搭建复杂 AI 业务流程。

  • 发布时间:2022 年
  • 主流语言:Python(你一直在用)、JavaScript
  • 简单比喻:大模型 = 发动机;LangChain = 整车底盘 + 传动管线

解决什么痛点

原生直接调用大模型 API,有三个短板:

  1. 没有外部知识:不知道你的日志、PDF、项目文档
  2. 没有记忆:多轮对话会丢失上下文
  3. 不会分步干活:复杂任务(读日志→提取 JSON→生成报告),单轮 Prompt 很难一次性完成

LangChain 就是专门解决上面三类问题,实现 RAG 知识库问答、链式流水线、智能 Agent 等应用CSDN博…。

六大核心模块

  1. Models(模型层)
    统一接口调用各种大模型,你代码里面 ChatOpenAI(model="deepseek‑chat") 就属于这一块,OpenAI 兼容接口(DeepSeek、通义千问都能用)。
  2. Prompt Templates(提示词模板)
    把提示词做成模板,预留 {log_text}{format_instructions} 这类占位符,运行时再填充变量,就是你代码中 ChatPromptTemplate.from_messages
  3. Output Parsers(输出解析器)
    把大模型返回字符串,转成结构化对象;PydanticOutputParser 就是你日志解析项目用的,强制输出 JSON 并校验字段。
  4. Chains / LCEL(链,流水线)⭐你最近一直在练
    LCEL=LangChain Expression Language,用管道符 | 串联步骤,形成流水线
    ⼀句话定义:LangChain=把LLM应⽤的「拼prompt→调模型→解析输出→拿到对象」4步封装成3个可拼装的"积⽊"+⼀个 | 操作符。
    最简单的⼀句话:
    chain.invoke(input) =
    parser.invoke(model.invoke(prompt.invoke(input)))
    具体写法(W3第⼀个LangChain程序):
import os 
from dotenv import load_dotenv 
from langchain_openai import ChatOpenAI 
from langchain_core.prompts import ChatPromptTemplate 
from langchain_core.output_parsers import StrOutputParser 
load_dotenv() 
# 注:完整代码必须传 api_key + base_url(指向 DeepSeek), 
# 这⾥简化展⽰ chain 拼接——实际跑看 Step 3 的 chat_lc.py 
model = ChatOpenAI( 
model="deepseek-chat", 
api_key=os.getenv("DEEPSEEK_API_KEY"), base_url="https://api.deepseek.com", 
temperature=0.5, 
) 
prompt = ChatPromptTemplate.from_messages([("user", "解释:{concept}")]) 
parser = StrOutputParser() 
chain = prompt | model | parser          # ← 3 个积⽊拼成 1 个 chain result = chain.invoke({"concept": "Pydantic"})  # ← 1 ⾏调⽤ = 3 步执⾏ 
result = chain.invoke({"concept": "Pydantic"})
# 关键!打印输出结果
print("\n====返回结果====")
print(result)

实验结果
在这里插入图片描述

第2步:安装LangChain

pip install --only-binary=:all: langchain langchain-openai langchain-community

langchain 是核⼼包, langchain-openai 是OpenAI/DeepSeek等OpenAI-compatible模型的集成包, langchain-community 是社区维护的扩展(PDF解析、向量数据库等,W4-W7才会⽤)。
为什么加 --only-binary=:all: ?
langchain依赖 SQLAlchemy → greenlet (C扩展)。如果你的macOS没装Xcode
CommandLineTools,pip会尝试本地编译 greenlet 然后失败。 --only-binary=:all: 强制只装PyPI上传好的 .whl ⽂件(Python3.9+macOS都有universal2wheel),跳过编译。
如果没加这个flag报错:
修法1(推荐):pip加 --only-binary=:all: 重装
修法2:装XcodeCLI: xcode-select --install (会弹窗,5-10分钟)

第3步:第⼀个LangChain程序

W1Day1你写过⼀个 chat.py ,今天⽤LangChain重构。创建 chat_lc.py :

# chat_lc.py 
# Week 03 Day 1 · LangChain 版 chat.py 
import os 
from dotenv import load_dotenv 
from langchain_openai import ChatOpenAI 
from langchain_core.prompts import ChatPromptTemplate 
from langchain_core.output_parsers import StrOutputParser 
load_dotenv() 
# === W1 Day 1 的代码需要 30 ⾏(client + 拼 messages + create + 解析)=== # === W3 Day 1 的代码只要 12 ⾏(3 个抽象拼起来)=== 
# 抽象 1:ChatModel(统⼀ LLM 客⼾端) 
model = ChatOpenAI( 
model="deepseek-chat", api_key=os.getenv("DEEPSEEK_API_KEY"), base_url="https://api.deepseek.com", temperature=0.5, 
) 
# 抽象 2:PromptTemplate(参数化提⽰词) 
prompt = ChatPromptTemplate.from_messages([ 
("system", "你是⼀个科研项⽬⽂档润⾊助⼿。任务是把⽤⼾给的项⽬说明润⾊成清晰版。要求:1) ⽤ Markdown 格式 2) 保留所有具体信息 3) 不要编造代码⽰例"), 
("user", "原始草稿:\n{rough_readme}\n\n请润⾊输出:"), 
]) 
# 抽象 3:OutputParser(把 LLM 输出转成可⽤对象) # StrOutputParser 直接拿字符串内容 
parser = StrOutputParser() 
# 3 个抽象⽤ | 连起来 = ⼀个 chain chain = prompt | model | parser 
# 调⽤ 
rough_readme = "这是「K同学啊」的项⽬。功能是处理 VisDrone 数据集。跑 baseline。之前⽤ YOLOv8,现在要试试 YOLOv9。代码有点乱。GPU 是⼀张 A100。" 
chain = prompt | model | parser   
result = chain.invoke({"rough_readme": rough_readme}) 
print(result) 

在这里插入图片描述

第4步:PromptTemplate实战

W1Day2/W2Day2你⽤f-string拼prompt――能⽤,但有3个问题:
•问题1:变量名拼错不报错
•问题2:变量重复填
•问题3:不能从外部动态改
如果你的prompt来⾃配置⽂件(YAML),f-string不能解析。 PromptTemplate解决这3个问题:
创建 prompt_demo.py :

# prompt_demo.py 
# 演⽰ PromptTemplate 的能⼒ 
from langchain_core.prompts import PromptTemplate, ChatPromptTemplate 
# 1. 变量名校验——填错会⽴即报错 
pt = PromptTemplate.from_template("你是 {agent_role} 助⼿") print("合法调⽤:", pt.format(agent_role="Python ⼯程师")) 
# print(pt.format(agent_role_X="..."))  # 取消注释看错误 
# 2. 同⼀变量复⽤——只写⼀次 
pt2 = PromptTemplate.from_template("⽤⼾问题:{q}\n\n请分三步回答:\n1. 思路\n2. 代码\n3. 复⽤问题:{q} 的边界") 
print(pt2.format(q="如何检测重复⽂件")) 
# 3. ChatPromptTemplate:⽀持多⻆⾊(system + user + assistant) chat_pt = ChatPromptTemplate.from_messages([ 
("system", "你是 {domain} ⽅向的资深⼯程师"), 
("user", "我的问题是:{question}"), 
]) 
msgs = chat_pt.format_messages(domain="Python", question="如何写⼀个⽂件去重⼩⼯具") 
print("\n⻆⾊消息:") 
for m in msgs: 
print(f"  {m.type}: {m.content}")

在这里插入图片描述

第5步:OutputParser+

W2Day2你写过 parse_log.py ,今天⽤LangChain的 PydanticOutputParser 重构。创建 parse_log_lc.py :

# parse_log_lc.py 
# Week 03 Day 1 · LangChain 版 parse_log.py 
import os 
from dotenv import load_dotenv 
from langchain_openai import ChatOpenAI 
from langchain_core.prompts import ChatPromptTemplate 
from pydantic import BaseModel, Field 
from langchain_core.output_parsers import PydanticOutputParser 
# 复⽤ W2 Day 1 的 schema 
import sys 
import os 
sys.path.insert(0, os.path.join(os.path.dirname(os.path.abspath(__file__)), "..")) 
from models import ExperimentRecord 
load_dotenv() 
# 抽象 1:模型(注意:让 LLM 输出 JSON 字符串) model = ChatOpenAI( 
model="deepseek-chat", api_key=os.getenv("DEEPSEEK_API_KEY"), base_url="https://api.deepseek.com", 
temperature=0.0,  # 0 保证稳定输出 model_kwargs={"response_format": {"type": "json_object"}}, ) 
# 抽象 2:PromptTemplate(⽤ parser.get_format_instructions() ⾃动⽣成 schema 描述) 
parser = PydanticOutputParser(pydantic_object=ExperimentRecord) 
prompt = ChatPromptTemplate.from_messages([ ("system", """你是⼀个实验⽇志解析助⼿。 
从⽤⼾提供的实验⽇志中提取关键信息,**只输出⼀个 JSON 对象**(不要 ```代码块标记,不要任何解释⽂字)。 
{format_instructions} 
特别处理: 
- 时间戳⽤ ISO 8601 格式 
- 超参数要全部列出来 
- 指标按列名分类"""), 
("user", "实验⽇志:\n```\n{log_text}\n```"), 
]) 
# 抽象 3:parser 已经在 prompt ⾥调⽤ 
chain = prompt.partial(format_instructions=parser.get_format_instructions()) | model | parser 
# 调⽤ 
SAMPLE_LOG = """ 
[2026-07-01 10:23:45] Starting training... 
[2026-07-01 10:24:12] Config: dataset=VisDrone, model=YOLOv8m, lr=0.001, batch_size=16, epochs=100 
[2026-07-01 11:05:33] Epoch 100/100: mAP@0.5=0.421 mAP@0.75=0.281 [2026-07-01 11:05:33] Training finished. Total runtime: 41m48s 
""" 
record = chain.invoke({"log_text": SAMPLE_LOG}) 
print(f"✓ 解析成功:{type(record).__name__}") 
print(f"  experiment_id = {record.experiment_id}") 
print(f"  model_name = {record.model_name}") 
print(f"  status = {record.status.value}") 
print(f"  metrics = {[(m.name, m.value) for m in record.metrics]}") 
print(f"\n完整 JSON:\n{record.model_dump_json(indent=2, ensure_ascii=False)}") 

在这里插入图片描述

Day2·LCEL⾼级组合(Passthrough/ Parallel/Lambda)

第1步:LCEL⾼级抽象

W3Day1你写了 prompt | model | parser 三段链——能⽤,但遇到真实场景就⿇烦:

场景1:第⼆步需要第⼀步的输⼊+第⼀步的输出

场景2:两个分⽀并⾏(⽐如:同时让GPT-4和Claude回答,选更好的)

场景3:在chain⾥塞⾃定义函数
这3个场景对应LCEL(LangChainExpressionLanguage,LangChain表达式语⾔)的3个⾼级抽象:
•RunnablePassthrough :把原始输⼊直接传给下⼀步。
•RunnableParallel :同时跑多个chain,输出dict。
•RunnableLambda :把任何Python函数包成Runnable。

第2步:RunnablePassthrough

⽰例:你想做⼀个"翻译+原⽂对照"chain

# 朴素写法:⾃⼰写 dict 
original_text = "Hello" 
translated = translator.invoke(original_text) 
output = {"原⽂": original_text, "译⽂": translated}  # ⼿动拼

LCEL写法:⽤RunnablePassthrough.assign

代码块
from langchain_core.runnables import RunnablePassthrough 
chain = ( 
RunnablePassthrough.assign(translation=translator)  # 原⽂不动 + 加 translation 字段 
| output_formatter                                 # 拿 dict 格式化 
) 
chain.invoke({"text": "Hello"}) 
# 输出: {"text": "Hello", "translation": "你好"}

RunnablePassthrough.assign接收上⼀步的dict(或input),加上新字段,传给下⼀步。
创建 passthrough_demo.py :

# passthrough_demo.py 
# 演⽰ RunnablePassthrough 
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser 
from langchain_openai import ChatOpenAI 
from langchain_core.runnables import RunnablePassthrough 
import os 
from dotenv import load_dotenv 
load_dotenv() 
model = ChatOpenAI( 
model="deepseek-chat", api_key=os.getenv("DEEPSEEK_API_KEY"), base_url="https://api.deepseek.com", temperature=0.3, 
) 
# 演⽰ 1:让原⽂"穿过"整个 chain 
translate_prompt = ChatPromptTemplate.from_messages([ 
("system", "你是中英翻译助⼿,把英⽂翻成中⽂。"), 
("user", "{text}") 
]) 
translate_chain = translate_prompt | model | StrOutputParser() 
# ⽤ RunnablePassthrough.assign 保留原⽂ 
final_chain = ( 
RunnablePassthrough.assign(translation=translate_chain) 
| ChatPromptTemplate.from_messages([ 
("system", "你是排版助⼿,给出原⽂和译⽂的对照格式。"), 
("user", "原⽂:{text}\n译⽂:{translation}\n\n请输出 '原⽂: 译⽂' 的对照格式。") 
]) 
| model 
| StrOutputParser() 
) 
result = final_chain.invoke({"text": "Hello, world. LangChain makes LLM apps easy."}) 
print(result) 

在这里插入图片描述

第3步:RunnableParallel――多分⽀并⾏

很多场景需要同时跑多个独⽴的步骤——⽐如:
•同⼀问题问多个模型,对⽐答案
•同⼀个⽂档同时做"总结+翻译+抽取关键点"
•RAG检索:问题+检索+历史会话同时取
RunnableParallel把多个Runnable并⾏执⾏,返回dict。
创建 parallel_demo.py

# parallel_demo.py
# 演示 RunnableParallel 多分支
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_openai import ChatOpenAI
from langchain_core.runnables import RunnableParallel, RunnablePassthrough
import os
from dotenv import load_dotenv
import time

load_dotenv()

model = ChatOpenAI(
    model="deepseek-chat",
    api_key=os.getenv("DEEPSEEK_API_KEY"),
    base_url="https://api.deepseek.com",
    temperature=0.5,
)

# 三个独立的分析任务
summarize_prompt = ChatPromptTemplate.from_messages([
    ("user", "用 1 句话总结:{text}")
])
key_points_prompt = ChatPromptTemplate.from_messages([
    ("user", "列出 3 个关键点:{text}")
])
sentiment_prompt = ChatPromptTemplate.from_messages([
    ("user", "判断情感(正面/负面/中性):{text}")
])

# 串行需要 6 秒,并行只要 2 秒
summarize_chain = summarize_prompt | model | StrOutputParser()
key_points_chain = key_points_prompt | model | StrOutputParser()
sentiment_chain = sentiment_prompt | model | StrOutputParser()

# 并行 chain
parallel_chain = RunnableParallel(
    summary=summarize_chain,
    key_points=key_points_chain,
    sentiment=sentiment_chain,
)

# 测性能
text = """LangChain is a framework for developing applications powered by language models.
It enables applications that are:
- Data‑aware: connect a language model to other sources of data
- Agentic: allow a language model to interact with its environment
The main value props of LangChain are:
1. Component‑based: use modular components
2. Customizable: easily extend and customize
3. Production‑ready: deploy with confidence"""

t0 = time.time()
result = parallel_chain.invoke({"text": text})
print(f"✓ 并行完成({time.time() - t0:.1f}s)")
print(f"\n总结:{result['summary']}\n")
print(f"关键点:\n{result['key_points']}\n")
print(f"情感:{result['sentiment']}")



在这里插入图片描述

第4步:RunnableLambda――把Python函数塞进chain

# lambda_demo.py 作者:K同学啊
# 演示 RunnableLambda 在 chain 中插入 Python 函数
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import StrOutputParser
from langchain_openai import ChatOpenAI
from langchain_core.runnables import RunnableLambda
import os
import re
from dotenv import load_dotenv

load_dotenv()

model = ChatOpenAI(
    model="deepseek-chat",
    api_key=os.getenv("DEEPSEEK_API_KEY"),
    base_url="https://api.deepseek.com",
    temperature=0.3,
)

# 自定义函数 1:预处理(清洗日志)
def clean_log(log: str) -> str:
    """去掉 ANSI 颜色码 + 多余空行"""
    log = re.sub(r'\x1b\[[0-9;]*m', '', log)
    log = re.sub(r'\n\s*\n', '\n', log)
    return log.strip()

# 自定义函数 2:后处理(写文件 + 打印)
def save_to_file(content: str) -> str:
    """保存到 outputs/week03/report.md"""
    os.makedirs("outputs/week03", exist_ok=True)
    with open("outputs/week03/report.md", "w", encoding="utf-8") as f:
        f.write(content)
    print(f"✓ 已保存到 outputs/week03/report.md")
    return content

# Chain:清洗 → 调 LLM → 保存
prompt = ChatPromptTemplate.from_messages([
    ("system", "你是实验日志总结助手,输出 3 行 Markdown 摘要。"),
    ("user", "日志:\n{log}")
])

chain = (
    RunnableLambda(clean_log)              # 预处理
    | prompt
    | model
    | StrOutputParser()
    | RunnableLambda(save_to_file)         # 后处理
)

log = "\x1b[32m[2026-07-01 10:23:45]\x1b[0m Starting training...\n\n\n[2026-07-01 11:05:33] mAP=0.421\n"
chain.invoke(log)

在这里插入图片描述

Day3·LCEL多步Chain实战

第1步:5步模板

任何LLM任务都能套这5步:

chain = ( 
RunnableLambda(preprocess)         # 1. 预处理:清洗 / 读⽂件 / 查 DB 
| prompt_template                    # 2. 提⽰词模板 
| chat_model                         # 3. 调 LLM 
| output_parser                      # 4. 解析输出 
| RunnableLambda(postprocess)        # 5. 后处理:写⽂件 / 调 API / ⼊库 
) 

W2vsW3写法的差异:

# W2:按顺序执⾏什么 
data = load() 
clean_data = clean(data) resp = call_llm(clean_data) result = parse(resp) save(result) 
# W3:设计组件 + 拼装 
preprocess = RunnableLambda(clean) 
llm_step = prompt | model | parser postprocess = RunnableLambda(save) 
chain = preprocess | llm_step | postprocess 

W3写法的4个优势:每段独⽴测试、整段可复⽤、可替换、可观测。

第2步:数据契约

•多步chain⽤dict传数据——dict的key就是数据契约。
•每个chain的输⼊输出dict字段固定:

# chain 1:preprocess 
# 输⼊: {"experiment_dir": "..."} 
# 输出: {"experiment_dir": "...", "log_text": "..."} 
# chain 2:analyze 
# 输⼊: {"experiment_dir": "...", "log_text": "..."} 
# 输出: {"experiment_dir": "...", "log_text": "...", "analysis_record": 
ExperimentRecord, "analysis_text": "..."} 
# chain 3:report 
# 输⼊: {"analysis_record": ..., "analysis_text": "..."} # 输出: {"report_text": "..."}

创建 data_contract_demo.py :

# data_contract_demo.py
# 演示多步 chain 的数据契约
from langchain_core.runnables import RunnableLambda, RunnablePassthrough

# 步骤 1:preprocess
def step1_preprocess(input: dict) -> dict:
    return {
        "experiment_dir": input["experiment_dir"],
        "tools_messages": [{"role": "user", "content": f"探索 {input['experiment_dir']}"}]
    }

# 步骤 2:模拟 LLM 调用(实际是 chat_model)
def step2_explore(input: dict) -> dict:
    return {
        "experiment_dir": input["experiment_dir"],
        "tools_messages": input["tools_messages"],
        "exploration_result": "找到 log.txt 和 config.yaml"
    }

# 步骤 3:analyze
def step3_analyze(input: dict) -> dict:
    return {
        "experiment_dir": input["experiment_dir"],
        "exploration_result": input["exploration_result"],
        "analysis_record": {"experiment_id": "exp_001", "model_name": "YOLOv8m"},
        "analysis_text": "实验使用 YOLOv8m..."
    }

# 步骤 4:report
def step4_report(input: dict) -> dict:
    return {
        "report_text": f"# 实验报告\n\n{input['analysis_text']}"
    }

# 拼接
chain = (
    RunnableLambda(step1_preprocess)
    | RunnableLambda(step2_explore)
    | RunnableLambda(step3_analyze)
    | RunnableLambda(step4_report)
)

# 调试:单步执行
print("=== 单步调试 ===")
result_step1 = step1_preprocess({"experiment_dir": "outputs/exp_001"})
print(f"Step 1 输出 keys: {list(result_step1.keys())}")
result_step2 = step2_explore(result_step1)
print(f"Step 2 输出 keys: {list(result_step2.keys())}")

# 全链
print("\n=== 全链执行 ===")
final = chain.invoke({"experiment_dir": "outputs/exp_001"})
print(f"最终输出 keys: {list(final.keys())}")
print(f"report_text 前 50 字: {final['report_text'][:50]}")

在这里插入图片描述

第3步:完整3步pipeline

# pipeline_lc.py
# Week 03 Day 3 · 3 步 LCEL pipeline(纯 Chain,不用 Agent)
import os
import sys
import argparse
from dotenv import load_dotenv
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import PydanticOutputParser, StrOutputParser
from langchain_core.runnables import RunnableLambda

sys.path.insert(0, ".")
from models import ExperimentRecord

load_dotenv()

model = ChatOpenAI(
    model="deepseek-chat",
    api_key=os.getenv("DEEPSEEK_API_KEY"),
    base_url="https://api.deepseek.com",
    temperature=0.3,
)

# ============================================
# Step 1:预处理(读 log.txt)
# ============================================
def read_log(input: dict) -> dict:
    """读实验目录下的 log.txt"""
    log_path = f"{input['experiment_dir']}/log.txt"
    if not os.path.exists(log_path):
        return {**input, "log_text": ""}
    with open(log_path, "r", encoding="utf-8") as f:
        log_text = f.read()
    return {**input, "log_text": log_text}

preprocess_step = RunnableLambda(read_log)

# ============================================
# Step 2:分析(Pydantic 强校验)
# ============================================
analyze_parser = PydanticOutputParser(pydantic_object=ExperimentRecord)

analyze_chain = (
    analyze_prompt := ChatPromptTemplate.from_messages([
        ("system", "你是实验日志分析助手。\n{format_instructions}"),
        ("user", "实验日志:\n{log_text}\n\n请输出结构化数据。")
    ]).partial(format_instructions=analyze_parser.get_format_instructions())
    | model
    | analyze_parser
)

# ============================================
# Step 3:报告(流式 + 双写)
# ============================================
report_prompt = ChatPromptTemplate.from_messages([
    ("system", "你是科研报告写作助手。报告结构:1. 实验概述 2. 配置详情 3. 结果分析 4. 改进建议"),
    ("user", "请基于实验数据生成报告:\n{experiment_data}")
])

def stream_to_file(input_iter, output_path: str):
    """流式 chunk 边打印边写文件"""
    os.makedirs(os.path.dirname(output_path) or ".", exist_ok=True)
    full = ""
    with open(output_path, "w", encoding="utf-8") as f:
        for chunk in input_iter:
            print(chunk, end="", flush=True)
            f.write(chunk)
            full += chunk
    print(f"\n\n✓ 报告已保存:{output_path}({len(full)} 字符)")
    return full

# ============================================
# 完整 pipeline
# ============================================
def run_pipeline(experiment_dir: str):
    # Step 1:读 log
    print(f"[Step 1] 读 log.txt...")
    state = preprocess_step.invoke({"experiment_dir": experiment_dir})
    log_len = len(state.get("log_text", ""))
    print(f"  ✓ 读取完成({log_len} 字)")
    if log_len == 0:
        print(f"  ✗ 未找到 log.txt,跳过")
        return None

    # Step 2:分析
    print(f"[Step 2] 分析生成 ExperimentRecord...")
    record = analyze_chain.invoke({"log_text": state["log_text"]})
    print(f"  ✓ 解析成功:{record.experiment_id}")

    # 持久化
    record_path = f"{experiment_dir}/experiment_record.json"
    with open(record_path, "w", encoding="utf-8") as f:
        f.write(record.model_dump_json(indent=2, ensure_ascii=False))
    print(f"  ✓ 入库:{record_path}")

    # Step 3:报告
    print(f"[Step 3] 流式生成报告...")
    output_path = f"{experiment_dir}/REPORT.md"
    report_chain = report_prompt | model | StrOutputParser()
    stream_to_file(
        report_chain.stream({"experiment_data": record.model_dump_json(indent=2, ensure_ascii=False)}),
        output_path
    )
    return record

if __name__ == "__main__":
    parser = argparse.ArgumentParser()
    parser.add_argument("--dir", required=True, help="实验目录")
    args = parser.parse_args()
    run_pipeline(args.dir)

Day4·完整80⾏Pipeline

第1步:设计80⾏的研究pipeline

阶段项⽬=3个step,每个⽤W3不同能⼒:

[实验⽬录][Step 1: 读 log]  ← RunnableLambda 
↓ 
[Step 2: 分析]  ← PydanticOutputParser + ExperimentRecord ↓ 
[Step 3: 报告]  ← 流式 + 双写 + RunnableLambda 
↓ 
[REPORT.md + experiment_record.json]

设计原则:

代码块
# 5 步模板 + W3 能⼒分配 
chain = ( 
# Step 1: 读 log (Lambda) 
RunnableLambda(read_log) 
# Step 2: 分析 (Chain) 
| analyze_prompt | model | PydanticOutputParser 
# Step 3: 报告 (Chain + 流式) 
| report_prompt | model | StrOutputParser 
) 

第2步:完整pipeline

# research_pipeline_lc.py
# Week 03 Day 4 · 完整 LCEL pipeline(80 行 vs W2 200 行)
import os
import sys
from pathlib import Path
from dotenv import load_dotenv
from langchain_openai import ChatOpenAI
from langchain_core.prompts import ChatPromptTemplate
from langchain_core.output_parsers import PydanticOutputParser, StrOutputParser
from langchain_core.runnables import RunnableLambda

sys.path.insert(0, ".")
from models import ExperimentRecord

# ============================================
# 0. 全局配置
# ============================================
load_dotenv()

model = ChatOpenAI(
    model="deepseek-chat",
    api_key=os.getenv("DEEPSEEK_API_KEY"),
    base_url="https://api.deepseek.com",
    temperature=0.3,
)

# ============================================
# 1. 预处理:读 log.txt
# ============================================
def read_log(input: dict) -> dict:
    """读实验目录下的 log.txt"""
    log_path = f"{input['experiment_dir']}/log.txt"
    if not os.path.exists(log_path):
        return {**input, "log_text": ""}
    with open(log_path, "r", encoding="utf-8") as f:
        log_text = f.read()
    return {**input, "log_text": log_text}

# ============================================
# 2. 分析 chain(Pydantic 强校验)
# ============================================
analyze_parser = PydanticOutputParser(pydantic_object=ExperimentRecord)

analyze_prompt = ChatPromptTemplate.from_messages([
    ("system", "你是实验日志分析助手。\n{format_instructions}"),
    ("user", "实验日志:\n{log_text}\n\n请输出结构化数据。")
])

analyze_chain = (
    analyze_prompt.partial(format_instructions=analyze_parser.get_format_instructions())
    | model
    | analyze_parser
)

# ============================================
# 3. 报告 chain(流式 + 双写)
# ============================================
report_prompt = ChatPromptTemplate.from_messages([
    ("system", "你是科研报告写作助手。报告结构:1. 实验概述 2. 配置详情 3. 结果分析 4. 改进建议"),
    ("user", "请基于实验数据生成报告:\n{experiment_data}")
])

def stream_to_file(input_iter, output_path: str):
    """流式 chunk 边打印边写文件"""
    os.makedirs(os.path.dirname(output_path) or ".", exist_ok=True)
    full = ""
    with open(output_path, "w", encoding="utf-8") as f:
        for chunk in input_iter:
            print(chunk, end="", flush=True)
            f.write(chunk)
            full += chunk
    print(f"\n\n✓ 报告已保存:{output_path}({len(full)} 字符)")
    return full

# ============================================
# 4. 完整 pipeline(用 LCEL | 拼)
# ============================================
# 注意:analyze_chain 输出是 ExperimentRecord 对象(不是 dict)
# 用 RunnableLambda 把 Pydantic 转 dict,再喂给下一步
def record_to_dict(record: ExperimentRecord) -> dict:
    return {"experiment_data": record.model_dump_json(indent=2, ensure_ascii=False)}

# 完整 chain:read_log → analyze → 转 dict → report(流式)
def run_pipeline(experiment_dir: str, output_dir: str = None):
    """实验目录 → 报告 + 入库"""
    output_dir = output_dir or experiment_dir

    # Step 1+2:read_log + analyze(合并 invoke)
    print(f"[Step 1+2] 读 log + 分析...")
    state = read_log({"experiment_dir": experiment_dir})
    if not state.get("log_text"):
        print(f"  ✗ 未找到 log.txt,跳过")
        return None
    print(f"  ✓ 读取完成({len(state['log_text'])} 字)")

    record = analyze_chain.invoke({"log_text": state["log_text"]})
    print(f"  ✓ 解析成功:{record.experiment_id}")

    # 持久化
    record_path = f"{output_dir}/experiment_record.json"
    with open(record_path, "w", encoding="utf-8") as f:
        f.write(record.model_dump_json(indent=2, ensure_ascii=False))
    print(f"  ✓ 入库:{record_path}")

    # Step 3:报告(流式 + 双写)
    print(f"[Step 3] 流式生成报告...")
    report_chain = report_prompt | model | StrOutputParser()
    report_path = f"{output_dir}/REPORT.md"
    stream_to_file(
        report_chain.stream({"experiment_data": record.model_dump_json(indent=2, ensure_ascii=False)}),
        report_path
    )
    return record

if __name__ == "__main__":
    import argparse
    parser = argparse.ArgumentParser()
    parser.add_argument("--dir", required=True, help="实验目录")
    parser.add_argument("--out", default=None, help="输出目录(默认与输入相同)")
    args = parser.parse_args()
    run_pipeline(args.dir, args.out)

在这里插入图片描述

第3步:单步调试+替换

# === 单步 1:只看 read_log 结果 === 
from research_pipeline_lc import read_log 
state = read_log({"experiment_dir": "outputs/week02/sample_experiment"}) print(state["log_text"][:500]) 
# === 单步 2:测试 analyze chain 单独跑 === 
from research_pipeline_lc import analyze_chain 
record = analyze_chain.invoke({ 
"log_text": "实验使⽤ YOLOv8m 在 VisDrone 上..."  # ⼿动写假输⼊ 
}) 
print(record.experiment_id, record.model_name) 
# === 单步 3:换模型测 === 
from langchain_openai import ChatOpenAI 
gpt_model = ChatOpenAI(model="gpt-4o-mini", temperature=0) 
# 把 model 替换到 chain ⾥——其他代码不变 
test_chain = analyze_prompt.partial(format_instructions=...) | gpt_model | PydanticOutputParser(pydantic_object=ExperimentRecord) 

更多推荐