【问题标题】:Apache Airflow - trigger/schedule DAG rerun on completion (File Sensor)Apache Airflow - 在完成时触发/安排 DAG 重新运行(文件传感器)
【发布时间】:2017-11-29 21:42:45
【问题描述】:

早安。

我也在尝试设置 DAG

  1. 监视/感知文件以访问网络文件夹
  2. 处理文件
  3. 存档文件

使用在线教程和 stackoverflow,我已经能够想出成功实现目标的以下 DAG 和 Operator,但是我希望 DAG 在完成时重新安排或重新运行,以便它开始监视/感应另一个文件.

我尝试设置变量max_active_runs:1,然后设置schedule_interval: timedelta(seconds=5),这是重新安排 DAG 但开始排队任务并锁定文件。

关于如何在 archive_task 之后重新运行 DAG,欢迎提出任何想法?

谢谢

DAG 代码

from airflow import DAG
from airflow.operators import PythonOperator, OmegaFileSensor, ArchiveFileOperator
from datetime import datetime, timedelta
from airflow.models import Variable

default_args = {
    'owner': 'glsam',
    'depends_on_past': False,
    'start_date': datetime.now(),
    'provide_context': True,
    'retries': 100,
    'retry_delay': timedelta(seconds=30),
    'max_active_runs': 1,
    'schedule_interval': timedelta(seconds=5),
}

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

filepath = Variable.get("soucePath_Test")
filepattern = Variable.get("filePattern_Test")
archivepath = Variable.get("archivePath_Test")

sensor_task = OmegaFileSensor(
    task_id='file_sensor_task',
    filepath=filepath,
    filepattern=filepattern,
    poke_interval=3,
    dag=dag)


def process_file(**context):
    file_to_process = context['task_instance'].xcom_pull(
        key='file_name', task_ids='file_sensor_task')
    file = open(filepath + file_to_process, 'w')
    file.write('This is a test\n')
    file.write('of processing the file')
    file.close()


proccess_task = PythonOperator(
    task_id='process_the_file', 
    python_callable=process_file,
    provide_context=True,
    dag=dag
)

archive_task = ArchiveFileOperator(
    task_id='archive_file',
    filepath=filepath,
    archivepath=archivepath,
    dag=dag)

sensor_task >> proccess_task >> archive_task

文件传感器操作员

    import os
    import re

    from datetime import datetime
    from airflow.models import BaseOperator
    from airflow.plugins_manager import AirflowPlugin
    from airflow.utils.decorators import apply_defaults
    from airflow.operators.sensors import BaseSensorOperator


    class ArchiveFileOperator(BaseOperator):
        @apply_defaults
        def __init__(self, filepath, archivepath, *args, **kwargs):
            super(ArchiveFileOperator, self).__init__(*args, **kwargs)
            self.filepath = filepath
            self.archivepath = archivepath

        def execute(self, context):
            file_name = context['task_instance'].xcom_pull(
                'file_sensor_task', key='file_name')
            os.rename(self.filepath + file_name, self.archivepath + file_name)


    class OmegaFileSensor(BaseSensorOperator):
        @apply_defaults
        def __init__(self, filepath, filepattern, *args, **kwargs):
            super(OmegaFileSensor, self).__init__(*args, **kwargs)
            self.filepath = filepath
            self.filepattern = filepattern

        def poke(self, context):
            full_path = self.filepath
            file_pattern = re.compile(self.filepattern)

            directory = os.listdir(full_path)

            for files in directory:
                if re.match(file_pattern, files):
                    context['task_instance'].xcom_push('file_name', files)
                    return True
            return False


    class OmegaPlugin(AirflowPlugin):
        name = "omega_plugin"
        operators = [OmegaFileSensor, ArchiveFileOperator]

【问题讨论】:

    标签: triggers airflow directed-acyclic-graphs


    【解决方案1】:

    Dmitris 方法效果很好。

    我还在阅读设置中发现schedule_interval=None,然后使用 TriggerDagRunOperator 也同样有效

    trigger = TriggerDagRunOperator(
        task_id='trigger_dag_RBCPV99_rerun',
        trigger_dag_id="RBCPV99_v2",
        dag=dag)
    
    sensor_task >> proccess_task >> archive_task >> trigger
    

    【讨论】:

    • 在不使用python_callable 的情况下,你是如何让它工作的?我有一个类似的用例,想知道你是怎么做到的。谢谢!
    • @CodingInCircles AFAICS python_callableTriggerDagRunOperator中的一个可选参数
    【解决方案2】:

    设置schedule_interval=None 并使用来自BashOperatorairflow trigger_dag 命令在前一个执行完成时启动下一个执行。

    trigger_next = BashOperator(task_id="trigger_next", 
               bash_command="airflow trigger_dag 'your_dag_id'", dag=dag)
    
    sensor_task >> proccess_task >> archive_task >> trigger_next
    

    您可以使用相同的airflow trigger_dag 命令手动启动您的第一次运行,然后trigger_next 任务将自动触发下一次运行。我们现在在生产中使用了很多个月,并且运行良好。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2018-07-01
      • 1970-01-01
      • 1970-01-01
      • 2020-08-14
      • 1970-01-01
      • 2022-06-24
      • 2021-12-27
      相关资源
      最近更新 更多