You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何将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"
        }
    ]
}

验证步骤

  1. 替换代码和权限策略中的所有占位符后,保存Lambda函数
  2. 手动触发Lambda测试,确认SNS订阅端(邮箱、企业微信机器人等)能正常收到告警通知
  3. 验证EventBridge定时规则的触发时间配置正确,次日确认能自动收到通知

注意事项

  • 如果查询结果数据量较大,建议不要直接把全量结果放在SNS消息中,改为在消息中附结果文件的S3访问链接即可,避免超出SNS单条消息大小限制
  • 可以根据需要调整消息格式,比如改为JSON格式方便自动化消费
  • 轮询等待查询的超时时间不要超过Lambda函数的最大执行时长限制

内容的提问来源于stack exchange,提问作者FJ777

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.23 19:36:01