【问题标题】:How to choose how often Apache Airflow scheduler updates a DAG?如何选择 Apache Airflow 调度程序更新 DAG 的频率?
【发布时间】:2021-06-30 11:52:40
【问题描述】:

Apache Airflow documentation 中所述,我可以通过在airflow.cfg 文件中设置配置变量min_file_process_interval 来控制DAG 的更新频率:

min_file_process_interval

解析 DAG 文件的秒数。 DAG 文件每隔 min_file_process_interval 秒解析一次。对 DAG 的更新会在此时间间隔后反映出来。将此数字保持在较低水平会增加 CPU 使用率。

但是,我没有找到任何线索或最佳实践来说明我应该为 min_file_process_interval 设置哪个值。

示例

我的 DAG 每天更改一次。默认情况下,min_file_process_interval 设置为 30 秒。这意味着大多数时候更新 DAG 是没有用的:只要 DAG 不改变,更新的 DAG 和之前的 DAG 是一样的。它消耗资源并生成日志。但是,如果我每天只更新一次 DAG,如果 DAG 在每日 DAG 更新之后发生更改,或者 DAG 在运行前也更新了,我是否有运行错误 DAG 的风险?

在这种情况下我应该为min_file_process_interval 设置什么值?

编辑:正如Elad's answer 回复此问题的先前版本中所述,应避免使用动态 DAG。但是,如果我有动态 DAG,如何选择min_file_process_interval

【问题讨论】:

    标签: airflow airflow-scheduler


    【解决方案1】:

    你正在混合两种不同的东西。 min_file_process_interval 表示 Airflow 多久扫描一次 .py 文件并更新 Airflow 中的 DAG。考虑一下,当您部署新的 .py 文件时,Airflow 需要读取它并在元存储数据库中创建它 - 所以设置是关于它发生的频率。

    对于您的用例,DAG 代码不应该每天都更新 - 事实上它根本不应该更新。它应该每天运行。您的 dag 只需要能够在每个日期处理正确的文件。你的代码可以是这样的:

    from airflow.providers.ftp.sensors.ftp import FTPSensor
    
    with DAG(dag_id='stackoverflow',
             default_args=default_args,
             schedule_interval="@daily",
             catchup=False
             ) as dag:
        # Waits for a file or directory to be present on FTP.
        sensor_op = FTPSensor(
            task_id='sensor_task',
            path='/my_folder/{{ ds }}/file.csv', #path to your file in the server
            fail_on_transient_errors=False,
            ftp_conn_id='ftp_default'
        )
        
        # Operator to process the file
        operator_op = SomeOperator()
    
        sensor_op >> operator_op
    

    在该 DAG 中,它将每天开始运行 - 第一个操作员是传感器,因此如果当天的文件不存在,工作流将等到它出现时,工作流将继续到第二个操作员,这应该处理文件。 请注意,FTPSensor 的path 参数是模板化的。这意味着您可以像 {{ ds }} 一样使用 macros 这将为您提供一个包含每天日期的路径,例如:

    /my_folder/2021-05-01/file.csv
    /my_folder/2021-05-02/file.csv
    /my_folder/2021-05-03/file.csv
    

    您也可以使用path='/my_folder/{{ ds }}.csv',这将提供:

    /my_folder/2021-05-01.csv
    /my_folder/2021-05-02.csv
    /my_folder/2021-05-03.csv
    

    【讨论】:

    • 感谢@Elad!在问我的问题之前,我简化了我的示例用例,但是这样做我让你回答了错误的问题。我同意你的观点,我应该避免使用动态 DAG,而 Airflow 为我的简化示例提供了解决方案。我将删除我的问题中有关文件的所有引用。但是如果我每天都有动态 DAG 更新,那么在设置 min_file_process_interval 时应该注意什么?如果动态 DAG 不好,为什么默认每 30 秒更新一次?
    • 动态 dag 没问题。在您的问题中,您提到了一个设计,其中 dag 本身的代码每次都在变化 - 这是另一回事。动态 dag 只是一个代码,您可以从 1 个 py 文件中创建多个 dag。至于您的问题,您应该考虑的唯一与 min_file_process_interval 相关的事情是不要在运营商范围之外进行昂贵的工作。这意味着例如顶级代码中没有数据库调用。例如,不要让 task_id 依赖于您从数据库中获取的数据。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2018-12-09
    • 1970-01-01
    • 2021-09-19
    • 1970-01-01
    • 2018-07-01
    • 2018-05-13
    相关资源
    最近更新 更多