大模型训练营第三周
- 🍨 本文为🔗365天深度学习训练营 中的学习记录博客
- 🍖 原作者:K同学啊
我使用的环境是Windows、vscode、conda,使用的命令根据环境做调整
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,有三个短板:
- 没有外部知识:不知道你的日志、PDF、项目文档
- 没有记忆:多轮对话会丢失上下文
- 不会分步干活:复杂任务(读日志→提取 JSON→生成报告),单轮 Prompt 很难一次性完成
LangChain 就是专门解决上面三类问题,实现 RAG 知识库问答、链式流水线、智能 Agent 等应用CSDN博…。
六大核心模块
- Models(模型层)
统一接口调用各种大模型,你代码里面ChatOpenAI(model="deepseek‑chat")就属于这一块,OpenAI 兼容接口(DeepSeek、通义千问都能用)。 - Prompt Templates(提示词模板)
把提示词做成模板,预留{log_text}、{format_instructions}这类占位符,运行时再填充变量,就是你代码中ChatPromptTemplate.from_messages。 - Output Parsers(输出解析器)
把大模型返回字符串,转成结构化对象;PydanticOutputParser就是你日志解析项目用的,强制输出 JSON 并校验字段。 - 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)
更多推荐
所有评论(0)