【问题标题】:Airflow variables getting updated even if the DAG is not running即使 DAG 未运行,气流变量也会更新
【发布时间】:2021-07-05 14:19:28
【问题描述】:

我正在从气流变量中读取一个整数变量,然后每次 DAG 运行时将该值加一并再次将其设置为该变量。

但是在下面的代码之后,每次刷新页面时,UI 处的变量都会发生变化。 不知道是什么导致了这种行为

counter = Variable.get('counter')
s = BashOperator(
    task_id='echo_start_variable',
    bash_command='echo ' + counter,
    dag=dag,
)
Variable.set("counter", int(counter) + 1)

sql_query = "SELECT * FROM UNNEST(SEQUENCE({start}, {end}))"
sql_query = sql_query.replace('{start}', start).replace('{end}', end)
submit_query = PythonOperator(
    task_id='submit_athena_query',
    python_callable=run_athena_query,
    op_kwargs={'query': sql_query, 'db': 'db',
               's3_output': 's3://s3-path/rohan/date=' + current_date + '/'},
    dag=dag)

e = BashOperator(
    task_id='echo_end_variable',
    bash_command='echo ' + counter,
    dag=dag,
)

s >> submit_query >> e

【问题讨论】:

    标签: python variables operators airflow directed-acyclic-graphs


    【解决方案1】:

    气流进程每 30 秒处理一次 DAG 文件(默认为 min_file_process_interval 设置),这意味着您拥有的任何顶级代码每 30 秒运行一次,所以 Variable.set("counter", int(counter) + 1) 将导致变量计数器每 30 秒增加 1。

    在顶级代码中与变量交互是一种不好的做法(无论增加值问题如何)。它每 30 秒打开一个到 Metastore 数据库的连接,这可能会导致严重的问题并使数据库不堪重负。

    要获取变量的值,您可以使用 Jinja:

    e = BashOperator(
        task_id='echo_end_variable',
        bash_command='echo {{ var.value.counter }}',
        dag=dag,
    )
    

    这是一种使用变量的安全方法,因为只有在执行运算符时才会检索值。

    如果您想将变量的值增加 1,请使用PythonOpeartor

    def increase():
        counter = Variable.get('counter')
        Variable.set("counter", int(counter) + 1)
    
    increase_op = PythonOperator(
        task_id='increase_task',
        python_callable=increase,
        dag=dag)
    

    只有在操作符运行时才会执行python callable。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2022-11-02
      • 2021-05-03
      • 2019-01-08
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2019-11-18
      相关资源
      最近更新 更多