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

如何通过Airflow的DatabricksSubmitRunOperator提交带自定义集群的Databricks工作流?

Airflow DatabricksSubmitRunOperator 自定义新集群完整示例

下面是使用DatabricksSubmitRunOperator提交带自定义新集群的Databricks工作流的完整示例,包含核心参数配置:

from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator
from airflow.utils.dates import days_ago

with DAG(
    dag_id='databricks_custom_cluster_workflow',
    start_date=days_ago(1),
    schedule_interval=None,
    catchup=False
) as dag:

    submit_databricks_job = DatabricksSubmitRunOperator(
        task_id='submit_custom_cluster_job',
        databricks_conn_id='databricks_default',  # 替换为你的Airflow Databricks连接ID
        json={
            # 自定义新集群配置
            'new_cluster': {
                'spark_version': '13.3.x-scala2.12',  # 替换为你的Spark版本
                'node_type_id': 'Standard_DS3_v2',  # 替换为你的节点类型
                'driver_node_type_id': 'Standard_DS3_v2',
                'num_workers': 2,  # 固定工作节点数量,启用自动伸缩时需移除该参数
                'autoscale': {  # 可选:自动伸缩配置
                    'min_workers': 1,
                    'max_workers': 5
                },
                'spark_conf': {
                    'spark.sql.shuffle.partitions': '200',
                    'spark.driver.memory': '8g'
                },
                'aws_attributes': {  # 云环境属性,根据你的云提供商调整(AWS/Azure/GCP)
                    'instance_profile_arn': 'arn:aws:iam::123456789012:instance-profile/databricks-instance-profile',
                    'availability': 'ON_DEMAND'
                },
                'custom_tags': {
                    'Project': 'DataPipeline',
                    'Environment': 'Production'
                },
                'runtime_engine': 'STANDARD'
            },
            # 任务配置(以Spark Python任务为例,可替换为Jar/Notebook等类型)
            'spark_python_task': {
                'python_file': 'dbfs:/path/to/your/script.py',  # 替换为你的脚本路径
                'parameters': ['--input-path', 'dbfs:/input/data', '--output-path', 'dbfs:/output/result']
            },
            # 其他可选配置
            'run_name': 'Airflow_Submitted_Custom_Cluster_Job',
            'timeout_seconds': 3600,  # 任务超时时间(秒)
            'email_notifications': {
                'on_start': ['user@example.com'],
                'on_success': ['user@example.com'],
                'on_failure': ['user@example.com']
            }
        }
    )

    submit_databricks_job

关键参数说明

  • new_cluster:必填字段,定义自定义集群的所有核心属性,需根据你的Databricks workspace环境调整版本、节点类型、云配置等。
  • 任务类型:除示例中的spark_python_task,还可使用spark_jar_task(运行Jar包)、notebook_task(运行Notebook)等,对应不同的任务配置字段。
  • databricks_conn_id:Airflow中已配置的Databricks连接ID,需确保连接已正确配置实例URL和访问令牌。
  • timeout_seconds:可选,设置任务最长运行时间,避免任务无限挂起。
  • email_notifications:可选,配置任务各状态的通知接收邮箱。

注意事项

  • 确认Spark版本和节点类型在你的Databricks workspace中可用。
  • 脚本或资源路径需正确指向DBFS或对应云存储的有效位置。
  • 启用自动伸缩时,需删除num_workers参数,仅保留autoscale配置。

内容的提问来源于stack exchange,提问作者GLC

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 21:45:29