1. 为什么我最终选择用 LangChain + Airflow 搭一套 AI 工作流
先说结论:如果你手头有一堆零散的 AI 调用——比如每天定时抓一批数据、丢给大模型做摘要、再自动生成报告发出去——那用 LangChain 管“智能逻辑”、用 Airflow 管“调度和依赖”,是目前最省心、也最容易扩展的组合。我从去年开始陆续帮几个团队落地过类似的东西,踩过的坑不算少,今天就把整套流程从头到尾拆一遍,代码和操作指引都给到,你照着抄基本能跑通。
这套方案解决的核心问题其实就一个:把“人手动点一下、等结果、再点下一步”的重复劳动,变成一条自动流转的流水线。适合谁看?会一点 Python、听说过 LangChain 但没真正跑起来、或者用过扣子工作流、dify 工作流这类可视化工具但觉得不够灵活、想自己写代码控制细节的人。如果你完全没碰过 Python,也没关系,我会把安装和环境配置的部分写得足够细。
先明确一下两个东西的分工,很多人一开始会搞混:
- LangChain:负责“智能”那一层。它把大模型的调用、提示词模板、输出解析、工具调用、记忆管理这些东西封装成可组合的模块。你可以理解成它是流水线上的“加工工位”。
- Airflow:负责“调度”那一层。它管的是任务什么时候跑、任务之间谁依赖谁、失败了重试几次、跑完发不发通知。它是流水线的“总控室”。
有人会问,那 LangGraph 呢?LangGraph 和 LangChain 的区别在于,LangGraph 更适合做有状态、有循环、需要人工介入(human in the loop)的复杂 agent 流程,而 LangChain 更适合线性的、步骤清晰的链式调用。我这条工作流大部分是线性步骤,偶尔有分支,用 LangChain 的 LCEL 表达式就够了,没必要上 LangGraph 增加复杂度。等你需要“模型自己决定下一步调哪个工具、调完再回来继续想”这种循环逻辑时,再切到 LangGraph 不迟。
至于为什么不用扣子工作流、dify 工作流这种可视化平台?它们上手确实快,拖拖拽拽就能出东西,但一旦你要接自己的私有数据、要写自定义的解析逻辑、要做复杂的错误处理,可视化平台就会开始“卡脖子”。代码方案的自由度是它们比不了的。当然,如果你只是做个简单的 markdown 转 word 工作流或者简历筛选工作流,可视化平台完全够用,不必杀鸡用牛刀。
2. 环境准备:Python 安装到依赖配置的完整避坑指南
2.1 Python 版本选择与安装
这一步看着简单,但新手翻车率极高。我的建议是:直接用 Python 3.10 或 3.11。别用 3.12 以上的最新版,因为 LangChain 生态里有些依赖包对最新版 Python 的兼容性还没跟上,你会遇到各种莫名其妙的编译错误。也别用 3.8 以下,太老了,很多新特性用不了。
Windows 用户去官网下载安装包时,务必勾选“Add Python to PATH”,这个选项不勾,后面在命令行里敲 python 会提示找不到命令,很多人卡在这。macOS 用户如果用 Homebrew,直接brew install python@3.11就行。Linux 用户一般系统自带,但版本可能偏低,建议用 pyenv 管理多版本。
安装完验证一下:
python --version pip --version两条命令都能正常输出版本号,说明基础环境 OK。如果 pip 版本太老,先升级一下:
python -m pip install --upgrade pip2.2 虚拟环境:别偷懒,一定要建
我见过太多人把所有包装在全局环境里,结果项目 A 和项目 B 的依赖版本打架,最后谁也跑不起来。虚拟环境就是给每个项目一个独立的“房间”,互不干扰。
# 创建虚拟环境 python -m venv ai_workflow_env # 激活(Windows) ai_workflow_env\Scripts\activate # 激活(macOS / Linux) source ai_workflow_env/bin/activate激活后命令行前面会出现(ai_workflow_env)的标识,说明你已经在虚拟环境里了。后面所有 pip 安装都在这上面操作。
2.3 核心依赖安装
这条工作流需要的核心包如下,我按重要程度排了序:
pip install langchain langchain-openai langchain-community pip install apache-airflow pip install pandas requests python-dotenv逐个说一下它们的作用:
langchain:核心框架,提供链、提示词模板、输出解析器等基础组件。langchain-openai:OpenAI 模型接口的封装。如果你用的是其他模型(比如本地部署的开源模型),换成对应的包即可,接口逻辑基本一致。langchain-community:社区贡献的各种集成,比如文档加载器、向量库连接器等。apache-airflow:调度框架。注意 Airflow 的安装比较重,它会带一堆依赖,装的时候耐心等。pandas:数据处理,抓下来的数据总得清洗一下。requests:发 HTTP 请求,抓数据用。python-dotenv:管理 API Key 等敏感配置,别把密钥硬编码在代码里。
注意:Airflow 在 Windows 上原生支持不太好,官方推荐在 Linux 或 macOS 上跑,或者用 WSL。如果你只有 Windows,建议装个 WSL2,在 Ubuntu 环境里操作,会顺畅很多。我早期在 Windows 上直接装 Airflow,遇到过数据库初始化失败、调度器起不来等一堆问题,换到 WSL 后一次通过。
2.4 API Key 的安全管理
建一个.env文件放在项目根目录:
OPENAI_API_KEY=你的密钥 OPENAI_BASE_URL=你的接口地址然后在代码里用python-dotenv加载:
from dotenv import load_dotenv import os load_dotenv() api_key = os.getenv("OPENAI_API_KEY")这样做的好处是,代码可以提交到 Git 仓库,但.env文件加到.gitignore里,密钥不会泄露。我见过有人直接把密钥写在代码里然后推到公开仓库,结果被人盗刷,这个教训一定要记住。
3. 用 LangChain 搭建智能处理链的核心细节
3.1 提示词模板的设计逻辑
LangChain 的提示词模板(PromptTemplate)不是简单的字符串拼接,它的价值在于把变量和固定指令分离,让同一条链可以处理不同的输入。举个例子:
from langchain_openai import ChatOpenAI from langchain_core.prompts import ChatPromptTemplate from langchain_core.output_parsers import StrOutputParser llm = ChatOpenAI(model="gpt-4o-mini", temperature=0.3) prompt = ChatPromptTemplate.from_messages([ ("system", "你是一名专业的内容分析助手,擅长从原始文本中提取关键信息并结构化输出。"), ("human", "请对以下内容进行摘要,并提取3到5个关键词:\n\n{raw_content}") ]) chain = prompt | llm | StrOutputParser()这里有几个设计决策值得说明:
为什么用ChatPromptTemplate而不是老的PromptTemplate?因为现在主流模型都是对话式的,system 和 human 角色分开写,模型对指令的遵循度更高。老的PromptTemplate把所有内容揉成一段,效果会差一些。
为什么temperature设成 0.3?摘要和提取任务需要稳定、可复现的输出,温度太高会让每次结果都不一样。如果是创意写作类任务,可以调到 0.7 以上。这个参数的本质是控制模型输出的随机性,0 最确定,1 最随机。
StrOutputParser是干嘛的?模型返回的是一个消息对象,里面包含元数据、token 统计等信息。解析器负责把真正有用的文本内容抽出来,变成纯字符串,方便后续处理。
3.2 输出解析:让模型返回结构化数据
纯文本摘要有时候不够用,你需要模型返回 JSON 格式的结构化数据,方便程序进一步处理。LangChain 提供了PydanticOutputParser:
from langchain_core.output_parsers import PydanticOutputParser from pydantic import BaseModel, Field from typing import List class ArticleAnalysis(BaseModel): summary: str = Field(description="文章摘要,不超过200字") keywords: List[str] = Field(description="3到5个关键词") sentiment: str = Field(description="情感倾向:正面/中性/负面") parser = PydanticOutputParser(pydantic_object=ArticleAnalysis) prompt = ChatPromptTemplate.from_messages([ ("system", "你是一名专业的内容分析助手。\n{format_instructions}"), ("human", "请分析以下内容:\n\n{raw_content}") ]) prompt = prompt.partial(format_instructions=parser.get_format_instructions()) chain = prompt | llm | parser这样chain.invoke({"raw_content": "..."})返回的就是一个ArticleAnalysis对象,你可以直接.summary、.keywords访问字段。关键点在于get_format_instructions(),它会把 Pydantic 模型的字段定义自动转成一段说明文字,塞进提示词里告诉模型该返回什么格式。这比手写“请返回 JSON,包含 summary、keywords、sentiment 三个字段”要可靠得多。
实操心得:即使有格式说明,模型偶尔还是会返回不合法的 JSON。建议在解析器外面包一层重试逻辑,或者在提示词里加一句“只返回 JSON,不要有任何其他文字”。我一般会在 system 消息末尾强调这一点,能显著降低解析失败率。
3.3 把多个处理步骤串成链
实际工作流往往不止一步。比如先摘要,再根据摘要生成一份简报,最后翻译成英文。用 LCEL 的管道符|可以很优雅地串起来:
summary_chain = summary_prompt | llm | StrOutputParser() report_chain = report_prompt | llm | StrOutputParser() translate_chain = translate_prompt | llm | StrOutputParser() full_chain = summary_chain | (lambda x: {"summary": x}) | report_chain | (lambda x: {"report": x}) | translate_chain这里的lambda是做一个数据格式的转换,因为每一步的输入输出字段名可能不一样,需要对齐。这种写法比传统的SequentialChain更直观,也更灵活。LCEL 的好处是它天然支持流式输出、异步调用、批量处理,而且每一步的输入输出都是透明的,调试起来方便。
4. 用 Airflow 编排调度:从 DAG 定义到任务依赖
4.1 Airflow 初始化与目录结构
Airflow 装好后,先初始化数据库:
airflow db init然后创建管理员账户:
airflow users create \ --username admin \ --firstname Admin \ --lastname User \ --role Admin \ --email admin@example.comAirflow 默认会在~/airflow目录下生成配置文件和 DAG 目录。你的工作流定义文件(DAG 文件)就放在~/airflow/dags/下面。每个 DAG 文件就是一个 Python 脚本,Airflow 会自动扫描这个目录并加载。
启动两个核心服务:
# 启动调度器 airflow scheduler # 启动 Web 界面(另开一个终端) airflow webserver --port 8080浏览器打开localhost:8080,用刚才创建的账户登录,就能看到 DAG 列表了。
4.2 定义一个完整的 DAG
下面是一个真实可跑的 DAG 示例,包含数据抓取、AI 处理、结果保存三个任务:
from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator import requests import json from langchain_openai import ChatOpenAI from langchain_core.prompts import ChatPromptTemplate from langchain_core.output_parsers import StrOutputParser from dotenv import load_dotenv load_dotenv() default_args = { "owner": "ai_workflow", "retries": 2, "retry_delay": timedelta(minutes=5), "email_on_failure": False, } dag = DAG( dag_id="ai_content_pipeline", default_args=default_args, description="抓取数据并用 AI 处理的自动化工作流", schedule="0 8 * * *", start_date=datetime(2024, 1, 1), catchup=False, tags=["ai", "langchain"], ) def fetch_data(**context): resp = requests.get("https://api.example.com/articles", timeout=30) resp.raise_for_status() articles = resp.json()[:10] context["ti"].xcom_push(key="articles", value=articles) return len(articles) def process_with_ai(**context): articles = context["ti"].xcom_pull(key="articles", task_ids="fetch_data") llm = ChatOpenAI(model="gpt-4o-mini", temperature=0.3) prompt = ChatPromptTemplate.from_messages([ ("system", "你是一名内容分析助手,请对以下文章进行摘要和关键词提取。"), ("human", "{content}") ]) chain = prompt | llm | StrOutputParser() results = [] for article in articles: output = chain.invoke({"content": article["body"]}) results.append({"title": article["title"], "analysis": output}) context["ti"].xcom_push(key="results", value=results) return len(results) def save_results(**context): results = context["ti"].xcom_pull(key="results", task_ids="process_with_ai") with open("/tmp/ai_results.json", "w", encoding="utf-8") as f: json.dump(results, f, ensure_ascii=False, indent=2) return "saved" task_fetch = PythonOperator(task_id="fetch_data", python_callable=fetch_data, dag=dag) task_process = PythonOperator(task_id="process_with_ai", python_callable=process_with_ai, dag=dag) task_save = PythonOperator(task_id="save_results", python_callable=save_results, dag=dag) task_fetch >> task_process >> task_save4.3 任务间数据传递:XCom 机制详解
上面代码里反复出现的xcom_push和xcom_pull是 Airflow 的任务间通信机制,叫 XCom(cross-communication)。一个任务把数据推上去,下游任务拉下来。
但这里有个大坑:XCom 的数据是存在 Airflow 的元数据库里的,默认是 SQLite,存小数据没问题,但如果你推一个几百 MB 的 DataFrame 上去,数据库直接爆掉。所以我的做法是:XCom 只传文件路径或者小量的元数据,真正的大数据写到磁盘或对象存储上,下游任务从路径去读。
比如上面process_with_ai如果处理结果很大,就应该把结果写到/tmp/results_{date}.json,然后 XCom 里只推这个路径字符串。
4.4 调度策略与依赖管理
schedule="0 8 * * *"是 cron 表达式,表示每天早上 8 点跑一次。cron 表达式五个字段分别是分、时、日、月、周。几个常用例子:
| 表达式 | 含义 |
|---|---|
0 8 * * * | 每天早上 8 点 |
0 */6 * * * | 每 6 小时一次 |
30 9 * * 1-5 | 工作日早上 9:30 |
0 0 1 * * | 每月 1 号零点 |
catchup=False这个参数很关键。如果设成 True,Airflow 会把你start_date到当前时间之间所有“错过”的调度都补跑一遍。比如你 start_date 设的是半年前,那它一启动就会补跑 180 次,直接把你的 API 额度刷爆。所以除非你确实需要回补历史数据,否则一律设 False。
任务依赖用>>表示,task_fetch >> task_process >> task_save就是串行执行。如果需要并行,可以这样:
task_fetch >> [task_process_a, task_process_b] >> task_save这表示 fetch 完成后,process_a 和 process_b 同时跑,两个都跑完了才执行 save。
5. 常见问题与排查技巧实录
5.1 LangChain 相关的高频问题
问题一:模型返回的内容解析失败,报 JSONDecodeError。
这是最常见的。原因通常是模型在 JSON 外面包了 ```json 这样的 markdown 代码块标记,或者加了一句“好的,以下是分析结果”。解决办法有三个:一是在提示词里明确要求“只返回 JSON,不要任何额外文字”;二是用StrOutputParser先拿到文本,再用正则把 JSON 部分抠出来;三是用 LangChain 的OutputFixingParser,它会在解析失败时自动把错误信息发回给模型,让模型修正后重新输出。
from langchain.output_parsers import OutputFixingParser from langchain_openai import ChatOpenAI fixing_parser = OutputFixingParser.from_llm(parser=parser, llm=ChatOpenAI())问题二:调用超时或者速率限制。
批量处理时特别容易遇到。我的做法是加一个简单的重试装饰器:
import time from functools import wraps def retry_with_backoff(max_retries=3, base_delay=2): def decorator(func): @wraps(func) def wrapper(*args, **kwargs): for attempt in range(max_retries): try: return func(*args, **kwargs) except Exception as e: if attempt == max_retries - 1: raise delay = base_delay * (2 ** attempt) time.sleep(delay) return wrapper return decorator指数退避的意思是第一次等 2 秒,第二次等 4 秒,第三次等 8 秒。这样既不会频繁冲击接口,又能给服务端足够的恢复时间。
5.2 Airflow 相关的踩坑记录
问题一:DAG 文件不显示在 Web 界面里。
排查顺序:先确认文件确实在~/airflow/dags/目录下;然后检查文件里有没有语法错误,Airflow 解析失败是不会报错的,只会静默忽略;再检查dag_id有没有和已有的重复。最直接的办法是在命令行跑airflow dags list,看能不能列出来。如果列不出来,就是文件本身有问题。
问题二:任务一直卡在 running 状态。
大概率是调度器没在跑,或者卡死了。先确认airflow scheduler进程还在。如果用的是 SQLite 作为元数据库,并发稍微高一点就容易锁库,建议开发阶段用 SQLite,生产环境换成 PostgreSQL。我早期用 SQLite 跑十几个任务并行,经常出现数据库锁死,换 PostgreSQL 后再没遇到过。
问题三:时区问题导致调度时间不对。
Airflow 默认用 UTC 时间。你写schedule="0 8 * * *",它会在 UTC 8 点跑,也就是北京时间下午 4 点。解决办法是在 DAG 定义里指定时区:
import pendulum dag = DAG( dag_id="ai_content_pipeline", schedule="0 8 * * *", start_date=pendulum.datetime(2024, 1, 1, tz="Asia/Shanghai"), ... )这样 8 点就是北京时间 8 点。
5.3 常见问题速查表
| 问题现象 | 可能原因 | 解决方法 |
|---|---|---|
| 模型输出解析失败 | 返回内容含额外文字或代码块标记 | 提示词强调纯 JSON,或用 OutputFixingParser |
| 调用频繁超时 | 触发速率限制 | 加指数退避重试,降低并发 |
| DAG 不显示 | 文件语法错误或路径不对 | 用airflow dags list排查 |
| 任务卡 running | 调度器未运行或数据库锁 | 检查 scheduler 进程,换 PostgreSQL |
| 调度时间偏移 | 默认 UTC 时区 | 用 pendulum 指定本地时区 |
| XCom 数据丢失 | 数据量过大导致数据库写入失败 | 只传路径,大数据写磁盘 |
| 依赖包冲突 | 全局环境装太多包 | 每个项目独立虚拟环境 |
独家避坑技巧:Airflow 的日志默认存在
~/airflow/logs/下,按 DAG ID 和任务 ID 分目录。任务失败时第一时间去看日志,90% 的问题日志里都有明确报错。另外,开发阶段把retries设小一点(比如 1 次),不然一个 bug 会重试好几次才暴露出来,浪费时间。
6. 工作流扩展:从单机到生产级的演进思路
跑通基础版本之后,你可能会想扩展。几个方向供参考。
接入本地知识库做问答。用 LangChain 的文档加载器把 PDF、Word、Markdown 读进来,切块后用向量库(比如 Chroma 或 FAISS)存起来,检索时先查向量库再把相关片段塞进提示词。这就是 RAG 的基本套路。Airflow 这边可以加一个定时任务,每天增量更新知识库。
加入人工审核环节。有些场景模型输出不能直接发布,需要人看一眼。LangGraph 的 human in the loop 就是干这个的——流程走到某一步暂停,等人确认后再继续。Airflow 里可以用Sensor或者TriggerDagRunOperator配合外部信号来实现类似效果。
多平台数据聚合。比如跨境电商场景,需要从多个平台抓订单数据,每个平台一个抓取任务,并行跑,最后汇总。Airflow 的动态任务映射(Dynamic Task Mapping)可以很优雅地处理这种“数量不确定的并行任务”。
监控与告警。Airflow 支持配置邮件告警、Slack 告警等。在default_args里设email_on_failure=True并配好 SMTP,任务失败时自动发邮件。更完善的做法是接 Prometheus + Grafana,把任务成功率、耗时、重试次数都做成看板。
我自己在实际操作中的体会是,这套组合最大的价值不在于“能跑”,而在于“跑得稳、看得见、改得动”。LangChain 让智能逻辑的迭代变得很快,改个提示词、换个模型、加个解析器,几行代码的事;Airflow 让整个流程的可靠性有了保障,失败了自动重试、出问题有日志可查、要加步骤就加个任务节点。两者配合,基本上覆盖了从原型到小规模生产的大部分需求。如果你刚开始搭,建议先把最小可跑的版本跑通——一个抓取任务、一个 AI 处理任务、一个保存任务,三节点串起来,确认整条链路通了,再往上加东西。