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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 04:15:02