【问题标题】:How do I pass custom data into the DatabricksRunNowOperator in airflow如何将自定义数据传递到气流中的 DatabricksRunNowOperator
【发布时间】:2022-08-10 20:44:31
【问题描述】:

我正在尝试创建一个使用 DatabricksRunNowOperator 运行 pyspark 的 DAG。 但是我无法弄清楚如何访问 pyspark 脚本中的气流配置。

parity_check_run = DatabricksRunNowOperator(
    task_id=\'my_task\',
    databricks_conn_id=\'databricks_default\',
    job_id=\'1837\',
    spark_submit_params=[\"file.py\", \"pre-defined-param\"],
    dag=dag,
)

我尝试通过kwargs 访问它,但这似乎不起作用。

  • 工作是如何定义的——它是笔记本、python 文件、轮子还是其他东西?

标签: apache-spark pyspark airflow databricks


【解决方案1】:

您可以使用 notebook_params 参数,如 documentation 中所示。

例如:

job_id=42

notebook_params = {
    "dry-run": "true",
    "oldest-time-to-consider": "1457570074236"
}

notebook_run = DatabricksRunNowOperator(
    job_id=job_id,
    notebook_params=notebook_params,

)

然后您可以通过 PySpark 代码中的dbutils.widgets.get("oldest-time-to-consider") 访问该值。

【讨论】:

    【解决方案2】:

    DatabricksRunNowOperator 支持为现有作业提供参数的不同方式,具体取决于作业的定义方式 (doc):

    • notebook_params 如果您使用笔记本 - 它是小部件名称的字典 -> 值。您可以使用dbutils.widgets.get 获取参数
    • python_params - 将传递给 Python 任务的参数列表 - 您可以通过 sys.argv 获取它们
    • jar_params - 将传递给 Jar 任务的参数列表。您可以像往常一样为 Java/Scala 程序获取它们
    • spark_submit_params - 将传递给 spark-submit 的参数列表

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2023-03-05
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-05-18
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多