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

通过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)

第二步:验证并调整查询语句

  1. 在Athena控制台手动执行原查询,确认是否返回数据。
  2. 若时区不匹配,调整日期条件。比如业务时区为东八区:
    checkin_date = date_add('day', -1, date_trunc('day', current_timestamp AT TIME ZONE 'Asia/Shanghai'))
    
  3. 确认cancelled_flag字段类型,调整条件中的值类型。

其他注意事项

  • Lambda执行角色需要包含athena:GetQueryExecution权限,以及S3读写权限。
  • 改用list_objects_v2替代已过时的list_objects。
  • 按查询执行ID过滤S3文件,避免误操作旧结果文件。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.20 20:24:32