
1. 从脚本到流程为什么数据工程师需要一个“调度器”如果你写过数据处理的脚本无论是用Python的pandas清洗CSV还是用requests拉取API数据最终都会遇到一个现实问题脚本写好了谁来定时跑今天凌晨1点要跑一个ETL任务明天早上9点要生成日报后天还要在A任务跑完后自动触发B任务。最开始你可能会用操作系统的crontab或者Windows的任务计划程序。这确实能解决“定时”问题但很快你就会发现当任务多了、依赖复杂了、失败需要重试或者要查看历史运行记录时crontab就显得力不从心了。它就像一个只会看钟表干活的工人钟点一到就执行命令至于任务成功与否、上下游是谁、中间出了什么状况它一概不知也不关心。这就是工作流调度器Workflow Scheduler要解决的问题。它不是一个简单的定时触发器而是一个编排、调度、监控和运维复杂任务依赖关系的平台。想象一下你要构建一个数据仓库的每日更新流程先从多个数据库抽取数据Extract然后进行清洗转换Transform最后加载到目标表Load。这个ETL流程中的每一步都可能是一个独立的脚本或程序它们之间有严格的先后顺序和依赖关系。转换任务必须等所有抽取任务都成功完成才能开始加载任务又必须等转换任务成功。如果某个抽取任务失败了整个流程应该暂停并告警而不是让后续任务拿着错误的数据继续运行。Apache Airflow正是为了解决这类问题而生的。它用Python代码来定义工作流将一个个任务Task及其依赖关系Dependencies描述成一个有向无环图DAG。你写的不是配置而是代码。这意味着你可以用Python的所有能力变量、循环、条件判断、从外部API获取参数等来动态生成你的工作流。这对于数据工程师来说无疑是如虎添翼。它把调度这个“脏活累活”系统化、可视化、可运维化让你能从繁琐的cron管理和日志排查中解放出来更专注于数据逻辑本身。因此说Airflow是数据工程师的“标配工具”并非夸大其词它几乎成了现代数据栈中承上启下的“中枢神经系统”。2. Airflow核心架构拆解元数据库、调度器与执行器如何协同要玩转Airflow不能只停留在写DAG的层面理解其核心组件如何协同工作是解决实际部署和运维中各种“怪现象”的关键。Airflow的架构清晰地区分了“定义”、“调度”和“执行”这三个关注点。2.1 元数据库Metastore这是Airflow的大脑皮层存储了所有的状态信息。默认使用SQLite仅用于测试生产环境通常用PostgreSQL或MySQL。它里面存了些什么DAG定义你写的Python DAG文件解析后的元数据。任务实例Task Instances每一次DAG运行DAG Run中每一个任务的具体实例及其状态如success、failed、running、up_for_retry。变量Variables和连接Connections全局的键值对配置和外部系统如数据库、API的连接信息敏感信息通常由环境变量或外部Secret管理工具注入。执行历史与日志路径所有任务运行的记录和日志文件索引。这个数据库是调度器、执行器和Web Server之间通信的枢纽。如果它性能不佳整个Airflow都会卡顿。一个常见的优化就是为task_instance等核心表建立合适的索引。2.2 调度器Scheduler这是Airflow的心脏一个持续运行的守护进程。它的工作流程是个循环解析DAGs定期默认约30秒扫描DAG_FOLDER目录下的Python文件解析出其中的DAG对象并将其序列化后存入元数据库。这里有个坑如果你的DAG文件顶部有耗时的导入比如导入一个巨大的机器学习库会严重拖慢调度器解析速度。最佳实践是将业务逻辑封装在任务函数内DAG文件顶部只保留必要的轻量级导入。检查任务状态根据DAG中设定的调度时间schedule_interval和依赖关系检查哪些任务满足了执行条件上游任务成功、时间已到等。创建任务实例将满足条件的任务实例化状态标记为queued排队中并将其放入执行队列。调度器本身不执行任务它只负责派活。调度器是单进程的虽然从Airflow 2.0开始支持高可用HA部署即运行多个调度器实例但它们需要协调以避免重复调度。这通常通过数据库行级锁来实现。如果调度器挂了新的任务不会被调度但已提交到执行器的任务会继续运行。2.3 执行器Executor这是真正干活的肌肉。调度器决定“什么时间做什么事”执行器负责“找人把事做了”。Airflow支持多种执行器这是其灵活性的体现LocalExecutor在调度器同一台机器上使用多进程或线程池来执行任务。适合中小规模、任务类型单一主要是Python任务的部署。缺点是任务会竞争调度器机器的资源。CeleryExecutor最经典的生产级选择。它使用Celery作为分布式任务队列。调度器将任务推送到消息队列如Redis/RabbitMQ一群独立的Worker进程可以在不同机器上从队列中拉取任务并执行。这实现了水平扩展Worker可以按需增减。KubernetesExecutor云原生时代的首选。每个任务实例Task Instance都会被调度器提交为Kubernetes集群中的一个独立的Pod。任务完成后Pod被销毁。这提供了极致的资源隔离和弹性但需要一定的K8s运维知识。你的任务需要被打包进Docker镜像。2.4 Web服务器Web Server这是一个Flask应用提供了图形化界面UI。你可以在这里查看DAG运行状态、触发手动运行、查看任务日志、管理变量和连接等。它是运维人员的主要交互界面。Web Server通常是无状态的可以水平扩展通过负载均衡器对外提供服务。2.5 工作线程Worker当使用CeleryExecutor或KubernetesExecutor时Worker是实际执行任务代码的进程或Pod。它们需要能够访问到DAG文件通常通过共享存储如NFS、Git同步或分布式文件系统以及任务代码所依赖的Python环境。整个工作流程可以概括为你写Python DAG文件 - 调度器解析并监控 - 满足条件后调度器将任务实例放入队列 - 执行器分配队列中的任务给Worker - Worker执行具体任务代码 - 结果和日志写回元数据库 - 你在Web UI上查看一切。理解这个流程当任务卡在“排队”状态时你就知道该去检查执行器Celery和Worker当DAG不显示时就知道该去检查调度器的解析日志。3. 编写你的第一个DAG超越“Hello World”的实用入门看过太多打印“Hello Airflow”的教程但那些离真实场景太远。让我们直接从一个贴近实际的数据工程师日常任务开始每天凌晨2点从某个API获取JSON格式的天气数据解析后存入PostgreSQL数据库如果失败则重试3次每次间隔5分钟。3.1 环境准备与项目结构首先确保你有一个可用的Airflow环境。对于本地开发最快捷的方式是使用官方提供的docker-compose文件。这里假设你已安装Docker并初始化了Airflow。一个清晰的DAG项目目录结构很重要your_project/ ├── dags/ # 存放所有DAG文件 │ └── fetch_weather_dag.py ├── plugins/ # 可选存放自定义Operator、Hook等 ├── config/ # 可选存放配置文件或SQL模板 └── requirements.txt # 项目Python依赖将你的fetch_weather_dag.py放在dags/目录下Airflow调度器会自动发现它。3.2 DAG定义与参数解析打开fetch_weather_dag.py我们开始编写from datetime import datetime, timedelta from airflow import DAG from airflow.operators.python import PythonOperator from airflow.providers.postgres.operators.postgres import PostgresOperator from airflow.providers.postgres.hooks.postgres import PostgresHook import requests import json # 1. 定义默认参数字典 default_args { owner: data_team, # 任务负责人 depends_on_past: False, # 是否依赖上一次DAG运行的成功 email: [alertyourcompany.com], # 失败时通知的邮箱列表 email_on_failure: True, # 失败时发邮件 email_on_retry: False, # 重试时不发邮件 retries: 3, # 失败后重试次数 retry_delay: timedelta(minutes5), # 重试间隔 start_date: datetime(2023, 10, 27), # **重要**调度开始的锚点日期 } # 2. 实例化DAG对象 dag DAG( daily_weather_etl, # DAG的唯一ID在UI中显示 default_argsdefault_args, # 注入默认参数 descriptionFetch daily weather data and load to PostgreSQL, schedule_interval0 2 * * *, # 每天UTC时间2点运行 (cron表达式) catchupFalse, # 是否补跑从start_date到现在错过的任务 tags[weather, etl], # 用于UI分类的标签 )这里有几个关键点start_date这是Airflow调度逻辑的基石。调度器会计算start_date schedule_interval的时间序列。假设start_date是2023-10-27schedule_interval是daily那么第一个DAG Run的执行日期execution_date是2023-10-27但实际运行时间是在2013-10-28 02:00因为要等27号的数据在28号凌晨2点可用。这个概念初学极易混淆。catchup设为False可以避免在部署DAG或Airflow服务停机后一次性触发大量历史任务导致“追赶”行为。生产环境通常建议关闭除非你有明确的补数据需求。schedule_interval可以用cron表达式如0 2 * * *也可以用预置的宏如daily、hourly。3.3 编写任务函数与使用Hook接下来我们定义具体的任务函数。任务间传递数据可以使用Airflow的XCom跨任务通信机制但对于大量数据更推荐使用外部存储如S3、数据库。这里我们演示一个简单的流程。def fetch_weather_data(**context): 从公开天气API获取数据。 使用context可以获取到execution_date等运行时信息。 # 假设我们使用一个公开的天气API api_url https://api.open-meteo.com/v1/forecast # 使用execution_date来决定获取哪天的数据通常是前一天 # execution_date是逻辑上的数据日期 execution_date context[execution_date] target_date execution_date - timedelta(days1) # 获取前一天的天气 params { latitude: 39.9042, # 例如北京 longitude: 116.4074, start_date: target_date.strftime(%Y-%m-%d), end_date: target_date.strftime(%Y-%m-%d), hourly: temperature_2m } try: response requests.get(api_url, paramsparams, timeout30) response.raise_for_status() # 如果状态码不是200抛出HTTPError weather_data response.json() # 这里可以进行一些简单的数据解析 extracted_data { date: target_date.strftime(%Y-%m-%d), location: Beijing, avg_temp: sum(weather_data[hourly][temperature_2m]) / len(weather_data[hourly][temperature_2m]), raw_json: json.dumps(weather_data) # 存储原始JSON } # 将数据通过XCom传递给下一个任务适用于小数据 context[task_instance].xcom_push(keyweather_extract, valueextracted_data) print(fSuccessfully fetched weather data for {target_date}) except requests.exceptions.RequestException as e: # 异常会被Airflow捕获触发重试或标记失败 print(fFailed to fetch weather data: {e}) raise def load_weather_to_db(**context): 将数据加载到PostgreSQL。 使用PostgresHook来管理连接避免在代码中硬编码密码。 # 从XCom中拉取上一个任务推送的数据 ti context[task_instance] extracted_data ti.xcom_pull(task_idsfetch_weather, keyweather_extract) if not extracted_data: raise ValueError(No data received from fetch_weather task.) # 使用PostgresHook连接信息在Airflow UI的Connections中配置 postgres_hook PostgresHook(postgres_conn_idweather_postgres) conn postgres_hook.get_conn() cursor conn.cursor() insert_sql INSERT INTO public.weather_daily (data_date, location, average_temperature, raw_data) VALUES (%s, %s, %s, %s) ON CONFLICT (data_date, location) DO UPDATE SET average_temperature EXCLUDED.average_temperature, raw_data EXCLUDED.raw_data, updated_at CURRENT_TIMESTAMP; cursor.execute(insert_sql, ( extracted_data[date], extracted_data[location], extracted_data[avg_temp], extracted_data[raw_json] )) conn.commit() cursor.close() conn.close() print(fWeather data for {extracted_data[date]} loaded successfully.)3.4 组装任务并设置依赖最后用Operator包装这些函数并定义它们之间的执行顺序。# 3. 定义任务 # 任务1创建目标表如果不存在。这是一个一次性或幂等的SQL任务。 create_table_task PostgresOperator( task_idcreate_weather_table, postgres_conn_idweather_postgres, sql CREATE TABLE IF NOT EXISTS public.weather_daily ( id SERIAL PRIMARY KEY, data_date DATE NOT NULL, location VARCHAR(100) NOT NULL, average_temperature DECIMAL(5,2), raw_data JSONB, created_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, updated_at TIMESTAMP DEFAULT CURRENT_TIMESTAMP, UNIQUE(data_date, location) ); , dagdag, ) # 任务2提取天气数据 fetch_task PythonOperator( task_idfetch_weather, python_callablefetch_weather_data, provide_contextTrue, # 传递context参数 dagdag, ) # 任务3加载数据到数据库 load_task PythonOperator( task_idload_weather_to_db, python_callableload_weather_to_db, provide_contextTrue, dagdag, ) # 4. 设置任务依赖关系 # 表示create_table_task 先执行成功后执行 fetch_task再成功后执行 load_task create_table_task fetch_task load_task现在将这个DAG文件放入你的dags/文件夹在Airflow Web UI上刷新你应该能看到一个名为daily_weather_etl的新DAG。将其激活Unpause调度器会在下一个调度周期明天凌晨2点触发它。你也可以手动触发一次来测试。注意在实际生产环境中API密钥、数据库密码等敏感信息绝对不要硬编码在DAG文件中。应该使用Airflow的Variable对于非敏感配置或Connections对于连接信息密码部分会被加密存储来管理或者更好的是使用像HashiCorp Vault这样的外部密钥管理服务通过环境变量注入。4. 生产环境部署与调优实战指南让一个DAG在本地跑起来是一回事让几十上百个DAG在生产环境中稳定、高效、可维护地运行是另一回事。以下是基于真实踩坑经验总结的部署与调优要点。4.1 执行器选型Celery还是K8sCeleryExecutor优点成熟、稳定、社区资料丰富。利用消息队列Redis/RabbitMQ解耦调度器和WorkerWorker可以水平扩展。适合大多数传统数据中心或虚拟机环境。缺点需要维护消息队列和Worker集群。任务运行在常驻的Worker进程/容器中存在环境依赖冲突的风险例如两个DAG需要不同版本的pandas。资源隔离性相对较弱。部署建议使用docker-compose或helm在K8s上部署一套独立的Airflow包含Scheduler, Web, Worker, Redis。为Worker设置资源限制CPU/Memory。使用CeleryKubernetesExecutor混合模式将资源密集型或特殊环境需求的任务路由到K8s Pod执行普通任务由Celery Worker执行。KubernetesExecutor优点资源隔离与弹性伸缩的终极解决方案。每个任务都在独立的Pod中运行环境通过Docker镜像绝对隔离。Pod规格可以按任务定制用完即焚不残留任何状态。天然契合云原生环境。缺点复杂度高需要K8s运维能力。每个任务启动Pod会有一定的开销镜像拉取、容器启动对于超短时任务秒级可能不划算。需要为所有可能的任务环境提前构建好镜像或者使用工具动态构建。部署建议使用官方Helm Chart部署。精心设计基础镜像包含常用依赖。利用K8s的ResourceQuota和LimitRange控制资源。对于需要访问特定工具如kubectl,awscli的任务考虑使用Sidecar容器或定制镜像。4.2 关键配置参数调优airflow.cfg文件中有数百个配置项以下几个对性能影响巨大parallelism控制整个Airflow实例允许同时运行的任务实例总数。根据数据库和整体资源设置。dag_concurrency控制单个DAG内可以同时运行的任务实例数。通常小于parallelism。max_active_runs_per_dag控制单个DAG同时可以有多少个DAG Run处于活动状态。防止一个出错的DAG无限重试产生海量任务实例。scheduler_heartbeat_sec调度器心跳间隔。在集群部署中用于判断调度器是否存活。worker_precheckCelery Worker在执行任务前是否检查数据库连接。生产环境建议设为False以避免不必要的开销。executor根据你的选择设为CeleryExecutor或KubernetesExecutor等。4.3 监控、告警与日志没有监控的系统就是在裸奔。监控指标使用StatsD导出Airflow指标任务成功/失败数、调度延迟、执行时间等到Prometheus用Grafana展示。关键看板包括DAG运行时长趋势、任务失败率、调度器健康状态、队列深度。告警基于上述指标设置告警规则。例如某个关键DAG连续失败、调度器心跳丢失、任务排队时间超过阈值。日志管理默认日志存储在本地文件难以查询。生产环境必须配置远程日志存储如S3、GCS、Elasticsearch。在airflow.cfg中配置remote_logging相关选项。这样在Web UI上可以直接查看集中存储的日志便于排查问题。4.4 高可用HA部署对于关键业务需要避免单点故障。元数据库使用云托管的、支持高可用的PostgreSQL或MySQL服务如AWS RDS Multi-AZ。调度器Airflow 2.0支持多调度器。运行2-3个调度器实例它们会通过数据库锁进行协调。使用进程管理器如systemd, supervisord或K8s Deployment确保调度器进程挂掉后能自动重启。Web服务器无状态可以通过负载均衡器如Nginx, ALB后面部署多个实例。WorkerCelery部署多个Worker节点并确保它们能自动加入Celery集群。使用K8s Deployment或类似工具管理。消息队列Celery使用Redis Cluster或RabbitMQ集群来保证消息队列的高可用。4.5 DAG版本控制与CI/CDDAG即代码必须纳入版本控制如Git。并建立CI/CD流水线代码检查在CI中运行pylint、black、mypy等工具进行代码风格和静态类型检查。DAG完整性测试编写简单的Python脚本导入DAG文件检查是否有语法错误验证DAG ID是否唯一任务ID是否重复等。部署通过CD流程将DAG文件同步到所有Airflow Worker节点共享的存储如NFS、S3或通过Git-Sync sidecar容器同步到K8s Pod内。回滚确保能快速回滚到上一个可用的DAG版本。5. 高级模式与最佳实践让Airflow发挥最大威力当你熟练掌握了基础DAG编写后以下高级模式和最佳实践能帮助你构建更健壮、更灵活的数据管道。5.1 动态DAG生成有时你需要创建大量结构相似、仅参数不同的DAG。例如为公司的每个产品线运行相同的ETL流程。手动复制粘贴DAG文件是灾难。这时可以使用动态DAG生成。def create_dag_for_product(product_id, schedule): 为每个产品动态创建一个DAG的工厂函数。 default_args {...} # 通用默认参数 dag_id fetl_product_{product_id} with DAG(dag_iddag_id, schedule_intervalschedule, default_argsdefault_args) as dag: start DummyOperator(task_idstart) extract PythonOperator( task_idfextract_{product_id}, python_callableextract_product_data, op_kwargs{product_id: product_id} # 将参数传递给任务函数 ) transform PythonOperator(...) load PythonOperator(...) start extract transform load return dag # 从配置文件或数据库读取产品列表 product_list get_all_products() for product in product_list: globals()[fetl_product_{product.id}] create_dag_for_product(product.id, product.schedule)Airflow调度器在导入模块时会执行顶层的Python代码因此这些globals()赋值操作会动态创建出多个DAG对象。务必注意动态生成的DAG ID必须稳定不能每次导入都变化否则会导致UI中出现大量“僵尸”DAG。5.2 使用TaskGroup和SubDAGs组织复杂流程对于包含几十上百个任务的巨型DAGUI上会变成一团乱麻。Airflow 2.0引入了TaskGroup可以将一组相关的任务在UI上折叠显示极大提升了可读性。from airflow.utils.task_group import TaskGroup with DAG(...) as dag: with TaskGroup(group_iddata_validation) as validation_group: task_a PythonOperator(task_idvalidate_schema, ...) task_b PythonOperator(task_idcheck_null, ...) task_c PythonOperator(task_idanomaly_detection, ...) task_a [task_b, task_c] # TaskGroup内部依赖 with TaskGroup(group_idreporting) as reporting_group: ... validation_group reporting_group # TaskGroup之间的依赖至于SubDAGs子DAG官方已不推荐使用因为它存在执行器隔离、死锁等问题TaskGroup是更好的替代方案。5.3 利用XCom进行任务间小数据通信XCom允许任务间传递小的数据片段键值对。如上文示例我们用它传递了提取的数据。但需牢记XCom数据存储在元数据库中不适合传递大型数据如DataFrame这会导致数据库迅速膨胀。传递大数据时应使用外部存储如S3路径、数据库记录ID然后通过XCom传递这个引用。5.4 传感器Sensors与智能调度传感器是一种特殊类型的Operator它会持续“感知”某个外部条件是否满足如文件是否到达S3、数据库分区是否就绪、API端点是否可用条件满足后才触发下游任务。这实现了基于事件的调度而不仅仅是基于时间。from airflow.providers.amazon.aws.sensors.s3 import S3KeySensor from airflow.sensors.external_task import ExternalTaskSensor # 等待S3上出现特定文件 wait_for_file S3KeySensor( task_idwait_for_input_file, bucket_keys3://my-bucket/input/{{ ds }}/data.csv, # 使用Jinja模板 aws_conn_idaws_default, timeout3600, # 等待超时时间秒 poke_interval60, # 每次检查间隔秒 modepoke, # 模式poke默认或reschedule dagdag, ) # 等待另一个DAG在特定日期的运行成功 wait_for_upstream_dag ExternalTaskSensor( task_idwait_for_daily_aggregation, external_dag_iddaily_aggregation_dag, external_task_idNone, # None表示等待整个DAG运行成功 allowed_states[success], execution_date_fnlambda dt: dt, # 通常等待同一天 timeout7200, dagdag, )使用modereschedule的传感器在检查间隔期内会释放Worker插槽更节省资源。5.5 错误处理与重试策略除了在default_args中定义全局重试策略还可以在任务级别覆盖。更精细的控制可以通过设置回调函数实现def alert_on_failure(context): 任务失败时的回调函数可以发送更详细的告警如Slack、PagerDuty。 task_instance context[task_instance] dag_id task_instance.dag_id task_id task_instance.task_id exception context.get(exception) error_message fTask {task_id} in DAG {dag_id} failed. Exception: {exception} # 调用发送告警的函数如发送到Slack webhook send_slack_alert(error_message) def retry_delay_calculator(retry_number): 自定义指数退避的重试延迟。 return timedelta(seconds10 * (2 ** retry_number)) # 10, 20, 40秒... my_task PythonOperator( task_idmy_task, python_callable..., retries5, retry_delaytimedelta(minutes1), on_failure_callbackalert_on_failure, # 失败回调 retry_exponential_backoffTrue, # 启用指数退避 max_retry_delaytimedelta(minutes30), # 最大重试间隔 dagdag, )5.6 测试你的DAG测试是保证数据管道可靠性的基石。单元测试单独测试你的任务函数python_callable像测试普通Python函数一样使用pytest。模拟mock外部依赖如requests.get,PostgresHook。集成测试在独立的测试环境如Docker化的Airflow中运行完整的DAG。可以使用airflow dags test dag_id execution_date命令在本地运行一次DAG不依赖调度器但这不会真正与数据库交互。更可靠的是使用airflow tasks test来测试单个任务。数据质量测试在DAG中嵌入数据质量检查任务例如使用Great Expectations库或自定义的PythonOperator来验证数据的完整性、一致性和准确性失败则阻断流程。掌握这些模式和实践你的Airflow将不再是简单的任务触发器而是一个真正强大、可靠、可维护的数据管道编排平台。从定义依赖关系到处理复杂错误场景从本地开发到生产部署每一个环节的深入理解都能让你在构建数据基础设施时更加得心应手。