数据量超Glue内存时,AWS Glue合并小文件为固定大小任务失败求助
解决AWS Glue合并大分区小Parquet文件并转String类型的内存问题
一、优化读取策略:避免全量加载
大分区直接全量读取会触发元数据过载,导致Executor内存溢出。改用文件批量分组读取,每次仅处理固定大小的文件批次,不一次性加载全部分区元数据。
- 用
boto3遍历S3分区文件,按预估大小(如1GB/批次)分组,逐批次处理:
import boto3 from awsglue.context import GlueContext from pyspark.context import SparkContext sc = SparkContext() glueContext = GlueContext(sc) spark = glueContext.spark_session s3_client = boto3.client('s3') bucket = "your-bucket" prefix = "path/to/partition/date=20240101/" # 按指定大小批量生成文件组 def get_file_batches(bucket, prefix, batch_size_gb=1): batch_size_bytes = batch_size_gb * 1024**3 current_batch = [] current_size = 0 paginator = s3_client.get_paginator('list_objects_v2') for page in paginator.paginate(Bucket=bucket, Prefix=prefix): for obj in page.get('Contents', []): file_size = obj['Size'] if current_size + file_size > batch_size_bytes and current_batch: yield current_batch current_batch = [] current_size = 0 current_batch.append(f"s3://{bucket}/{obj['Key']}") current_size += file_size if current_batch: yield current_batch # 逐批次处理:读取→转String→按固定大小写入 for file_batch in get_file_batches(bucket, prefix, batch_size_gb=1): # 读取批次文件 df = spark.read.parquet(*file_batch) # 转换所有列为string类型 df_string = df.select([df[col].cast("string").alias(col) for col in df.columns]) # 估算目标文件数(按压缩后1GB/文件计算,Parquet压缩率约1:5) total_raw_size = sum(s3_client.head_object(Bucket=bucket, Key=f.split('//')[1])['ContentLength'] for f in file_batch) target_file_count = max(1, int(total_raw_size / (1024**3 * 5))) # 用coalesce合并分区(无shuffle,低内存开销) df_string.coalesce(target_file_count).write.mode("append").parquet("s3://your-target-bucket/merged/date=20240101/")
二、调整Glue任务内存与Spark配置
不要盲目堆叠DPU,重点优化Spark内存参数,降低容器崩溃概率:
- 在Glue任务的Job parameters中添加以下配置:
--conf spark.executor.memory=16g(G.2x实例单Executor配16G内存,对应2 DPU)--conf spark.executor.cores=4(匹配G.2x实例的核数,避免资源浪费)--conf spark.driver.memory=8g--conf spark.sql.shuffle.partitions=64(减少shuffle分区数,降低内存占用)--conf spark.sql.parquet.enableVectorizedReader=false(关闭矢量化读取,避免大分区元数据加载的内存峰值)--conf spark.hadoop.fs.s3a.connection.timeout=300000(延长S3连接超时,减少网络导致的Executor断开)
- Worker类型选G.2x,内存/核比更高,适合大内存任务;40GB分区配20-30 DPU即可,过多DPU会增加集群协调开销。
三、基于Glue Catalog的分区优化读取
如果已创建Glue Catalog表,利用分区谓词下推,仅加载目标分区数据,避免全表扫描:
from awsglue.job import Job from awsglue.transforms import Map job = Job(glueContext) job.init(args['JOB_NAME'], args) # 仅读取指定分区,触发谓词下推 datasource = glueContext.create_dynamic_frame.from_catalog( database="your-db", table_name="your-table", push_down_predicate="date='20240101'", additional_options={"enableCatalogPartitionPredicatePushdown": "true"} ) # 转换所有列为string类型 def cast_to_string(rec): for key in rec: rec[key] = str(rec[key]) return rec datasource_string = Map.apply(frame=datasource, f=cast_to_string) # 按固定大小控制输出文件数 estimated_size_gb = 40 compression_ratio = 5 target_file_size_gb = 1 num_partitions = int(estimated_size_gb / compression_ratio / target_file_size_gb) # 用coalesce合并分区,避免shuffle datasource_string = datasource_string.coalesce(max(1, num_partitions)) glueContext.write_dynamic_frame.from_options( frame=datasource_string, connection_type="s3", connection_options={"path": "s3://your-target-bucket/merged/date=20240101/"}, format="parquet" ) job.commit()
四、避坑要点
- 优先用
coalesce而非repartition:repartition会触发全量shuffle,内存开销极大;coalesce仅合并现有分区,无shuffle操作。 - 不在Driver端处理大量元数据:用boto3分页获取文件列表,避免一次性加载数万条文件元数据撑爆Driver内存。
- 关闭自动广播:添加
--conf spark.sql.autoBroadcastJoinThreshold=-1,防止大分区数据被广播导致内存溢出。
内容的提问来源于stack exchange,提问作者codeduck
相关产品推荐
相关产品推荐

