相关疑难解决方法(0)

气流ExternalTask​​Sensor卡住了

我正在尝试使用ExternalTask​​Sensor,并且它已经陷入了另一个已经成功完成的DAG任务.

这里,第一个DAG"a"完成其任务,之后应该触发通过ExternalTask​​Sensor的第二个DAG"b".相反,它陷入了寻找a.first_task的困境.

第一个DAG:

import datetime
from airflow import DAG
from airflow.operators.python_operator import PythonOperator

dag = DAG(
    dag_id='a',
    default_args={'owner': 'airflow', 'start_date': datetime.datetime.now()},
    schedule_interval=None
)

def do_first_task():
    print('First task is done')

PythonOperator(
    task_id='first_task',
    python_callable=do_first_task,
    dag=dag)
Run Code Online (Sandbox Code Playgroud)

第二个DAG:

import datetime
from airflow import DAG
from airflow.operators.python_operator import PythonOperator
from airflow.operators.sensors import ExternalTaskSensor

dag = DAG(
    dag_id='b',
    default_args={'owner': 'airflow', 'start_date': datetime.datetime.now()},
    schedule_interval=None
)

def do_second_task():
    print('Second task is done')

ExternalTaskSensor(
    task_id='wait_for_the_first_task_to_be_completed',
    external_dag_id='a',
    external_task_id='first_task',
    dag=dag) >> \
PythonOperator(
    task_id='second_task',
    python_callable=do_second_task,
    dag=dag)
Run Code Online (Sandbox Code Playgroud)

我在这里错过了什么?

python airflow

10
推荐指数
2
解决办法
7393
查看次数

气流外部传感器卡在外面

我想在另一个dag完成后开始一个dag.一个解决方案是使用外部传感器功能,下面你可以找到我的解决方案 我遇到的问题是依赖的dag卡在戳,我检查了这个答案 ,并确保两个dags运行在相同的时间表,我的简化代码如下:任何帮助将不胜感激.领导者dag:

from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta


default_args = {
   'owner': 'airflow',
   'depends_on_past': False,
   'start_date': datetime(2015, 6, 1),
   'retries': 1,
   'retry_delay': timedelta(minutes=5),



 }

 schedule = '* * * * *'

 dag = DAG('leader_dag', default_args=default_args,catchup=False, 
 schedule_interval=schedule)

t1 = BashOperator(
   task_id='print_date',
   bash_command='date',
   dag=dag)
Run Code Online (Sandbox Code Playgroud)

依赖的dag:

from airflow import DAG
from airflow.operators.bash_operator import BashOperator
from datetime import datetime, timedelta
from airflow.operators.sensors import ExternalTaskSensor


default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2018, 10, 8), …
Run Code Online (Sandbox Code Playgroud)

airflow airflow-scheduler

8
推荐指数
1
解决办法
1727
查看次数

将顶级DAG连接在一起

我需要有几个相同(只有在不同参数)的顶级DAG小号那也可以用以下限制/假设一起触发:

  • 各个顶级DAG将会拥有,schedule_interval=None因为它们仅需要偶尔的手动触发
  • 但是,一系列DAG需要每天运行
  • 序列中DAG的顺序数量是固定的(在编写代码之前已知),并且很少更改(几个月内一次)
  • 无论DAG是失败还是成功,触发链都不得中断
  • 当前,它们必须串联运行。将来他们可能需要并行触发

因此,我为dags目录中的每个DAG创建了一个文件,现在必须将它们连接起来以便顺序执行。我确定了两种方法可以完成此操作:

  1. SubDagOperator

  2. TriggerDagRunOperator

    • 可在我的演示中使用,并行运行(不按顺序运行),因为它不等待触发的DAG完成才移至下一个
    • ExternalTaskSensor 可能有助于克服上述限制,但会使事情变得很混乱

我的问题是

  • 如何克服的局限性parent_id前缀dag_idSubDagS'
  • 如何迫使TriggerDagRunOperatorS 等待DAG完成
  • 是否有其他替代/更好的方法可以独立的(顶级)DAG连接在一起?
  • 我为每个顶级DAG 创建单独文件(仅在输入方面有所不同的DAG)的方法是否有解决方法?

我正在使用puckel / …

airflow

6
推荐指数
1
解决办法
1073
查看次数

标签 统计

airflow ×3

airflow-scheduler ×1

python ×1