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

数据量超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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.06 16:00:54