通过AWS Lambda运行Python脚本时S3生成文件为空问题求助
Athena查询生成空S3文件的问题排查与解决
问题背景
使用以下Python脚本通过AWS Lambda调用Athena执行查询,结果S3中生成的文件为空,怀疑查询语句存在问题:
import json import urllib.parse from datetime import datetime import boto3 import botocore.session as bc print('Loading function') client_s3=boto3.client('s3') def run_athena(BUCKET_NAME, PREFIX): client=boto3.client('athena') database = 'data_lake' query ="SELECT buid, property_id, checkin_date, count(*) number_of_direct_bookings FROM data_lake.vw_qs_bookings_v1 where source_logic in ('sources_phone', 'sources_email', 'sources_website','sources_front_desk') and cancelled_flag=0 and checkin_date= date_add('day', -1, current_date) group by 1,2,3"; s3_output = 's3://gainsight-file/gainsight/' response = client.start_query_execution(QueryString=query, QueryExecutionContext={'Database': database}, ResultConfiguration={'OutputLocation': s3_output}) today = datetime.today().strftime('%Y-%m-%d') response_s3 = client_s3.list_objects( Bucket=BUCKET_NAME, Prefix=PREFIX, ) name = response_s3["Contents"][0]["Key"] client_s3.copy_object(Bucket=BUCKET_NAME, CopySource=BUCKET_NAME+'/'+name, Key=PREFIX+"vw_qs_bookings_v1_"+today) client_s3.delete_object(Bucket=BUCKET_NAME, Key=name) def lambda_handler(event, context): print("Entered lambda_handler") run_athena('gainsight-file', 'gainsight/')
核心问题分析
1. 未等待Athena查询完成就操作S3
start_query_execution是异步调用,提交查询后立刻返回响应,但此时Athena可能还在处理查询,S3中的结果文件尚未生成或还未写入数据。这时候直接去S3复制文件,大概率拿到的是空文件,这是导致空文件的首要原因。
2. 查询语句的潜在问题
- 日期条件准确性:
date_add('day', -1, current_date)基于UTC时区返回日期,如果业务时区不是UTC,可能导致日期匹配错误,查不到数据。 - 数据存在性验证:视图
vw_qs_bookings_v1中是否存在符合source_logic、cancelled_flag=0且checkin_date为昨日的记录?可在Athena控制台执行简化查询验证:SELECT COUNT(*) FROM data_lake.vw_qs_bookings_v1 WHERE source_logic IN ('sources_phone', 'sources_email', 'sources_website','sources_front_desk') AND cancelled_flag=0 AND checkin_date= date_add('day', -1, current_date) - 字段类型匹配:确认
cancelled_flag是数值类型还是字符串类型。如果是字符串,条件应改为cancelled_flag='0',否则会因类型不匹配查不到数据。 - 视图权限与同步:确认Lambda执行角色有权限访问该视图,且视图底层数据是最新的,无延迟同步问题。
修复方案
第一步:添加Athena查询等待逻辑
修改run_athena函数,提交查询后循环检查状态,直到查询完成或失败:
def run_athena(BUCKET_NAME, PREFIX): client=boto3.client('athena') database = 'data_lake' query ="SELECT buid, property_id, checkin_date, count(*) number_of_direct_bookings FROM data_lake.vw_qs_bookings_v1 where source_logic in ('sources_phone', 'sources_email', 'sources_website','sources_front_desk') and cancelled_flag=0 and checkin_date= date_add('day', -1, current_date) group by 1,2,3"; s3_output = 's3://gainsight-file/gainsight/' # 提交查询 response = client.start_query_execution( QueryString=query, QueryExecutionContext={'Database': database}, ResultConfiguration={'OutputLocation': s3_output} ) query_execution_id = response['QueryExecutionId'] # 等待查询完成 while True: query_status = client.get_query_execution(QueryExecutionId=query_execution_id) status = query_status['QueryExecution']['Status']['State'] if status in ['SUCCEEDED', 'FAILED', 'CANCELLED']: break import time time.sleep(1) # 处理查询失败情况 if status != 'SUCCEEDED': error_msg = query_status['QueryExecution']['Status'].get('StateChangeReason', 'Unknown error') raise Exception(f"Athena query failed: {error_msg}") today = datetime.today().strftime('%Y-%m-%d') # 按查询ID过滤结果文件,避免拿到旧文件 response_s3 = client_s3.list_objects_v2( Bucket=BUCKET_NAME, Prefix=f"{PREFIX}{query_execution_id}" ) if 'Contents' not in response_s3 or len(response_s3['Contents']) == 0: raise Exception("No result file found in S3") name = response_s3["Contents"][0]["Key"] client_s3.copy_object( Bucket=BUCKET_NAME, CopySource=f"{BUCKET_NAME}/{name}", Key=f"{PREFIX}vw_qs_bookings_v1_{today}" ) client_s3.delete_object(Bucket=BUCKET_NAME, Key=name)
第二步:验证并调整查询语句
- 在Athena控制台手动执行原查询,确认是否返回数据。
- 若时区不匹配,调整日期条件。比如业务时区为东八区:
checkin_date = date_add('day', -1, date_trunc('day', current_timestamp AT TIME ZONE 'Asia/Shanghai')) - 确认
cancelled_flag字段类型,调整条件中的值类型。
其他注意事项
- Lambda执行角色需要包含
athena:GetQueryExecution权限,以及S3读写权限。 - 改用
list_objects_v2替代已过时的list_objects。 - 按查询执行ID过滤S3文件,避免误操作旧结果文件。
内容的提问来源于stack exchange,提问作者Caroline Silva
相关产品推荐
相关产品推荐

