Airflow DatabricksSubmitRunOperator是否支持Python Wheel及替代方案咨询
用Airflow的DatabricksSubmitRunOperator提交Python Wheel作业的方案
完全可以实现,虽然Airflow官方文档里的DatabricksSubmitRunOperator示例只展示了python_file类型,但该Operator本质是封装了Databricks的Run Submit API,只要传入符合API规范的参数,就能提交Python Wheel类型的作业。
具体实现代码
from airflow.providers.databricks.operators.databricks import DatabricksSubmitRunOperator submit_wheel_task = DatabricksSubmitRunOperator( task_id="run_python_wheel_job", databricks_conn_id="databricks_default", # 可选:使用现有集群,或者替换为new_cluster配置新集群 existing_cluster_id="your_cluster_id_here", libraries=[ # 指定wheel包位置,支持DBFS、S3、ADLS或PyPI包名 {"whl": "dbfs:/wheels/your_package-1.0.0-py3-none-any.whl"} ], run_name="Airflow Triggered Python Wheel Job", # 配置Python Wheel任务参数 python_wheel_task={ "package_name": "your_package", "entry_point": "main", # 包中定义的入口函数/脚本 "parameters": ["--env", "prod", "--input-path", "dbfs:/data/input"] # 传递给入口的参数 } )
关键说明
python_wheel_task是Databricks API原生支持的任务类型,Airflow的Operator会直接将该参数传递给Databricks后台,无需额外配置。- 如果需要动态创建集群,将
existing_cluster_id替换为new_cluster参数,格式与创建集群的API参数一致,例如:new_cluster={ "spark_version": "13.3.x-scala2.12", "node_type_id": "m5.xlarge", "num_workers": 2, "spark_conf": {"spark.databricks.delta.preview.enabled": "true"} } libraries中的wheel路径可以是DBFS上的文件,或者公网PyPI包名(如{"pypi": {"package": "your-package==1.0.0"}}),只要集群能访问到即可。
替代方案(若遇到版本兼容问题)
如果使用的Airflow Databricks Provider版本过低,不支持直接传递python_wheel_task,可以通过raw_json参数直接传入完整的API请求体:
submit_wheel_task = DatabricksSubmitRunOperator( task_id="run_python_wheel_job", databricks_conn_id="databricks_default", raw_json={ "existing_cluster_id": "your_cluster_id_here", "libraries": [{"whl": "dbfs:/wheels/your_package-1.0.0-py3-none-any.whl"}], "run_name": "Airflow Raw API Wheel Job", "python_wheel_task": { "package_name": "your_package", "entry_point": "main", "parameters": ["--env", "prod"] } } )
内容的提问来源于stack exchange,提问作者linpingta
相关产品推荐
相关产品推荐

