AWS Lambda调用Amazon Athena查询未返回结果,求解决方案
使用AWS Lambda获取Amazon Athena查询结果的解决方法
你当前的Lambda代码仅提交了Athena查询请求,而Athena是异步执行查询的,因此start_query_execution返回的只是查询执行ID,并非实际查询结果。要拿到目标数据,需要等待查询完成后再获取结果。
一、核心修改逻辑
- 提交查询后,轮询查询状态,直到查询进入完成/失败/取消状态
- 根据结果规模,选择直接通过API获取数据或从指定S3路径读取输出文件
二、小结果集场景代码示例
如果查询结果行数较少(≤1000行),可直接用get_query_results API获取数据:
import boto3 import time # 配置参数 query = 'SELECT DISTINCT awsaccountid, digeststarttime FROM smxtech.cloudtrail_digest' DATABASE = 'xtech' output = 's3://xtech-destination/' athena_client = boto3.client('athena') def lambda_handler(event, context): # 启动Athena查询 response = athena_client.start_query_execution( QueryString=query, QueryExecutionContext={'Database': DATABASE}, ResultConfiguration={'OutputLocation': output} ) query_execution_id = response['QueryExecutionId'] # 轮询等待查询完成 while True: query_status = athena_client.get_query_execution(QueryExecutionId=query_execution_id)['QueryExecution']['Status']['State'] if query_status in ['SUCCEEDED', 'FAILED', 'CANCELLED']: break time.sleep(2) # 每2秒检查一次状态 if query_status == 'SUCCEEDED': # 获取并解析查询结果 result_data = athena_client.get_query_results(QueryExecutionId=query_execution_id) # 提取表头 columns = [col['Label'] for col in result_data['ResultSet']['ResultSetMetadata']['ColumnInfo']] # 提取数据行(跳过表头行) rows = [] for row in result_data['ResultSet']['Rows'][1:]: rows.append({columns[i]: row['Data'][i]['VarCharValue'] for i in range(len(columns))}) return {'status': 'success', 'data': rows} else: # 返回错误信息 error_msg = athena_client.get_query_execution(QueryExecutionId=query_execution_id)['QueryExecution']['Status']['StateChangeReason'] return {'status': query_status, 'error_message': error_msg}
三、大结果集场景代码示例
如果查询结果数据量较大(超过1000行或文件大于1MB),get_query_results会有返回限制,建议直接从S3输出路径读取CSV文件:
import boto3 import time import csv from io import StringIO # 初始化客户端 athena_client = boto3.client('athena') s3_client = boto3.client('s3') # 配置参数 query = 'SELECT DISTINCT awsaccountid, digeststarttime FROM smxtech.cloudtrail_digest' DATABASE = 'xtech' output = 's3://xtech-destination/' def lambda_handler(event, context): # 启动查询并等待完成(同小结果集的轮询逻辑) response = athena_client.start_query_execution( QueryString=query, QueryExecutionContext={'Database': DATABASE}, ResultConfiguration={'OutputLocation': output} ) query_execution_id = response['QueryExecutionId'] while True: query_status = athena_client.get_query_execution(QueryExecutionId=query_execution_id)['QueryExecution']['Status']['State'] if query_status in ['SUCCEEDED', 'FAILED', 'CANCELLED']: break time.sleep(2) if query_status == 'SUCCEEDED': # 获取S3输出文件路径 output_location = athena_client.get_query_execution(QueryExecutionId=query_execution_id)['QueryExecution']['ResultConfiguration']['OutputLocation'] # 解析S3桶和文件键 bucket = output_location.split('/')[2] file_key = '/'.join(output_location.split('/')[3:]) # 读取并解析S3中的CSV文件 s3_response = s3_client.get_object(Bucket=bucket, Key=file_key) csv_content = s3_response['Body'].read().decode('utf-8') reader = csv.DictReader(StringIO(csv_content)) rows = list(reader) return {'status': 'success', 'data': rows} else: error_msg = athena_client.get_query_execution(QueryExecutionId=query_execution_id)['QueryExecution']['Status']['StateChangeReason'] return {'status': query_status, 'error_message': error_msg}
四、必要的权限配置
确保Lambda执行角色拥有以下权限:
athena:StartQueryExecutionathena:GetQueryExecutionathena:GetQueryResultss3:GetObject(针对Athena输出桶xtech-destination)s3:PutObject(针对Athena输出桶)- (若数据源存储在S3)
s3:GetObject(针对CloudTrail数据所在的S3桶)
内容的提问来源于stack exchange,提问作者Marks
相关产品推荐
相关产品推荐

