如何将Lambda执行Athena查询的结果自动推送至SNS主题
具体实现方案
前置准备
提前确认以下资源已就绪:
- 接收通知的SNS主题,记录其ARN值
- Athena查询使用的工作组、对应查询结果输出S3桶路径
- 已绑定EventBridge定时触发规则的Lambda函数,运行环境推荐使用Python 3.9及以上版本(内置AWS SDK boto3无需额外安装依赖)
Lambda函数逻辑改造
核心逻辑分为触发Athena查询→等待查询完成→拉取查询结果→格式化后推送SNS四个环节,参考代码如下:
import boto3 import time # 配置参数,替换为你自己的实际值 ATHENA_WORKGROUP = "primary" ATHENA_OUTPUT_S3 = "s3://你的Athena结果桶路径/" SNS_TOPIC_ARN = "arn:aws-cn:sns:区域:账号ID:主题名" QUERY_SQL = 'SELECT instance_id, disk_path, "use%" FROM custom_diskutilization WHERE "use%" > 70' athena_client = boto3.client('athena') sns_client = boto3.client('sns') def lambda_handler(event, context): # 1. 提交Athena查询 query_response = athena_client.start_query_execution( QueryString=QUERY_SQL, WorkGroup=ATHENA_WORKGROUP, ResultConfiguration={"OutputLocation": ATHENA_OUTPUT_S3} ) query_execution_id = query_response['QueryExecutionId'] # 2. 轮询等待查询完成 while True: query_status = athena_client.get_query_execution(QueryExecutionId=query_execution_id) state = query_status['QueryExecution']['Status']['State'] if state in ['SUCCEEDED', 'FAILED', 'CANCELLED']: break time.sleep(2) if state != 'SUCCEEDED': # 查询失败发送错误通知 error_msg = f"磁盘使用率查询失败,错误原因:{query_status['QueryExecution']['Status']['StateChangeReason']}" sns_client.publish(TopicArn=SNS_TOPIC_ARN, Subject="磁盘使用率查询异常", Message=error_msg) return # 3. 拉取查询结果 query_results = athena_client.get_query_results(QueryExecutionId=query_execution_id) rows = query_results['ResultSet']['Rows'] # 4. 格式化结果 if len(rows) <= 1: message = "当前无磁盘使用率超过70%的实例" else: message = "磁盘使用率超过70%的实例列表如下:\n" # 提取表头 headers = [col['VarCharValue'] for col in rows[0]['Data']] message += " | ".join(headers) + "\n" message += "-"*50 + "\n" # 提取数据行 for row in rows[1:]: row_data = [col.get('VarCharValue', '') for col in row['Data']] message += " | ".join(row_data) + "\n" # 5. 发送SNS通知 sns_client.publish( TopicArn=SNS_TOPIC_ARN, Subject="每日高磁盘使用率实例告警", Message=message )
IAM权限配置
给Lambda的执行角色添加以下权限策略,替换其中的占位符为你的实际资源ARN:
{ "Version": "2012-10-17", "Statement": [ { "Effect": "Allow", "Action": [ "athena:StartQueryExecution", "athena:GetQueryExecution", "athena:GetQueryResults" ], "Resource": "*" }, { "Effect": "Allow", "Action": "s3:GetObject", "Resource": "arn:aws-cn:s3:::你的Athena结果桶/*" }, { "Effect": "Allow", "Action": "sns:Publish", "Resource": "你的SNS主题ARN" } ] }
验证步骤
- 替换代码和权限策略中的所有占位符后,保存Lambda函数
- 手动触发Lambda测试,确认SNS订阅端(邮箱、企业微信机器人等)能正常收到告警通知
- 验证EventBridge定时规则的触发时间配置正确,次日确认能自动收到通知
注意事项
- 如果查询结果数据量较大,建议不要直接把全量结果放在SNS消息中,改为在消息中附结果文件的S3访问链接即可,避免超出SNS单条消息大小限制
- 可以根据需要调整消息格式,比如改为JSON格式方便自动化消费
- 轮询等待查询的超时时间不要超过Lambda函数的最大执行时长限制
内容的提问来源于stack exchange,提问作者FJ777
相关产品推荐
相关产品推荐

