如何在Go SDK中正确等待AWS Athena查询执行并处理失败?
问题描述
我有一段可运行的Go代码,通过轮询GetQueryResults的错误返回值来等待Athena查询完成,代码如下:
func GetQueryResults(client *athena.Client, QueryID *string) []types.Row { params := &athena.GetQueryResultsInput{ QueryExecutionId: QueryID, } data, err := client.GetQueryResults(context.TODO(), params) for err != nil { println(err.Error()) time.Sleep(time.Second) data, err = client.GetQueryResults(context.TODO(), params) } return data.ResultSet.Rows }
问题在于当查询失败时,我完全无法跳出循环。
例如在Python中,我可以这样实现:
while athena.get_query_execution(QueryExecutionId=execution_id)["QueryExecution"][ "Status" ]["State"] in ["RUNNING", "QUEUED"]: sleep(2)
我可以在Go的循环中通过strings.Contains(err.Error(),"FAILED")来检查,但希望有更优雅的方式。我尝试寻找Go中的等效方法但未成功,请问Go SDK是否有返回执行状态的函数?除了err != nil之外,是否有更好的错误检查方式?
解决方案
1. 使用GetQueryExecution API主动查询状态
Go SDK和Python SDK一致,提供了GetQueryExecution方法直接获取查询的执行状态,这是最可靠、优雅的方式,无需依赖GetQueryResults的错误返回值判断。
你可以轮询该API的返回结果,根据QueryExecution.Status.State的枚举值处理不同状态:
types.QueryExecutionStateRunning:查询仍在运行types.QueryExecutionStateQueued:查询在排队types.QueryExecutionStateSucceeded:查询成功,此时可调用GetQueryResults获取结果types.QueryExecutionStateFailed/types.QueryExecutionStateCancelled:查询失败/取消,直接终止循环并返回错误
2. 用错误类型断言替代字符串匹配
如果仍需要通过GetQueryResults的错误做判断,不要用字符串匹配,而是利用AWS SDK的错误类型断言。AWS Go SDK的错误通常实现了smithy.APIError接口,可以提取精准的错误码:
import ( "errors" "github.com/aws/aws-sdk-go-v2/service/athena" "github.com/aws/smithy-go" ) // 在循环中处理错误 if err != nil { var smithyErr smithy.APIError if errors.As(err, &smithyErr) { switch smithyErr.ErrorCode() { case "QueryExecutionStillRunningException": // 查询仍在运行,继续轮询 time.Sleep(time.Second) continue case "QueryExecutionNotFoundException", "InvalidQueryExecutionIdException": // 无效查询ID,直接返回错误 return nil, err default: // 其他错误(如查询失败),终止循环 return nil, err } } // 非AWS服务错误,直接返回 return nil, err }
3. 改进后的完整代码示例
结合状态查询和结果获取的健壮实现:
import ( "context" "errors" "time" "github.com/aws/aws-sdk-go-v2/aws" "github.com/aws/aws-sdk-go-v2/service/athena" "github.com/aws/aws-sdk-go-v2/service/athena/types" ) func WaitForQueryAndGetResults(client *athena.Client, queryID *string) ([]types.Row, error) { ctx := context.TODO() // 轮询查询状态 for { execResp, err := client.GetQueryExecution(ctx, &athena.GetQueryExecutionInput{ QueryExecutionId: queryID, }) if err != nil { return nil, err } state := execResp.QueryExecution.Status.State switch state { case types.QueryExecutionStateSucceeded: // 查询成功,获取结果 resultResp, err := client.GetQueryResults(ctx, &athena.GetQueryResultsInput{ QueryExecutionId: queryID, }) if err != nil { return nil, err } return resultResp.ResultSet.Rows, nil case types.QueryExecutionStateFailed: return nil, errors.New("query failed: " + aws.ToString(execResp.QueryExecution.Status.StateChangeReason)) case types.QueryExecutionStateCancelled: return nil, errors.New("query was cancelled") case types.QueryExecutionStateRunning, types.QueryExecutionStateQueued: // 等待后继续轮询 time.Sleep(2 * time.Second) continue default: return nil, errors.New("unknown query state: " + string(state)) } } }
方案优势
- 主动查询状态,逻辑清晰,不会因
GetQueryResults的错误语义变化而失效 - 精准区分查询的各种状态,处理逻辑更严谨
- 避免字符串匹配的脆弱性,错误处理更可靠
内容的提问来源于stack exchange,提问作者Mohamed Yasser
相关产品推荐
相关产品推荐

