【问题标题】:How to trigger Airflow DAG from AWS SQS?如何从 AWS SQS 触发 Airflow DAG?
【发布时间】:2021-05-28 19:18:51
【问题描述】:

我想根据 SQS 消息触发 Airflow DAD。我对 Airflow 很陌生,但我认为应该这样做:

选项 1

使用Airflow SQS Sensor。据我了解,这会等待 SQS 消息继续执行已经触发 DAG。这是否意味着 DAG 总是需要运行并等待 SQS 消息捕获任何最终的新消息并处理它们?这是否也意味着我应该将我的 DAG 安排在一个非常短的时间间隔内,以便当一个 SQS 消息被一个 DAG 处理时,另一个 DAG 被创建来处理下一个 SQS 消息?

选项 2

添加一个 lambda 或其他东西来监视 SQS 消息并在需要时使用 Airflow API 触发 DAG。

最后,我想尽量减少触发 DAG 所需的交互次数,因此我想使用 Airflow 内置方式来观察 SQS。

谢谢

【问题讨论】:

    标签: airflow amazon-sqs


    【解决方案1】:

    这两个选项都有效,但是选项 2 基本上是传感器的替代实现。我认为更好的解决方案是选项 1 进行一些修改:

    使用SQSSensor,但使用mode='reschedule',每隔一段时间,传感器就会“唤醒”检查是否满足标准。请注意,这与sleep(x) 不同。当不满足条件时,Airflow 将释放工作人员以执行其他需要运行的任务并将SQSSensor 返回到调度队列。 您可以在docs 中阅读有关传感器模式的更多信息。

    from airflow.providers.amazon.aws.sensors.sqs import SQSSensor
    SQSSensor(
        task_id='test_task',
        dag=dag,
        sqs_queue='your_queue',
        aws_conn_id='aws_default',
        mode='reschedule')
    

    请注意,传感器将无限期运行,直到满足条件。您可以在传感器任务上设置timeout(还有其他可能的超时原因,如集群策略和其他默认值,但这是另一个主题)。

    【讨论】:

    • 传感器选项如何处理传入 SQS 消息的高峰负载?它会限制在预定的频率吗?这意味着如果它被安排为每秒运行一次并且每秒有 2 条传入消息到达,它就永远无法完全耗尽它?
    • 您混淆了 DAG schedule_interval 和传感器戳。一旦 DAG 运行,schedule_interval 就无关紧要了。您可以根据自己的需要将传感器设置为poke_interval
    猜你喜欢
    • 2018-01-16
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2021-12-08
    • 1970-01-01
    • 2020-05-25
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多