Apache Airflow 中使用 Amazon SQS 发送通知的完整指南(SqsNotifier 实战)
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
Amazon SQS(Simple Queue Service)通知器SqsNotifier允许用户利用 Dag 级别和 Task 级别的各种on_*_callbacks回调,向 Amazon SQS 队列推送消息。本文将基于 Apache Airflow 的 Amazon Provider 源码(sqs.py)与配套测试(test_sqs.py),讲解该通知器的完整配置方式、核心参数、模板渲染能力以及底层实现原理,让你能直接在 DAG 中接入 SQS 通知。
引言:什么是 SqsNotifier
在 Apache Airflow 中,Notifier(通知器)是一种可复用的通知组件,可以挂载到 DAG 或单个 Task 的失败、重试、成功等回调点上。SqsNotifier是 Amazon Provider 提供的通知器实现,位于airflow.providers.amazon.aws.notifications.sqs模块,其作用是将一条消息发送到指定的 Amazon SQS 队列。
它基于BaseNotifier构建(见 providers/common/compat 的 notifier 兼容层),因此与 Airflow 内置的其他 Notifier 使用方式完全一致,可以在on_failure_callback、on_success_callback、on_retry_callback等位置直接引用。SQS 本身具备高可用、可持久化、削峰填谷的特性,非常适合作为任务失败告警的投递通道,再配合下游的 Lambda、EC2 轮询或告警平台完成最终消费。
环境准备
使用SqsNotifier需要满足以下前提:
已安装 Amazon Provider 包(当前仓库中该 Provider 的版本为 9.36.0,最低要求 Apache Airflow
>=2.11.0,依赖boto3>=1.41.0、botocore>=1.41.0,详见 providers/amazon/docs/index.rst),安装命令:pip install apache-airflow-providers-amazon配置一个可用的 AWS 连接(默认使用
aws_conn_id="aws_default"),用于提供访问 SQS 的凭证。在 AWS 控制台或通过
aws sqs create-queue预先创建目标队列,并获得其QueueUrl。
示例代码
原文档给出了一个完整的示例,展示了如何在 DAG 级和 Task 级同时使用 SQS 通知:
from datetime import datetime, timezone from airflow import DAG from airflow.providers.standard.operators.bash import BashOperator from airflow.providers.amazon.aws.notifications.sqs import send_sqs_notification dag_failure_sqs_notification = send_sqs_notification( aws_conn_id="aws_default", queue_url="https://sqs.eu-west-1.amazonaws.com/123456789098/MyQueue", message_body="The Dag {{ dag.dag_id }} failed", ) task_failure_sqs_notification = send_sqs_notification( aws_conn_id="aws_default", region_name="eu-west-1", queue_url="https://sqs.eu-west-1.amazonaws.com/123456789098/MyQueue", message_body="The task {{ ti.task_id }} failed", ) with DAG( dag_id="mydag", schedule="@once", start_date=datetime(2023, 1, 1, tzinfo=timezone.utc), on_failure_callback=[dag_failure_sqs_notification], catchup=False, ): BashOperator(task_id="mytask", on_failure_callback=[task_failure_sqs_notification], bash_command="fail")需要说明的是,代码中的send_sqs_notification实际上是SqsNotifier的别名。从 sqs.py 源码可见:
send_sqs_notification = SqsNotifier测试 test_sqs.py 也专门验证了这一点:
def test_class_and_notifier_are_same(self): assert send_sqs_notification is SqsNotifier因此,上述两个名字可以互换使用。
核心参数详解
SqsNotifier的构造函数(sqs.py)支持以下参数:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
aws_conn_id | str \| None | SqsHook.default_conn_name(即"aws_default") | 用于 AWS 凭证的连接 ID;若为 None 或空,则使用 boto3 默认行为 |
queue_url | str | 必填 | 目标 SQS 队列的 URL,消息将被发送至此 |
message_body | str | 必填 | 要发送的消息正文 |
message_attributes | dict \| None | None(内部转为{}) | 消息的附加属性,具体细节参见botocore.client.SQS.send_message |
message_group_id | str \| None | None | 仅适用于 FIFO(先进先出)队列,用于标识消息所属的分组 |
delay_seconds | int | 0 | 消息延迟投递的秒数 |
region_name | str \| None | None | AWS 区域名;未指定时使用 boto3 默认行为 |
这些参数都声明在template_fields元组中(sqs.py):
template_fields: Sequence[str] = ( "queue_url", "message_body", "message_attributes", "message_group_id", "delay_seconds", "aws_conn_id", "region_name", )这意味着所有参数都支持 Jinja 模板渲染——queue_url、message_body甚至aws_conn_id都可以在运行时根据 DAG 上下文动态解析。这一点在测试 test_sqs_notifier_templated 中得到验证:测试使用"{{ dag.dag_id }}"、"https://sqs.{{ var_region }}.amazonaws.com/{{ var_account }}/{{ var_queue }}"、"The {{ var_username|capitalize }} Show"等模板表达式,最终在运行时被渲染为具体的 dag_id、区域、账号、队列名与消息正文。
工作原理:从 Notifier 到 SQS 的调用链
SqsNotifier的底层实现非常简洁,核心是hook与notify方法:
@cached_property def hook(self) -> SqsHook: """Amazon SQS Hook (cached).""" return SqsHook(aws_conn_id=self.aws_conn_id, region_name=self.region_name) def notify(self, context): """Publish the notification message to Amazon SQS queue.""" self.hook.send_message( queue_url=self.queue_url, message_body=self.message_body, delay_seconds=self.delay_seconds, message_attributes=self.message_attributes, message_group_id=self.message_group_id, )hook是一个cached_property:同一通知器实例在生命周期内只会创建一个SqsHook,避免重复初始化连接。测试test_parameters_propagate_to_hook(test_sqs.py)断言了hook属性被缓存,且aws_conn_id、region_name会原样传递给SqsHook构造器。notify(context)是BaseNotifier定义的统一入口,Airflow 在执行相应回调(如 DAG 失败)时自动调用它。- 除同步的
notify外,SqsNotifier还实现了异步版本async_notify(sqs.py),调用hook.asend_message完成异步发送,对应测试test_async_notify(test_sqs.py)。
SqsHook本身继承自AwsBaseHook(sqs hook 源码),构造时将client_type固定为"sqs",因此底层对接的是 boto3 的SQS.Client。其send_message方法会构建如下参数并调用get_conn().send_message(**params)(sqs.py 第 80-114 行):
{ "QueueUrl": queue_url, "MessageBody": message_body, "DelaySeconds": delay_seconds, "MessageAttributes": message_attributes or {}, "MessageGroupId": message_group_id, "MessageDeduplicationId": message_deduplication_id, }_build_msg_params会通过prune_dict剔除值为None的键,确保 FIFO 专用的MessageGroupId、MessageDeduplicationId等参数仅在需要时才会传给 boto3。
参数传播与模板渲染的源码验证
为了让读者确信上述行为,以下给出测试中的关键断言(test_sqs.py):
@pytest.mark.parametrize("aws_conn_id", ["aws_test_conn_id", None, PARAM_DEFAULT_VALUE]) @pytest.mark.parametrize("region_name", ["eu-west-2", None, PARAM_DEFAULT_VALUE]) def test_parameters_propagate_to_hook(self, aws_conn_id, region_name): notifier = SqsNotifier(**notifier_kwargs, **SEND_MSG_KWARGS) with mock.patch("airflow.providers.amazon.aws.notifications.sqs.SqsHook") as mock_hook: hook = notifier.hook assert hook is notifier.hook, "Hook property not cached" mock_hook.assert_called_once_with( aws_conn_id=(aws_conn_id if aws_conn_id is not NOTSET else "aws_default"), region_name=(region_name if region_name is not NOTSET else None), ) notifier.notify({}) mock_hook.return_value.send_message.assert_called_once_with(**SEND_MSG_KWARGS)该测试同时验证了三件事:
aws_conn_id与region_name正确传播到SqsHook;未显式指定aws_conn_id时默认使用"aws_default"。hook属性被缓存(两次访问返回同一对象)。notify触发send_message,且queue_url、message_body、delay_seconds、message_attributes、message_group_id全部按构造时的值传递。
模板渲染测试(test_sqs.py 第 65-94 行)则证明:当把aws_conn_id="{{ dag.dag_id }}"、region_name="{{ var_region }}"、queue_url等参数写成模板时,最终 Hook 收到的是渲染后的实际值(如"test_sqs_notifier_templated"、"ca-central-1"、"https://sqs.ca-central-1.amazonaws.com/123321123321/AwesomeQueue"、"The Truman Show")。这意味着你可以把队列 URL、区域甚至连接 ID 都做成可配置的变量,实现环境无关的 DAG。
进阶实践建议
- DAG 级与 Task 级回调的组合使用:将通知器放入
on_failure_callback数组(注意要写成[notifier]列表形式),可以同时挂多个通知器,例如同时向 SQS 和 Slack 发送告警。 - FIFO 队列注意事项:如果目标队列是 FIFO 类型,必须指定
message_group_id(必要时配合message_deduplication_id,当前 Notifier 参数中message_deduplication_id尚未直接暴露,可通过扩展或直接使用SqsHook实现)。SQS 对 FIFO 队列要求MessageGroupId非空,否则send_message会报错。 - 延迟投递:
delay_seconds允许 0~900 秒的延迟,适合对告警做分级降噪,例如失败后延迟一段时间再投递,观察是否为瞬时抖动。 - 异步执行环境:在启用了 Deferrable/异步执行的环境中,
async_notify可以避免阻塞事件循环;它通过SqsHook.asend_message使用异步连接发送(sqs hook 第 116-156 行)。 - 模板化环境隔离:利用
template_fields的特性,将queue_url、region_name、message_body全部模板化,配合 Airflow Variables 或{{ conn.xxx }}语法,一套 DAG 即可在多个环境间复用。
参考链接
- 通知器源码:providers/amazon/src/airflow/providers/amazon/aws/notifications/sqs.py
- 底层 Hook:providers/amazon/src/airflow/providers/amazon/aws/hooks/sqs.py
- 单元测试:providers/amazon/tests/unit/amazon/aws/notifications/test_sqs.py
- Amazon Provider 文档索引与安装要求:providers/amazon/docs/index.rst
- Amazon Provider 通知指南索引:providers/amazon/docs/notifications/index.rst
【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考