【问题标题】:AirFlow DatabricksSubmitRunOperator does not take in notebook parametersAirFlow DatabricksSubmitRunOperator 不接受笔记本参数
【发布时间】:2020-05-01 12:37:00
【问题描述】:

我正在尝试从 Airflow 触发笔记本。笔记本有定义为小部件的参数,我试图通过 notebook_params 参数将值传递给它,虽然它触发了,但当我查看提交的作业时,似乎没有传递参数。

例如代码

new_cluster = {'spark_version': '6.5.x-cpu-ml-scala2.11',
                        'node_type_id': 'Standard_DS3_v2',
                        'num_workers': 4
                        }

notebook_task = DatabricksSubmitRunOperator(task_id='notebook_task',
             json={'new_cluster': new_cluster,
                                'notebook_task': {
                                    'notebook_path': '/Users/abc@test.com/Demo',
                                    'notebook_parameters':'{"fromdate":"20200420","todate":"20200420", "datalakename":"exampledatalake", "dbname": "default", "filesystem":"refined" , "tablename":"ntcsegmentprediction", "modeloutputpath":"curated"}'
                                },
                            })

但是,DatabricksRunNowOperator 支持它,并且它可以工作

notebook_run = DatabricksRunNowOperator(task_id='notebook_task',
            job_id=24,
            notebook_params={"fromdate":"20200420","todate":"20200420", "datalakename":"exampledatalake", "dbname": "default", "filesystem":"refined" , "tablename":"ntcsegmentprediction", "modeloutputpath":"curated"}
        )

here中DatabricksSubmitRunOperator的文档和源代码中

它说它可以接受 notebook_task。如果可以,不知道为什么不能接受参数

我错过了什么?

如果需要更多信息,我也可以提供。

【问题讨论】:

  • 你有想过这个吗?我也有同样的问题
  • 我不得不使用 DatabricksRunNowOperator。创建了一个 Databricks 作业并使用它调用它。然后参数正确传递。不确定 DatabricksSubmitRunOperator 有什么问题。您可能还想使用 DatabricksRunNowOperator。
  • 嘿 Saugat 我也在尝试从 Airflow 触发笔记本,请指导您如何解决此问题。请建议这将是巨大的帮助。
  • 请像我说的那样使用 DatabricksRunNowOperator,并在下面提供了一个示例。创建一个作业,然后传递该作业的 id 和参数。再次 - 示例在问题本身中。

标签: airflow databricks azure-databricks


【解决方案1】:

您应该使用base_parameters 而不是notebook_params

https://docs.databricks.com/dev-tools/api/latest/jobs.html#jobsnotebooktask

【讨论】:

    【解决方案2】:

    要将其与 DatabricksSubmitRunOperator 一起使用,您需要将其作为额外参数添加到 json 参数中:ParamPair

    notebook_task_params = {
        'new_cluster': cluster_def,
        'notebook_task': {
            'notebook_path': 'path',
            'base_parameters':{
                "param1": "**",
                "param2": "**"}            
        }   
    }
    notebook_task = DatabricksSubmitRunOperator(
        task_id='id***',
        dag=dag,
        trigger_rule=TriggerRule.ALL_DONE,
        json=notebook_task_params)
    

    然后您可以只使用 dbutils.widgets.get 来检索值或设置默认值。

    param1 = getArgument("param1", "default")
    param2 = getArgument("param2", "default")
    

    getArgument > (DEPRECATED) 等价于 get

    【讨论】:

    • 如果你看我的问题,我也做了同样的事情。问题是,正如 Gustavo 在他的回答中指出的那样,我使用的是 notebook_params 而不是 base_parameters。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2017-12-10
    • 1970-01-01
    • 2018-10-22
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多