ARTICLE DETAIL

资讯详情

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

从 0 到 1 搭建 Airflow 数据管道:5 步搞定 DAG、并发与调度

从 0 到 1 搭建 Airflow 数据管道:5 步搞定 DAG、并发与调度 从 0 到 1 搭建 Airflow 数据管道5 步搞定 DAG、并发与调度【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guidesApache Airflow 是一款用 Python 代码定义数据工作流编排的开源引擎把数据管道拆成一个个任务Task按依赖关系自动调度、失败重试并记录每一步状态相当于给你的数据任务装上自动巡航 行车记录仪。本文不堆概念直接带你走完五个环节认清组件 → 跑通最小 DAG → 控制并发 → 选对调度 → 接上告警每一步都能照着做。一、先看清楚Airflow 背后是哪几个进程在干活 写 DAG 之前建议先花三分钟知道 Airflow 由哪些部件组成这样出问题时你知道该往哪查Webserver提供 Airflow 网页界面UI你平时看任务状态、暂停 DAG 都在这里。Scheduler调度器多进程 Python 守护进程负责决定哪个任务、什么时候、在哪里跑。Database元数据库存放 DAG 和任务状态记录通常用 Postgres也支持 MySQL、SQLite 等。Executor执行器真正派发任务干活的机制它决定了你的部署是单机还是分布式。另外还有两个按需启用的组件Worker按执行器决定是否独立存在和Triggerer只有使用可延迟/异步类 Operator 时才需要单独启动。想深入细节可以看 guides/airflow-components.md。二、最小可用示例10 分钟写出第一个能跑的 DAG 先别急着学高级功能一个能跑通的 DAG 只需要三样东西default_args公共参数、任务、依赖箭头读作上游完成后再执行下游from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator def fetch_orders(**context): print(拉取当日订单明细...) def load_orders(**context): print(写入数仓...) default_args { owner: data-team, start_date: datetime(2026, 1, 1), retries: 2, retry_delay: timedelta(minutes5), } with DAG( dag_idorder_pipeline, scheduledaily, catchupFalse, default_argsdefault_args, ) as dag: fetch PythonOperator(task_idfetch_orders, python_callablefetch_orders) load PythonOperator(task_idload_orders, python_callableload_orders) fetch load几个容易踩的坑官方最佳实践指南里有详细说明guides/dag-best-practices.mdstart_date要用固定日期别用datetime.now()否则补数据和重跑时容易出错。DAG 文件里别写顶层代码如直接发 API 请求因为调度器默认每 30 秒解析一遍dags目录顶层请求会被反复执行。依赖写法只选一种或set_downstream全文件保持一致可读性会好很多。三、并发失控怎么办用 Pool 给任务排队 ⚡任务一多你大概率会遇到这种情况20 个任务同时打同一个接口把上游 API 打挂了。Airflow 的Pool资源池就是干这个的——给一组任务设一个并发上限。所有任务默认都在default_pool里带 128 个槽位。创建自定义池有三条路UI 里走Admin→Pools填名称、槽位数、描述用airflow pools set等 CLI 命令还支持从 JSON 批量导入Airflow 2.0 可用 REST API 提交 POST 请求创建。然后在任务上指定pool参数槽位满了任务就排队有空位自动放行用priority_weight可以决定排队时谁先跑值越大越优先fetch PythonOperator( task_idfetch_orders, python_callablefetch_orders, poolapi_rate_limited_pool, # 只允许同时跑 3 个 priority_weight3, # 比别的 DAG 里的任务优先 )⚠️ 一个反直觉的坑如果pool写了个不存在的池名任务会静默地永远不执行UI 上也没有报错提示。建池子后务必核对拼写。更多实战案例比如跨两个 DAG 共享同一个池、按 DAG 分级限流在 guides/airflow-pools.md 里值得对照着改自己的代码。四、定时跑还是数据驱动三种调度方式怎么选 schedule参数决定了 DAG 什么时候运行按复杂度递进有三档方式写法示例适合场景cron 表达式/预设scheduledaily、schedule30 2 * * *固定时刻、固定周期timedelta 对象scheduletimedelta(hours6)按固定频率每 6 小时滚动Dataset 触发schedule[orders_ready]上游数据到了再跑不靠猜时间第三档是 Airflow 2.4 引入的数据驱动调度只要某个 Dataset 被更新依赖它的 DAG 就会被触发上游在哪个 DAG 里产出的数据都不影响from airflow.datasets import Dataset orders_ready Dataset(s3://lakehouse/orders/orders.csv) with DAG( dag_idorder_pipeline, schedule[orders_ready], catchupFalse, default_argsdefault_args, ) as dag: ...这类 DAG 在 UI 的 Schedule 列会显示为DatasetNext Run 列还会告诉你依赖的数据集更新了几个。另外提醒一下cron 调度下周一的逻辑日期运行实际要等到周二才执行数据区间是结束后再跑如果时间要求特殊Airflow 2.2 还支持自定义 Timetable 来突破 cron 限制细节见 guides/scheduling-in-airflow.md。五、失败不慌重试、告警与生产选型 ️分布式环境里任务失败是常态关键是把怎么活下来写进配置而不是靠人盯重试把retries建议 2~4 次和retry_delay放进default_args全 DAG 生效个别任务再单独覆盖。告警配置email_on_failure、on_failure_callback等参数把失败信息推到邮箱或 IM 群SLA 超期也能触发提醒完整做法见 guides/error-notifications-in-airflow.md。幂等保证同一个任务重跑多次结果一致比如写入用覆盖分区而非追加这样失败后重跑才是安全的。任务原子化抽取、转换、加载各写一个任务失败时只重跑坏的那一段。生产环境的执行器Executor按规模选执行器特点一句话建议SequentialExecutor单进程顺序执行无并行本地调试LocalExecutor单机多进程并行开发和单机部署CeleryExecutor消息队列 多 Worker 分布式任务量大的生产主力KubernetesExecutor每个任务独立 Pod长任务、资源隔离要求高的场景选型对比与更多细节见 guides/airflow-executors-explained.md。下一步跟着做这三件事把上面第二个代码块里的 DAG 改造成你自己的真实管道放到本地 Airflow 的dags目录里跑一遍打开 UI 确认任务状态流转正常给你打外部 API 的任务建一个 Pool故意把槽位数调到 1观察任务排队和priority_weight的生效顺序动手排查一次故障故意让任务失败验证重试与失败回调是否按预期触发然后翻 guides/debugging-dags.md 对照日志找原因。想系统学习可以克隆本仓库离线阅读全部指南git clone https://gitcode.com/gh_mirrors/ai/airflow-guides建议阅读顺序guides/get-started-airflow-2.md本地环境→ guides/dags.mdDAG 详解→ guides/templating.md模板与宏→ guides/data-quality.md数据质量校验从能跑逐步走到跑得稳。【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guides创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表