【问题标题】:Airflow: How to run a task on multiple workers气流:如何在多个工作人员上运行任务
【发布时间】:2018-09-26 04:00:45
【问题描述】:

我刚刚用 celery executor 设置了气流,这是我的 DAG 的骨架

dag = DAG('dummy_for_testing', default_args=default_args)

t1 = BashOperator(
    task_id='print_date',
    bash_command='date >> /tmp/dag_output.log',
    queue='test_queue',
    dag=dag)

t3 = BashOperator(
    task_id='print_host',
    bash_command='hostname >> /tmp/dag_output.log',
    queue='local_queue',
    dag=dag)

t2 = BashOperator(
    task_id='print_uptime',
    bash_command='uptime >> /tmp/dag_output.log',
    queue='local_queue',
    dag=dag)

t2.set_upstream(t3)
t2.set_upstream(t1)

我有 2 名工人。其中一个只运行一个名为local_queue 的队列,另一个运行两个名为local_queue,test_queue 的队列

我只想在 1 台机器上运行任务 1,但在两台机器上运行任务 2 和 3。即在仅运行 local_queue 的 worker 1 上,t2 和 t3 应该运行,而在同时运行 local_queue 和 test_queue 的 worker 2 上,所有 3(t1、t2 和 t3)都应该运行。 任务运行总数应为 5。

但是,当我运行它时,只运行了 3 个任务。 1) print_date 为工人 2 运行(这是正确的) 2) print_host 只为工人 1 运行(不正确。应该为两个工人运行)和 3) print_uptime 只为工人 2 运行(也不正确。应该为两个工人运行)

您能否指导我如何设置它以便运行 5 个任务。在生产中,我想通过将机器分组到队列中来管理机器,并且对于所有具有 QUEUE_A -> 执行 X 的机器和所有具有 QUEUE_B -> 执行 Y 的机器等。

谢谢

【问题讨论】:

    标签: celery airflow celery-task airflow-scheduler


    【解决方案1】:

    不是让一个工作人员处理两个队列,而是让每个工作人员处理一个队列。 所以worker命令应该是这样的:

    airflow worker -q test_queue
    airflow worker -q local_queue
    

    然后有两个相同的任务,但在不同的队列中。

    dag = DAG('dummy_for_testing', default_args=default_args)
    
    t1 = BashOperator(
        task_id='print_date',
        bash_command='date >> /tmp/dag_output.log',
        queue='test_queue',
        dag=dag)
    
    t3 = BashOperator(
        task_id='print_host',
        bash_command='hostname >> /tmp/dag_output.log',
        queue='local_queue',
        dag=dag)
    
    t3_2 = BashOperator(
        task_id='print_host_2',
        bash_command='hostname >> /tmp/dag_output.log',
        queue='test_queue',
        dag=dag)
    
    t2 = BashOperator(
        task_id='print_uptime',
        bash_command='uptime >> /tmp/dag_output.log',
        queue='local_queue',
        dag=dag)
    
    t2_2 = BashOperator(
        task_id='print_uptime_2',
        bash_command='uptime >> /tmp/dag_output.log',
        queue='test_queue',
        dag=dag)
    
    t2.set_upstream(t3)
    t2.set_upstream(t3_2)
    t2.set_upstream(t1)
    
    t2_2.set_upstream(t3)
    t2_2.set_upstream(t3_2)
    t2_2.set_upstream(t1)
    

    【讨论】:

    • 如果我有 2 台主机就可以了。拥有 10 多台主机或 100 台主机怎么样?我认为为每个主机制作任务很容易出错。如果我必须向此 DAG 设置添加主机,那么我宁愿不接触正在运行的 DAG,而是使用队列名称作为一种逻辑分组?
    • 气流不是这样工作的。如果您的任务在一台主机上失败但在另一台主机上失败怎么办?当您需要查看日志时会发生什么?查看日志是基于每个任务的,那么您正在查看哪些日志:主机 9 或主机 10?您可以在迭代队列数组时在 for 循环中创建这些任务,这样,如果您添加一个新工作人员,那么您只需在 DAG 顶部的数组中添加一个新队列。
    • 我明白你的意思。谢谢你。我想我也可以将队列名称与工作人员的主机名联系起来。我可以将工作人员定义存储在数据库中吗?然后在结果集上使用 for 循环来动态创建任务?气流会自行更新这些 DAG 吗?
    • 你可以这样做,但它可能会降低调度程序的性能(尤其是当你有很多 dag 这样做时)。调度程序会频繁处理 dag 文件,因此如果每次尝试查看 dag 的外观时都必须访问数据库,那么这可能会增加大量数据库调用。另一方面,dag 被频繁解析,因此调度程序将很快看到您所做的任何更改。请参阅调度程序设置:min_file_process_interval、max_threads 以调整 dag 解析。
    猜你喜欢
    • 2020-12-30
    • 2019-07-21
    • 1970-01-01
    • 2016-01-30
    • 2011-10-27
    • 2020-03-20
    • 1970-01-01
    • 2022-07-10
    • 1970-01-01
    相关资源
    最近更新 更多