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

如何通过SNS发送Step Functions中两个Athena查询的结果行数通知?

实现方案

一、调整Athena查询(两种可选方式)

方式1:直接返回行数(推荐,高效省资源)

将原查询改为嵌套结构,直接输出结果行数,无需额外处理完整结果集。例如:
原查询:

SELECT col1, col2 FROM my_table WHERE date = CURRENT_DATE

修改为:

SELECT COUNT(*) AS row_count FROM (
  SELECT col1, col2 FROM my_table WHERE date = CURRENT_DATE
) AS subquery

查询结果将直接返回一行一列的行数数据,后续提取更简便。

方式2:保留原查询结果,后续统计行数

如果需要保留完整查询结果,原查询保持不变,后续通过API或Lambda统计结果文件的行数。

二、扩展Step Functions流程

将状态机流程调整为:EventBridge触发 → 执行Athena查询1 → 获取查询1行数 → 执行Athena查询2 → 获取查询2行数 → 组装消息发送SNS

1. 确保Athena查询步骤输出QueryExecutionId

每个Athena查询任务调用StartQueryExecution后,会返回QueryExecutionId,需将其保存到状态机上下文中,作为后续获取结果的唯一标识。

2. 添加获取行数的步骤

对应方式1:调用Athena GetQueryResults API

在Step Functions中新增一个Task类型步骤,调用Athena的GetQueryResults接口,输入参数为上一步的QueryExecutionId:

{
  "QueryExecutionId.$": "$.QueryExecutionId"
}

通过ResultSelector直接提取行数(结果集中第二行第一列的值,第一行为列名):

{
  "row_count.$": "$.ResultSet.Rows[1].Data[0].VarCharValue"
}

对应方式2:用Lambda函数统计行数

编写Python Lambda函数,接收QueryExecutionId后完成以下操作:

import boto3

athena = boto3.client('athena')
s3 = boto3.client('s3')

def lambda_handler(event, context):
    query_execution_id = event['QueryExecutionId']
    # 获取查询结果的S3存储路径
    exec_response = athena.get_query_execution(QueryExecutionId=query_execution_id)
    s3_path = exec_response['QueryExecution']['ResultConfiguration']['OutputLocation']
    bucket, key = s3_path.replace('s3://', '').split('/', 1)
    
    # 读取S3文件并统计行数(减去表头行)
    obj = s3.get_object(Bucket=bucket, Key=key)
    row_count = sum(1 for _ in obj['Body'].iter_lines()) - 1
    return {'row_count': row_count}

在Step Functions中新增调用该Lambda的Task步骤,传入QueryExecutionId并保存返回的row_count。

3. 组装消息并发送至SNS

新增Task步骤调用SNS的Publish API,构造包含两行数的邮件内容:

{
  "TopicArn": "arn:aws:sns:your-region:your-account-id:your-topic-arn",
  "Message.$": "States.Format('每日Athena查询结果:\\n查询1行数:{} \\n查询2行数:{}', $.query1_row_count, $.query2_row_count)",
  "Subject": "每日查询结果行数通知"
}

三、权限配置

为Step Functions执行角色添加以下权限:

  • athena:GetQueryExecution、athena:GetQueryResults(方式1所需)
  • lambda:InvokeFunction(方式2所需)
  • sns:Publish
  • s3:GetObject(Lambda读取S3结果文件所需)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.04 07:50:22