如何在Athena的get_query_execution完成后调用get_query_results并消除报错
解决Athena查询未完成时调用GetQueryResults的报错问题
问题根源是你用了固定的15秒等待逻辑,但Athena查询的执行时间受数据量、集群负载等因素影响,无法保证15秒内一定完成。当查询仍处于RUNNING状态时调用get_query_results,就会触发InvalidRequestException报错。
正确的做法是轮询查询状态,直到它进入终态(成功、失败或取消),再执行结果获取操作。以下是修改后的完整代码:
import json import boto3 import time def lambda_handler(event, context): client = boto3.client('athena') # 执行查询获取Athena元数据 query_start = client.start_query_execution( QueryString = """SELECT * FROM information_schema.columns WHERE table_catalog = 'table-catalog-name' AND table_schema = 'schema' ORDER BY table_name, column_name""", QueryExecutionContext = { 'Database': 'database' }, ResultConfiguration = { 'OutputLocation': 's3://bucket-name/' } ) query_id = query_start['QueryExecutionId'] query_status = None # 轮询查询状态,直到进入终态 while query_status not in ['SUCCEEDED', 'FAILED', 'CANCELLED']: response = client.get_query_execution(QueryExecutionId=query_id) query_status = response['QueryExecution']['Status']['State'] if query_status == 'FAILED': raise Exception(f"Athena query failed: {response['QueryExecution']['Status']['StateChangeReason']}") elif query_status == 'CANCELLED': raise Exception("Athena query was cancelled") # 每次轮询间隔3秒,平衡响应速度与API调用频率 time.sleep(3) # 查询成功后获取结果 results = client.get_query_results(QueryExecutionId=query_id) for row in results['ResultSet']['Rows']: print(row)
关键改动说明:
- 替换固定等待:用循环轮询替代
time.sleep(15),确保只在查询真正完成后才执行结果获取 - 异常处理:捕获查询失败或取消的情况,抛出包含具体原因的异常,便于问题排查
- 合理轮询间隔:设置3秒的等待间隔,避免过于频繁调用Athena API导致限流
内容的提问来源于stack exchange,提问作者Forgottenluv
相关产品推荐
相关产品推荐

