尧图建网站 尧图建网站 YAOTU WEB BUILD 免费咨询
ARTICLE DETAIL

资讯详情

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

Apache Airflow 零基础实战:15 行代码搭起你的第一条数据管道

Apache Airflow 零基础实战:15 行代码搭起你的第一条数据管道 Apache Airflow 零基础实战15 行代码搭起你的第一条数据管道【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guidesApache Airflow 是一款开源的数据工作流编排工具用 Python 代码定义数据管道并自动调度。读完本文你能独立写出第一条 DAG并学会配重试、限流、动态任务和失败告警。Airflow 核心概念速览第一次接触 Airflow只需要记住 4 个词DAG就是一张「先做 A 再做 B」的流程图Airflow 按这张图自动调度任务图不能有环。Operator流程图上每个节点对应的执行单元负责完成一个具体动作比如跑一段 SQL、调一次 API。Hook和外部系统数据库、API的连接桥给 Operator 提供现成的连接方式你不用自己写连接代码。Pool一套「令牌」机制池子里有几个槽位就允许几个任务并发执行超出的任务排队等待。在 UI 里你会看到 Graph View 和 Grid View前者看依赖结构后者按时间轴看每次任务的执行状态排查问题时两者配合用。快速上手15 行写出第一条 DAG下面这段代码定义了一个每天运行、带 2 次重试的两步管道extract取数据transform加工数据。把它放进dags目录刷新 UI 就能看到它。from datetime import datetime from airflow.decorators import dag, task dag(scheduledaily, start_datedatetime(2024, 1, 1), catchupFalse, default_args{retries: 2}) def my_first_pipeline(): task def extract(): return [1, 2, 3] task def transform(data): return [x * 2 for x in data] transform(extract()) my_first_pipeline()task装饰器就是 TaskFlow API它会自动根据函数调用关系推断依赖不用手动写。运行后到 Grid View 里确认两个任务都变绿说明你的第一条数据管道已经跑通。典型场景实战场景一如何给任务配重试和失败告警遇到什么问题上游数据库偶发超时任务一失败就得人工盯着 UI 等通知。怎么配在default_args里统一声明重试策略和邮件告警整个 DAG 的所有任务自动继承。from datetime import timedelta default_args { retries: 3, retry_delay: timedelta(minutes5), email_on_failure: True, email: [opsyourteam.com], }效果瞬时故障每 5 分钟自动重试 3 次重试耗尽仍失败时邮件发到值班组半夜不用再刷 UI。场景二Pool 限流的正确打开方式遇到什么问题几十个任务同时打同一个第三方 API触发对端限流全线报错。怎么配在 UI 的Admin→Pools建一个 5 槽位的api_pool然后把调用该 API 的任务都指过去再用priority_weight给关键任务加优先级。task PythonOperator( task_idcall_api, python_callablecall_api, poolapi_pool, priority_weight2, )⚠️注意如果pool指向一个不存在的池任务永远不会被调度且 UI 不会报错建池子前先确认拼写。效果最多 5 个请求并发打 API其余任务排队限流消失关键任务先跑。场景三动态任务生成处理数量不定的文件遇到什么问题每天落盘的文件从 3 个到 300 个不等写死任务不现实。怎么配用partial()固定不变的参数用expand()对列表做映射运行时自动生成 N 个并行任务。task def load_file(name: str): print(floading {name}) return name files previous_task() # 返回今天需要处理的文件名列表 load_file.partial().expand(namefiles)效果UI 中任务显示为load_file [ ]方括号里的数字就是本次展开的实例数点开可逐个查看日志。Airflow 常见坑与修复坑 1DAG 文件每 30 秒执行一次数据库查询症状数据库连接数被 Airflow 打满。原因调度器每隔min_file_process_interval默认 30 秒会重新解析一遍dags目录你在文件顶层写的外部请求会跟着一起执行。修复把所有取数逻辑挪进 Operator 或任务函数内部文件顶层只保留 DAG 定义。坑 2重跑历史日期的任务时数据跑错日期症状手动重跑上周的 DAG任务却处理了今天的分区。原因代码里用了datetime.today()它取的是当前时间而不是本次运行的逻辑日期。修复改用模板变量{{ ds }}或{{ ds_nodash }}它们始终指向本次 DAG Run 的日期。坑 3新一次 DAG Run 迟迟不启动症状上一个 Run 还在跑新的 Run 一直处于排队。原因max_active_runs默认值为 1同一时刻只允许一个 Run。修复在 DAG 定义中设置max_active_runs3按需放开并发。坑 4UI 里找不到刚提交的 DAG症状代码没问题DAG 列表就是看不到。原因文件不在调度器监控的dags_folder里或文件名以test_、__开头被忽略。修复确认文件放在dags目录下且命名符合规范必要时调小min_file_process_interval加快解析。进阶调优执行器与调度策略跑通之后往下研究这三个方向选对 ExecutorSequentialExecutor单进程顺序执行适合本地调试LocalExecutor单机多进程并行适合开发机生产环境用CeleryExecutor或KubernetesExecutor横向扩 Worker任务执行能力不再受单机限制。选对调度方式固定时间用 cron 表达式schedule5 4 * * *Airflow 2.2 还支持自定义 Timetable 和基于 Dataset 的数据驱动调度——上游数据落地才触发下游比掐点跑更可靠。做好资源隔离dag_concurrency限制单个 DAG 的并发任务数worker_concurrency限制每个 Worker 的承载两者配合 Pool 从三个粒度管住并发避免大 DAG 把小 DAG 的资源挤没。仓库guides/目录下有 50 篇主题指南从 Executor 对比到 Databricks、Snowflake 集成都有专门文档建议按需取阅。想动手时先跑通本文的 15 行示例再按场景逐个叠加重试、限流和告警你的第一条生产级数据管道就成了。【免费下载链接】airflow-guidesGuides and docs to help you get up and running with Apache Airflow.项目地址: https://gitcode.com/gh_mirrors/ai/airflow-guides创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表