Apache Airflow 中使用 Amazon SQS 发送通知的完整指南SqsNotifier 实战【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowAmazon SQSSimple 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 Airflow2.11.0依赖boto31.41.0、botocore1.41.0详见 providers/amazon/docs/index.rst安装命令pip install apache-airflow-providers-amazon配置一个可用的 AWS 连接默认使用aws_conn_idaws_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_idaws_default, queue_urlhttps://sqs.eu-west-1.amazonaws.com/123456789098/MyQueue, message_bodyThe Dag {{ dag.dag_id }} failed, ) task_failure_sqs_notification send_sqs_notification( aws_conn_idaws_default, region_nameeu-west-1, queue_urlhttps://sqs.eu-west-1.amazonaws.com/123456789098/MyQueue, message_bodyThe task {{ ti.task_id }} failed, ) with DAG( dag_idmydag, scheduleonce, start_datedatetime(2023, 1, 1, tzinfotimezone.utc), on_failure_callback[dag_failure_sqs_notification], catchupFalse, ): BashOperator(task_idmytask, on_failure_callback[task_failure_sqs_notification], bash_commandfail)需要说明的是代码中的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_idstr \| NoneSqsHook.default_conn_name即aws_default用于 AWS 凭证的连接 ID若为 None 或空则使用 boto3 默认行为queue_urlstr必填目标 SQS 队列的 URL消息将被发送至此message_bodystr必填要发送的消息正文message_attributesdict \| NoneNone内部转为{}消息的附加属性具体细节参见botocore.client.SQS.send_messagemessage_group_idstr \| NoneNone仅适用于 FIFO先进先出队列用于标识消息所属的分组delay_secondsint0消息延迟投递的秒数region_namestr \| NoneNoneAWS 区域名未指定时使用 boto3 默认行为这些参数都声明在template_fields元组中sqs.pytemplate_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_idself.aws_conn_id, region_nameself.region_name) def notify(self, context): Publish the notification message to Amazon SQS queue. self.hook.send_message( queue_urlself.queue_url, message_bodyself.message_body, delay_secondsself.delay_seconds, message_attributesself.message_attributes, message_group_idself.message_group_id, )hook是一个cached_property同一通知器实例在生命周期内只会创建一个SqsHook避免重复初始化连接。测试test_parameters_propagate_to_hooktest_sqs.py断言了hook属性被缓存且aws_conn_id、region_name会原样传递给SqsHook构造器。notify(context)是BaseNotifier定义的统一入口Airflow 在执行相应回调如 DAG 失败时自动调用它。除同步的notify外SqsNotifier还实现了异步版本async_notifysqs.py调用hook.asend_message完成异步发送对应测试test_async_notifytest_sqs.py。SqsHook本身继承自AwsBaseHooksqs 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.pypytest.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允许 0900 秒的延迟适合对告警做分级降噪例如失败后延迟一段时间再投递观察是否为瞬时抖动。异步执行环境在启用了 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底层 Hookproviders/amazon/src/airflow/providers/amazon/aws/hooks/sqs.py单元测试providers/amazon/tests/unit/amazon/aws/notifications/test_sqs.pyAmazon Provider 文档索引与安装要求providers/amazon/docs/index.rstAmazon 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),仅供参考