如何优化查询大型DataFrame并保存至S3的内存占用?
高效读写方案解决大DataFrame内存问题
方案1:让Athena直接将查询结果写入S3(彻底规避本地内存加载)
不需要将全量数据加载到本地DataFrame,直接通过Athena的CTAS(Create Table As Select)语句把查询结果直接写入S3,全程在AWS端处理,本地仅需触发查询。
代码示例:
import awswrangler as wr import boto3 region = "" access_key = "" secret_key = "" database = "" s3_output_path = "s3://your-bucket/path/to/output/" # 初始化会话 boto3.setup_default_session(region_name=region, aws_access_key_id=access_key, aws_secret_access_key=secret_key) # 用CTAS语句直接将查询结果写入S3的CSV文件 ctas_sql = f""" CREATE TABLE athena_output WITH ( format = 'CSV', external_location = '{s3_output_path}', skip_header = 'true' ) AS SELECT * FROM tbl """ # 执行CTAS查询 wr.athena.run_query(sql=ctas_sql, database=database, wait=True) # 如果需要提取unique值,直接用Athena查询去重结果并写入S3 distinct_sql = f""" CREATE TABLE distinct_values WITH ( format = 'CSV', external_location = 's3://your-bucket/path/to/distinct-values/' ) AS SELECT DISTINCT value FROM tbl """ wr.athena.run_query(sql=distinct_sql, database=database, wait=True)
方案2:分块读取+分块写入S3
如果必须在本地处理数据,使用chunksize参数分块读取Athena结果,避免一次性加载全量数据到内存,同时分块追加写入S3。
代码示例:
from io import StringIO import pandas as pd import boto3 import awswrangler as wr region = "" access_key = "" secret_key = "" database = "" s3_bucket = "your-bucket" s3_file_path = "path/to/folder/values.csv" boto3.setup_default_session(region_name=region, aws_access_key_id=access_key, aws_secret_access_key=secret_key) # 分块读取数据,chunksize根据本地内存情况调整 chunks = wr.athena.read_sql_query(sql="SELECT * FROM tbl", database=database, chunksize=100000) # 累积去重值 unique_values = set() s3_client = boto3.client('s3') # 处理第一个chunk,初始化S3文件(覆盖模式) first_chunk = next(chunks) unique_values.update(first_chunk['value'].unique()) csv_buffer = StringIO() first_chunk.to_csv(csv_buffer, index=False) s3_client.put_object(Bucket=s3_bucket, Key=s3_file_path, Body=csv_buffer.getvalue()) # 处理后续chunk,追加内容到S3文件 for chunk in chunks: unique_values.update(chunk['value'].unique()) csv_buffer = StringIO() # 跳过表头,避免重复写入 chunk.to_csv(csv_buffer, index=False, header=False) # 读取现有文件内容并追加新数据 existing_content = s3_client.get_object(Bucket=s3_bucket, Key=s3_file_path)['Body'].read().decode('utf-8') new_content = existing_content + csv_buffer.getvalue() s3_client.put_object(Bucket=s3_bucket, Key=s3_file_path, Body=new_content) # 保存去重值到S3 unique_df = pd.Series(list(unique_values), name='value').to_frame() wr.s3.to_csv(df=unique_df, path=f"s3://{s3_bucket}/path/to/unique-values.csv")
方案3:使用Parquet格式替代CSV
Parquet是列存储格式,压缩率远高于CSV,内存占用更低,awswrangler支持直接将Athena查询结果写入Parquet文件到S3,无需本地加载全量数据。
代码示例:
import awswrangler as wr import boto3 region = "" access_key = "" secret_key = "" database = "" s3_parquet_path = "s3://your-bucket/path/to/parquet-output/" boto3.setup_default_session(region_name=region, aws_access_key_id=access_key, aws_secret_access_key=secret_key) # 直接将Athena查询结果以Parquet格式写入S3,自动分块处理 wr.athena.read_sql_query( sql="SELECT * FROM tbl", database=database, ctas_approach=True, s3_output=s3_parquet_path, dataset=True, format="parquet" ) # 提取去重值并写入Parquet wr.athena.read_sql_query( sql="SELECT DISTINCT value FROM tbl", database=database, ctas_approach=True, s3_output="s3://your-bucket/path/to/distinct-parquet/", dataset=True, format="parquet" )
内容的提问来源于stack exchange,提问作者codenoodles
相关产品推荐
相关产品推荐

