LangChain + Airflow 实战:构建自动化 AI 工作流全指南
发布时间:2026/9/20 0:59:22来源:尧图网络
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 python3.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-openaiOpenAI 模型接口的封装。如果你用的是其他模型比如本地部署的开源模型换成对应的包即可接口逻辑基本一致。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(modelgpt-4o-mini, temperature0.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 提供了PydanticOutputParserfrom 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(description3到5个关键词) sentiment: str Field(description情感倾向正面/中性/负面) parser PydanticOutputParser(pydantic_objectArticleAnalysis) prompt ChatPromptTemplate.from_messages([ (system, 你是一名专业的内容分析助手。\n{format_instructions}), (human, 请分析以下内容\n\n{raw_content}) ]) prompt prompt.partial(format_instructionsparser.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 adminexample.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(minutes5), email_on_failure: False, } dag DAG( dag_idai_content_pipeline, default_argsdefault_args, description抓取数据并用 AI 处理的自动化工作流, schedule0 8 * * *, start_datedatetime(2024, 1, 1), catchupFalse, tags[ai, langchain], ) def fetch_data(**context): resp requests.get(https://api.example.com/articles, timeout30) resp.raise_for_status() articles resp.json()[:10] context[ti].xcom_push(keyarticles, valuearticles) return len(articles) def process_with_ai(**context): articles context[ti].xcom_pull(keyarticles, task_idsfetch_data) llm ChatOpenAI(modelgpt-4o-mini, temperature0.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(keyresults, valueresults) return len(results) def save_results(**context): results context[ti].xcom_pull(keyresults, task_idsprocess_with_ai) with open(/tmp/ai_results.json, w, encodingutf-8) as f: json.dump(results, f, ensure_asciiFalse, indent2) return saved task_fetch PythonOperator(task_idfetch_data, python_callablefetch_data, dagdag) task_process PythonOperator(task_idprocess_with_ai, python_callableprocess_with_ai, dagdag) task_save PythonOperator(task_idsave_results, python_callablesave_results, dagdag) task_fetch task_process task_save4.3 任务间数据传递XCom 机制详解上面代码里反复出现的xcom_push和xcom_pull是 Airflow 的任务间通信机制叫 XComcross-communication。一个任务把数据推上去下游任务拉下来。但这里有个大坑XCom 的数据是存在 Airflow 的元数据库里的默认是 SQLite存小数据没问题但如果你推一个几百 MB 的 DataFrame 上去数据库直接爆掉。所以我的做法是XCom 只传文件路径或者小量的元数据真正的大数据写到磁盘或对象存储上下游任务从路径去读。比如上面process_with_ai如果处理结果很大就应该把结果写到/tmp/results_{date}.json然后 XCom 里只推这个路径字符串。4.4 调度策略与依赖管理schedule0 8 * * *是 cron 表达式表示每天早上 8 点跑一次。cron 表达式五个字段分别是分、时、日、月、周。几个常用例子表达式含义0 8 * * *每天早上 8 点0 */6 * * *每 6 小时一次30 9 * * 1-5工作日早上 9:300 0 1 * *每月 1 号零点catchupFalse这个参数很关键。如果设成 TrueAirflow 会把你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(parserparser, llmChatOpenAI())问题二调用超时或者速率限制。批量处理时特别容易遇到。我的做法是加一个简单的重试装饰器import time from functools import wraps def retry_with_backoff(max_retries3, base_delay2): 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 时间。你写schedule0 8 * * *它会在 UTC 8 点跑也就是北京时间下午 4 点。解决办法是在 DAG 定义里指定时区import pendulum dag DAG( dag_idai_content_pipeline, schedule0 8 * * *, start_datependulum.datetime(2024, 1, 1, tzAsia/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_failureTrue并配好 SMTP任务失败时自动发邮件。更完善的做法是接 Prometheus Grafana把任务成功率、耗时、重试次数都做成看板。我自己在实际操作中的体会是这套组合最大的价值不在于“能跑”而在于“跑得稳、看得见、改得动”。LangChain 让智能逻辑的迭代变得很快改个提示词、换个模型、加个解析器几行代码的事Airflow 让整个流程的可靠性有了保障失败了自动重试、出问题有日志可查、要加步骤就加个任务节点。两者配合基本上覆盖了从原型到小规模生产的大部分需求。如果你刚开始搭建议先把最小可跑的版本跑通——一个抓取任务、一个 AI 处理任务、一个保存任务三节点串起来确认整条链路通了再往上加东西。
网站建设高端定制企业官网