我正在尝试设置一个DAG,它将响应云发布/订阅消息。它需要我在我的DAG代码中添加以下导入语句: from airflow.providers.google.cloud.operators.pubsub import ()
from airflow.providers.google.cloud.sensors.pubsub import
在Google PubSub中,每个源都有一个单独的主题;当更新一个源时,它在相应的主题订阅中发送一条消息。当所有源被更新时(即当每个订阅中至少有一条新消息时),作业就可以启动。DAG从一系列并行任务开始,每个订阅都检查是否发布了新消息,但没有对其进行分类。下一个任务等待所有前面的任务,并使用XCOM查看 all 是否包含一条消息</em