You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 14:15:09