【问题标题】:(AWS) Athena: Query Results seem too short(AWS) Athena:查询结果似乎太短
【发布时间】:2018-01-18 19:26:16
【问题描述】:

我的 Athena 查询的结果似乎太短了。试图弄清楚为什么?

设置:

胶水目录(118.6 Gig 大小)。 数据:以 CSV 和 JSON 格式存储在 S3 中。 Athena Query:当我查询整个表的数据时,每个 Query 只得到 40K 的结果,一个月的数据平均应该有 1.21 亿条记录。

Athena Cap 是否查询结果数据?这是服务限制吗(文档不建议是这种情况)。

【问题讨论】:

    标签: amazon-web-services amazon-s3 amazon-athena aws-glue


    【解决方案1】:

    好像有1000个限制。 您应该使用NextToken 来迭代结果。

    引用GetQueryResults 文档

    MaxResults 在此返回的最大结果数(行) 请求。

    类型:整数

    有效范围:最小值0,最大值1000。

    必填:否

    【讨论】:

    • 从 SDK 的角度来看,这是有道理的。但是在 Console 中,查询限制远远超过 1000。对于 select * 有 121,000,000 条记录的表返回 40,000 条记录。
    • 好的。不知道您使用的是 CLI SDK。但是如果 CLI 返回 40000,方法是一样的。在NextToken 的帮助下迭代您的结果。
    【解决方案2】:

    因此,一次获得 1000 个结果显然无法扩展。值得庆幸的是,有一个简单的解决方法。 (或者也许这就是它应该一直这样做的方式。)

    当您运行 Athena 查询时,您应该得到一个 QueryExecutionId。此 Id 对应于您将在 S3 中找到的输出文件。

    这是我写的一个sn-p:

    s3 = boto3.resource("s3")
    athena = boto3.client("athena")
    response: Dict = athena.start_query_execution(QueryString=query, WorkGroup="<your_work_group>")
    execution_id: str = response["QueryExecutionId"]
    print(execution_id)
    
    # Wait until the query is finished
    while True:
        try:
            athena.get_query_results(QueryExecutionId=execution_id)
            break
        except botocore.exceptions.ClientError as e:
            time.sleep(5)
    
    local_filename: str = "temp/athena_query_result_temp.csv"
    s3.Bucket("athena-query-output").download_file(execution_id + ".csv", local_filename)
    return pd.read_csv(local_filename)
    

    确保对应的WorkGroup 设置了“查询结果位置”,例如"s3://athena-query-output/"

    另请参阅此主题的类似答案:How to Create Dataframe from AWS Athena using Boto3 get_query_results method

    【讨论】:

      【解决方案3】:

      另一种选择是分页和计数方法: 不知道是否有更好的方法,比如 select count(*) from table like...

      这是可以使用的完整示例代码。使用 python boto3 athena api 我使用 paginator 并将结果转换为 dict 列表,并与结果一起返回计数。

      以下是两种方法 第一个将分页 第二个会将分页结果转换为字典列表并计算计数。

      注意:在这种情况下,不需要转换为dict 列表。如果你不想这样..在代码中你可以修改为只有计数

      def get_athena_results_paginator(params, athena_client):
          """
      
          :param params:
          :param athena_client:
          :return:
          """
          query_id = athena_client.start_query_execution(
              QueryString=params['query'],
              QueryExecutionContext={
                  'Database': params['database']
              }
              # ,
              # ResultConfiguration={
              #     'OutputLocation': 's3://' + params['bucket'] + '/' + params['path']
              # }
              , WorkGroup=params['workgroup']
      
          )['QueryExecutionId']
          query_status = None
          while query_status == 'QUEUED' or query_status == 'RUNNING' or query_status is None:
              query_status = athena_client.get_query_execution(QueryExecutionId=query_id)['QueryExecution']['Status']['State']
              if query_status == 'FAILED' or query_status == 'CANCELLED':
                  raise Exception('Athena query with the string "{}" failed or was cancelled'.format(params.get('query')))
              time.sleep(10)
          results_paginator = athena_client.get_paginator('get_query_results')
          results_iter = results_paginator.paginate(
              QueryExecutionId=query_id,
              PaginationConfig={
                  'PageSize': 1000
              }
          )
          count, results = result_to_list_of_dict(results_iter)
          return results, count
      
      
      def result_to_list_of_dict(results_iter):
          """
      
          :param results_iter:
          :return:
          """
          results = []
          column_names = None
          count = 0
          for results_page in results_iter:
              print(len(list(results_iter)))
              for row in results_page['ResultSet']['Rows']:
                  count = count + 1
                  column_values = [col.get('VarCharValue', None) for col in row['Data']]
                  if not column_names:
                      column_names = column_values
                  else:
                      results.append(dict(zip(column_names, column_values)))
          return count, results
      
      

      【讨论】:

        猜你喜欢
        • 2018-12-17
        • 2019-11-07
        • 2021-10-18
        • 2017-06-17
        • 2020-09-10
        • 2020-07-14
        • 2020-01-13
        • 2019-05-21
        • 2020-09-12
        相关资源
        最近更新 更多