Airflow 2.4中仅在必要场景下删除Dataproc集群的实现方法
Airflow 2.4实现带资源自动清理的Dataproc ETL工作流
需求回顾
我们的ETL工作流逻辑如下:
检查文件 > 创建Dataproc集群 > 执行ETL操作 > 删除集群
核心约束:
- 若无待处理文件,跳过所有后续步骤以节省资源
- 若有文件待处理,无论ETL操作成功或失败,必须删除Dataproc集群以控制资源成本
- 使用Airflow 2.4版本,无法依赖Airflow 2.7新增的
as_setup()和as_tear_down()方法
实现方案
通过BranchPythonOperator实现分支判断,结合TriggerRule.ALL_DONE确保集群清理逻辑的执行,具体步骤如下:
1. 定义文件检查分支函数
编写Python函数判断是否存在待处理文件,返回对应分支的任务ID:
from airflow.models import Variable import os def check_pending_files(**context): # 替换为实际的文件检查逻辑,比如检查GCS/S3路径或本地目录 pending_files_path = Variable.get("pending_files_path") if os.listdir(pending_files_path): # 有文件待处理,进入创建集群分支 return "create_dataproc_cluster" else: # 无文件,进入跳过分支 return "skip_all_steps"
2. 构建完整DAG
以下是完整的DAG代码示例,包含所有任务和逻辑控制:
from airflow import DAG from airflow.operators.python import BranchPythonOperator from airflow.operators.dummy import DummyOperator from airflow.providers.google.cloud.operators.dataproc import ( DataprocCreateClusterOperator, DataprocDeleteClusterOperator, DataprocSubmitJobOperator ) from airflow.utils.trigger_rule import TriggerRule from datetime import datetime # 替换为你的Dataproc集群配置 CLUSTER_CONFIG = { "master_config": { "num_instances": 1, "machine_type_uri": "n1-standard-2", "disk_config": {"boot_disk_type": "pd-standard", "boot_disk_size_gb": 50}, }, "worker_config": { "num_instances": 2, "machine_type_uri": "n1-standard-2", "disk_config": {"boot_disk_type": "pd-standard", "boot_disk_size_gb": 50}, }, } # 替换为你的ETL作业配置 ETL_JOB_CONFIG = { "reference": {"project_id": "your-gcp-project"}, "placement": {"cluster_name": "your-dataproc-cluster"}, "pyspark_job": {"main_python_file_uri": "gs://your-bucket/etl_script.py"}, } with DAG( dag_id="dataproc_etl_with_cleanup", schedule_interval="@daily", start_date=datetime(2024, 1, 1), catchup=False, tags=["dataproc", "etl"] ) as dag: # 1. 检查待处理文件,分支判断 check_files = BranchPythonOperator( task_id="check_pending_files", python_callable=check_pending_files, provide_context=True ) # 2. 无文件时的跳过分支 skip_all = DummyOperator( task_id="skip_all_steps" ) # 3. 创建Dataproc集群 create_cluster = DataprocCreateClusterOperator( task_id="create_dataproc_cluster", project_id="your-gcp-project", cluster_config=CLUSTER_CONFIG, region="us-central1", cluster_name="your-dataproc-cluster" ) # 4. 执行ETL作业 run_etl = DataprocSubmitJobOperator( task_id="execute_etl", project_id="your-gcp-project", region="us-central1", job=ETL_JOB_CONFIG ) # 5. 删除Dataproc集群 delete_cluster = DataprocDeleteClusterOperator( task_id="delete_dataproc_cluster", project_id="your-gcp-project", region="us-central1", cluster_name="your-dataproc-cluster", # 关键:无论前面ETL成功/失败都执行 trigger_rule=TriggerRule.ALL_DONE ) # 设置任务依赖 check_files >> [create_cluster, skip_all] create_cluster >> run_etl >> delete_cluster
关键逻辑说明
- 分支控制:通过
BranchPythonOperator判断是否有文件,无文件时直接跳转到skip_all_steps任务,不会执行集群创建/删除逻辑 - 强制清理:删除集群任务设置
trigger_rule=TriggerRule.ALL_DONE,意味着只要创建集群任务执行完成(无论ETL成功或失败),都会触发集群删除,避免资源泄漏 - 兼容性:所有API均兼容Airflow 2.4版本,无需依赖高版本的
as_setup()/as_tear_down()特性
内容的提问来源于stack exchange,提问作者Khilesh Chauhan
相关产品推荐
相关产品推荐

