Amazon Athena查询返回无效数据问题排查求助
问题描述
创建外部表athena_test_table并插入3条数据后,直接在Athena控制台执行SELECT查询结果正常;但通过C#应用(基于AWSSDK.Core 3.7.106.2和AWSSDK.Athena 3.7.104.9)查询,或执行SHOW TABLES IN TESTDB语句时,目标表中会新增空行或垃圾数据,且重复请求会导致表数据量持续增长。
创建表语句
CREATE EXTERNAL TABLE IF NOT EXISTS athena_test_table( id int, c1 string, c2 string, c3 string) LOCATION 's3://s3-path/athena/'
插入数据语句
insert into "athena_test_table" values (1,'val1a','val1b','val1c'), (2,'val2a','val2b','val2c'), (3,'val3a','val3b','val3c')
C#查询代码
private async Task<ResultSet> SendAndResponseToAthena(string query, AthenaParameters data, ProxyData proxyData) { try { Uri.TryCreate(proxyData.Url, UriKind.Absolute, out var proxyUrl); var clientConfig = new AmazonAthenaConfig { ProxyHost = proxyUrl?.Host, ProxyPort = proxyUrl?.Port ?? 0, ProxyCredentials = new NetworkCredential(proxyData.Username, proxyData.Password), RegionEndpoint = RegionEndpoint.GetBySystemName(data.Region) }; var client = new AmazonAthenaClient(data.UserName, data.Password, clientConfig); var executionId = await SubmitAthenaQuery(query, client, data); await WaitForQueryToComplete(client, executionId); var resultSet = await ProcessResult(client, executionId); return resultSet; } catch (Exception ex) { LogException(ex.Message, ex, this); throw new Exception(ex.ToString()); } } private async Task<string> SubmitAthenaQuery(string query, AmazonAthenaClient client, AthenaParameters data) { var executionRequest = new StartQueryExecutionRequest { QueryString = query, ResultConfiguration = new ResultConfiguration { OutputLocation = data.S3Location }, QueryExecutionContext = new QueryExecutionContext { Database = data.DatabaseName } }; try { var response = await client.StartQueryExecutionAsync(executionRequest); var executionId = response.QueryExecutionId; LogInfo($"Submit Athena Query: {query} with Execution Id: {executionId}.", this); return executionId; } catch (Exception ex) { LogException(ex.Message, ex, this); throw new Exception(ex.ToString()); } } private async Task WaitForQueryToComplete(AmazonAthenaClient client, string executionId) { var getQueryExecutionRequest = new GetQueryExecutionRequest { QueryExecutionId = executionId }; var isQueryStillRunning = true; while (isQueryStillRunning) { var getQueryExecutionResponse = await client.GetQueryExecutionAsync(getQueryExecutionRequest); var queryState = getQueryExecutionResponse.QueryExecution.Status.State.ToString(); var changeReason = getQueryExecutionResponse.QueryExecution.Status.StateChangeReason; if (queryState.Equals(QueryExecutionState.FAILED.ToString())) { throw new Exception($"The Amazon Athena query failed to run with error message: {changeReason}"); } if (queryState.Equals(QueryExecutionState.CANCELLED.ToString())) { throw new Exception($"The Amazon Athena query was cancelled with change reason: {changeReason}"); } if (queryState.Equals(QueryExecutionState.SUCCEEDED.ToString())) { isQueryStillRunning = false; } else { // sleep an amount of time (500ms) before retrying again await Task.Delay(500); } } } private async Task<ResultSet> ProcessResult(AmazonAthenaClient client, string executionId) { var queryResultsRequest = new GetQueryResultsRequest { QueryExecutionId = executionId }; try { LogInfo($"Retrieve the results for Execution Id: {executionId}.", this); var results = await client.GetQueryResultsAsync(queryResultsRequest); var pageResult = results; while (pageResult.NextToken != null) // process pagination (>1000 entries) { var request = new GetQueryResultsRequest { QueryExecutionId = executionId, NextToken = pageResult.NextToken }; pageResult = await client.GetQueryResultsAsync(request); results.ResultSet.Rows.AddRange(pageResult.ResultSet.Rows); } return results.ResultSet; } catch (Exception ex) { LogException(ex.Message, ex, this); throw new Exception(ex.ToString()); } }
原因分析
- 数据目录与查询结果目录重叠:外部表的
LOCATION指向s3://s3-path/athena/,而C#代码中设置的查询结果输出目录(ResultConfiguration.OutputLocation)如果是该路径或其子路径,每次查询生成的结果文件会被Athena当作表的数据文件加载,导致空行/垃圾数据,且重复查询会持续增加文件数量。 - 未指定数据格式与SerDe:创建表时未明确
ROW FORMAT和存储格式,Athena默认使用LazySimpleSerDe处理文本数据,一旦查询结果文件(如CSV格式)混入数据目录,解析时会出现不匹配,生成垃圾数据。 - 元数据未及时刷新:混入的垃圾文件未被清理时,Athena会持续扫描并加载这些文件,导致数据量不断增长。
解决办法
1. 分离数据与查询结果目录
- 修改外部表的
LOCATION为专属数据路径,例如:s3://s3-path/athena/test-table-data/ - 在C#代码中,将
ResultConfiguration.OutputLocation设置为独立的查询结果路径,例如:s3://s3-path/athena/query-results/,确保两个路径完全无重叠。
2. 明确表的数据格式与SerDe
修改创建表语句,指定可靠的格式(如Parquet,避免文本格式的解析问题):
CREATE EXTERNAL TABLE IF NOT EXISTS athena_test_table( id int, c1 string, c2 string, c3 string ) ROW FORMAT SERDE 'org.apache.hadoop.hive.ql.io.parquet.serde.ParquetHiveSerDe' STORED AS INPUTFORMAT 'org.apache.hadoop.hive.ql.io.parquet.MapredParquetInputFormat' OUTPUTFORMAT 'org.apache.hadoop.hive.ql.io.parquet.MapredParquetOutputFormat' LOCATION 's3://s3-path/athena/test-table-data/'
3. 清理垃圾数据并刷新元数据
- 登录S3控制台,进入表的原
LOCATION路径,删除所有非插入操作生成的文件(如查询结果的CSV/TXT文件)。 - 执行
MSCK REPAIR TABLE athena_test_table;刷新表的元数据,确保Athena不再加载无效文件。
4. 优化C#代码配置
- 确保
ResultConfiguration.OutputLocation始终指向独立目录,避免配置错误; - 可为查询结果目录设置S3生命周期规则,自动清理旧的查询结果文件,减少存储占用。
内容的提问来源于stack exchange,提问作者leon22
相关产品推荐
相关产品推荐

