Параметры уровня DAG для строк и целых чисел ⇐ Python
-
Anonymous
Параметры уровня DAG для строк и целых чисел
У меня есть Airflow SparkSubmitOperator, и я хочу добавить параметры уровня DAG для настройки:
executor_cores executor_memory application_args Я написал следующий код:
dag = DAG( 'spark_submit_ex', параметры={ "executor_cores": Param(2, type="integer", минимум=2), "executor_memory": Param("4g", type="string"), "some_id": Param("file_sink", type="string") }, …. ) spark_job = SparkSubmitOperator( executor_cores = '{{ params.executor_cores }}', driver_memory = '{{ params.driver_memory }}', приложение = '0.0.1-SNAPSHOT.jar', application_args=['{{ params.some_id }}'], ….. В журналах воздушного потока я вижу правильное значение для application_args и неправильное для executor_cores и driver_memory:
- Spark-Submit cmd: spark-submit ….. --executor-cores {{ params.executor_cores }} --executor-memory {{ params.executor_memory }} …. 0.0.1-SNAPSHOT.jar --workspaceId file_sink Я также пытался использовать двойные кавычки ( executor_cores='{{ params.executor_cores }}'), но в этом случае DAG не запускался.
У меня есть Airflow SparkSubmitOperator, и я хочу добавить параметры уровня DAG для настройки:
executor_cores executor_memory application_args Я написал следующий код:
dag = DAG( 'spark_submit_ex', параметры={ "executor_cores": Param(2, type="integer", минимум=2), "executor_memory": Param("4g", type="string"), "some_id": Param("file_sink", type="string") }, …. ) spark_job = SparkSubmitOperator( executor_cores = '{{ params.executor_cores }}', driver_memory = '{{ params.driver_memory }}', приложение = '0.0.1-SNAPSHOT.jar', application_args=['{{ params.some_id }}'], ….. В журналах воздушного потока я вижу правильное значение для application_args и неправильное для executor_cores и driver_memory:
- Spark-Submit cmd: spark-submit ….. --executor-cores {{ params.executor_cores }} --executor-memory {{ params.executor_memory }} …. 0.0.1-SNAPSHOT.jar --workspaceId file_sink Я также пытался использовать двойные кавычки ( executor_cores='{{ params.executor_cores }}'), но в этом случае DAG не запускался.