ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Apache Airflow 完整指南:从零搭建一条可靠数据管道的 7 步实战路径

Apache Airflow 完整指南:从零搭建一条可靠数据管道的 7 步实战路径 Apache Airflow 完整指南从零搭建一条可靠数据管道的 7 步实战路径【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guides周日下午两点你被拉进一个线上群昨天的经营日报缺数了。顺着排查下去发现抽取订单 → 清洗去重 → 汇总入库这段活儿是三个 shell 脚本加 crontab 串起来的——中间挂了没人知道没有重试也没有任何通知。这类靠脚本和定时器堆出来的数据管道撑不起一个正经的数据团队。这正是Apache Airflow这个数据工作流编排平台要解决的问题把数据管道写成 Python 代码由它统一负责任务调度、依赖编排、失败重试和状态监控。本文按一个真实故障场景 → 最小可运行 DAG → 可靠化 → 资源控制 → 进阶玩法的顺序带你从零搭一条能上生产的管道。1. 从一个凌晨告警说起为什么裸脚本撑不住上面那个故障暴露了三个典型问题依赖靠人肉对齐脚本 A 跑完才手动/定时跑脚本 B一旦 A 延迟B 就会拿到半截数据失败无感知进程退出码非 0 之后什么都没有第二天看报表才发现问题无法回补想重跑昨天的数据只能凭记忆重放命令。Airflow 的解法是把管道 代码每个步骤是一个Task任务节点步骤之间的先后关系写成有向无环图DAG调度器Scheduler按图驱动执行并记录每一次DAG Run一次完整的运行实例和每个任务实例Task Instance的状态。失败自动重试、成功/失败可挂回调、历史运行随时可回看——这些能力是裸脚本给不了的。2. 最小可运行 DAG先把核心概念对上号写第一个 DAG 之前先花两分钟把这几个词对上号后面所有文档你都能读懂概念一句话解释类比DAG一次工作流的完整定义节点 依赖边一张施工图纸Task图中的一个节点一次原子操作一个工序Operator生成 Task 的任务类如BashOperator、PythonOperator工序的模板Hook与外部系统数据库、API交互的连接层专用扳手Sensor不干活、只轮询等待某条件成立的特殊任务等货到的看门人Pool给一组任务设的并发上限车间的工位数量下面是一个完整的、能直接跑的最小 DAG——三个任务、一条依赖链约 10 行from datetime import datetime, timedelta from airflow import DAG from airflow.operators.bash import BashOperator default_args {retries: 2, retry_delay: timedelta(minutes5)} with DAG(daily_order_pipeline, start_datedatetime(2026, 1, 1), scheduledaily, catchupFalse, default_argsdefault_args) as dag: fetch_orders BashOperator(task_idfetch_orders, bash_commandpython extract.py) clean_dedupe BashOperator(task_idclean_dedupe, bash_commandpython clean.py) load_summary BashOperator(task_idload_summary, bash_commandpython load.py, poolwarehouse_write_pool) fetch_orders clean_dedupe load_summary几个值得注意的细节是最直观的依赖写法等价于set_downstream/set_upstream链式写起来最清爽catchupFalse表示启动后不会把历史缺失的日期一次性补齐对新手更友好scheduledaily是 Airflow 内置的 cron 预设之一适合固定节奏的管道。3. 让管道真正跑起来四步启动与验证第 1 步安装并初始化元数据库通常用 PostgreSQL本机学习可用 SQLitepip install apache-airflow airflow db migrate第 2 步把 DAG 文件放进dags_folder默认是~/airflow/dags/文件名比如daily_order_pipeline.py。第 3 步分别启动调度器和 Web 服务开发期可后台挂起airflow scheduler airflow webserver第 4 步浏览器打开 Web UI确认三件事DAG 出现在列表里且状态为Active灰色Paused表示已暂停手动点Trigger DAG进入Grid View看三个任务依次变绿点进某个任务看Log页签确认你的extract.py输出出现在这里。⚠️ 坑DAG 文件有语法错误时调度器不会崩溃只是 UI 里看不到这条 DAG。去Admin → Jobs里查 Scheduler 的日志那里会有解析报错。这是新手最常卡住的一步先学会看它。另外记住一个默认值所有任务默认跑在default_pool里该池默认 128 个并发槽位。这决定了你的机器上最多同时有多少任务实例在执行后面讲 Pool 时会用到它。4. 可靠化改造如何配置重试与失败告警跑通之后管道会频繁失败——这不是坏事是提醒你该做容错了。生产管道的可靠化分三层第一层任务级重试。把网络抖动、限流、锁竞争这类可恢复失败交给 Airflow 自动重跑default_args { retries: 3, retry_delay: timedelta(minutes5), retry_exponential_backoff: True, # 重试间隔翻倍别把下游打得更狠 email_on_failure: True, on_failure_callback: notify_oncall, # 自写回调可发 Slack / 钉钉 / 飞书 }第二层DAG 级兜底。除了on_failure_callbackDAG 参数里还支持on_success_callback、on_execution_fail_callback整个 DAG Run 失败时触发。把重试 3 次还失败的人找出来比把每次抖动都拉群更有价值。第三层任务幂等。这一层最容易被忽略 最佳实践重试的前提是任务幂等——同一个 DAG Run 重跑多次对数据的影响和跑一次相同。写 ETL 时优先用按分区覆盖写而不是追加写否则重试一次重复数据就多一份。5. 给数据加一道闸数据质量检查放在哪一步管道跑得完不等于跑得对。质量检查任务Quality Gate的位置直接决定坏数据能传多远抽取后、清洗前拦截上游源数据缺失、空表、行数骤降入库前拦截主键冲突、字段越界——这是最常见的落点跨表比对核对订单表与支付表的金额对账。工具选型不必纠结常见四个都有现成玩法检查工具检查定义形式接入方式适合谁Great ExpectationsJSON 期望套件专用 Operator想系统化沉淀期望、要检查报告SodaYAML 检查项BashOperator跑 CLI想轻量快速起步dbt testSQL 测试BashOperator/PythonOperator跑dbt test已在用 dbt 的团队SQL Check OperatorPython 字典 SQL直接挂任务一两条简单检查不想引新依赖把检查任务插到依赖链中间如clean_dedupe quality_gate load_summary检查失败时下游自动不跑坏数据就被挡在仓库外面了。这是可靠性和正确性的分界线值得尽早加。6. Executor 选型单机还是分布式一张表讲清任务真正在哪台机器上跑由Executor决定。它是你扩容时的第一决策点Executor执行形态典型场景主要约束SequentialExecutor单进程串行本机调试 DAG 逻辑完全不并行别上生产LocalExecutor单机多进程小团队、一台服务器单机是单点机器挂全停CeleryExecutor分布式 worker 池经 Redis/RabbitMQ 消息队列多节点生产环境需要自己维护 broker 和 workerKubernetesExecutor每个任务动态起一个 Pod已有 K8s 平台的团队任务粒度小、Pod 启动有开销选型建议很直白学习和验证期用 LocalExecutor业务上量后迁移 CeleryExecutor 或 KubernetesExecutor。Executor 在airflow.cfg或环境变量AIRFLOW__CORE__EXECUTOR中配置切换后任务编排代码一行都不用改——这正是DAG as Code的好处。⚠️ 注意并发由多个参数共同钳制parallelism全局活跃任务数上限、dag_concurrency单 DAG 并发上限、worker 自身的并发度三者取最严格的一个生效。调并发之前先想清楚这一层不然为什么任务没跑起来会很难查。7. 用 Pool 控制并发不压垮下游接口假设有 12 个任务要写同一个数仓而数仓的并发连接上限只有 3——直接放开跑第 4 个开始就会超时甚至把服务打挂。解法是Pool给一组干同一件事的任务设并发上限。创建 Pool 有三种方式UI 的Admin → Pools手动加、airflow poolsCLI 命令支持从 JSON 批量导入、或 REST API。然后在任务上指定即可load_summary PythonOperator( task_idload_summary, python_callablewrite_to_warehouse, poolwarehouse_write_pool, # 3 个槽位 priority_weight2, # 槽位抢不满时谁先跑 )配套的两个参数值得了解pool_slots一个任务占用的槽位数默认 1。适合一个任务吃掉整个 GPU 节点这类独占场景priority_weight同一 Pool 内任务排队时的优先级数值大的先执行。⚠️ 坑Pool 名写错不会报错。任务会一直挂在待调度状态UI 也不做任何检查。把 Pool 名抽成常量集中管理是团队里踩坑之后的通用做法。Pool 只控制任务实例级的并发。如果你要限制的是同一个 DAG 最多几个 Run 并行那是max_active_runs的活儿别拿 Pool 硬套。8. 进阶玩法动态任务与基于数据的调度动态任务Airflow 2.3任务数量不再写死在代码里。比如每个大区导出一份文件大区列表是运行时才知道的task def export_region_file(region: str): ... region_list fetch_regions() # 运行期拿到 [cn-east, cn-north, us-west] export_region_file.expand(regionregion_list)partial()传所有任务共享的参数expand()传要展开成多个并行实例的那个参数。UI 里这类任务会显示task_id [3]这样的角标点进去能逐个查看映射出的实例——实例数量随数据变化DAG 代码不变这是 2.3 之后写可变规模管道的标准姿势。Dataset 触发Airflow 2.4当什么时候跑不由时间决定、而由数据到了没有决定时Cron 就力不从心了。生产者任务通过outlets声明我产出了这份数据消费者 DAG 把它写进scheduleorders_dataset Dataset(postgres://orders_db/fct_daily_order) with DAG(order_report_dag, schedule[orders_dataset], catchupFalse): build_report PythonOperator(task_idbuild_report, python_callablerender_report)两种调度方式怎么选维度Cron / 时间调度Dataset 数据调度触发依据固定时间点上游数据就绪典型场景每日汇总、对账跨团队的我依赖你的产出风险上游延迟时下游空跑生产者不声明 outlets 则永远不触发 最佳实践Dataset 的 URI 只会被原样存储不要把密码、密钥拼进 URI 字符串用环境变量或密钥后端托管。9. 落地学习路径7 步走到生产按下面的顺序动手每一步都有明确产出大约一到两个周末就能走完装好并跑通最小 DAG初始化数据库、起 Scheduler Webserver用本文第 2 节的三任务 DAG 在 Grid View 里看到全绿手动制造一次失败让clean_dedupe抛异常验证重试和日志是否按预期工作加上失败回调配一个最简单的on_failure_callback发消息到你的手机或群插入质量检查任务先放一条最简单的行数下限检查位置选在入库前建一个 Pool给所有写数仓的任务统一挂上确认排队行为符合预期把一个固定列表改成expand体会任务数由数据决定研究 Dataset 调度给生产者声明outlets把下游 DAG 从daily切到数据触发。配套资源本项目的guides/目录是一整套英文指南合集其中airflow-pools.mdPool 全貌、airflow-executors-explained.md四类 Executor 深度对比、dynamic-tasks.md动态任务映射、scheduling-in-airflow.md时间调度与 Timetable、data-quality.md质量工具选型几篇与本文主线一一对应可对照细读。最后一步把仓库拉到本地跟着练git clone https://gitcode.com/gh_mirrors/ai/airflow-guides写到这里回看开头那个凌晨两点的告警如果它当时跑在 Airflow 上retries可能已经自愈了没有自愈的话on_failure_callback会在第一次失败时就把你叫醒——而且日志、状态、重跑入口都现成。这就是数据工作流编排工具最朴素的价值把出事了人肉救火变成系统自愈 人只看该看的。【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guides创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表