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

使用PySpark创建50GB随机整数Parquet文件失败求助

AWS EMR集群YARN杀节点问题及50GB随机Parquet文件生成优化

问题详情

  • 尝试过不同规模的AWS EMR集群,均因YARN杀死所有节点导致任务失败;推测是内存需求过高,但使用10台m5.4xlarge实例(单节点64GB内存)的集群仍失败。
  • 当前使用的PySpark代码如下:
import os
import math
from pyspark.sql import SparkSession
from pyspark.sql.functions import rand

num_of_ints = int(size_in_mb * 1024 * 1024 / 4)
max_int = 2147483647

# 创建SparkSession
spark = SparkSession\
    .builder\
    .appName("GenerateRandomData") \
    .getOrCreate()

# 生成指定行数的DataFrame,包含0到max_int的随机整数列
df = spark.range(num_of_ints).withColumn("value", (rand(seed=42) * max_int).cast("integer"))

# 将DataFrame保存为Parquet文件
out_file = os.path.join(out_folder, 'random_list.parquet')
partitions = math.ceil(size_in_mb/10000) # 按10GB分片
df.repartition(partitions).write.mode("overwrite").parquet(out_file)

# 停止SparkSession
spark.stop()
  • 数据生成阶段仅由2个任务执行,与集群约140个核心的资源严重不匹配(Spark UI显示该现象)。
  • 需求:接受任何生成50GB随机整数Parquet文件的替代方法。

问题分析与解决方案

1. 核心问题:初始并行度不足导致单任务内存过载

spark.range(num_of_ints)默认使用spark.default.parallelism(默认值远低于集群核心数)生成初始分区,导致数据生成阶段只有少量任务。单任务需要处理数十GB的数据,内存占用远超YARN阈值,最终触发节点被杀死。后续的repartition是在数据生成完成后执行,无法解决生成阶段的内存压力。

2. 优化后的PySpark代码(直接提升并行度)

import os
import math
from pyspark.sql import SparkSession
from pyspark.sql.functions import rand

size_in_mb = 51200  # 对应50GB
num_of_ints = int(size_in_mb * 1024 * 1024 / 4)
max_int = 2147483647

# 初始化时指定并行度参数,匹配集群核心数
spark = SparkSession\
    .builder\
    .appName("GenerateRandomData")\
    .config("spark.default.parallelism", "140")  # 与集群140核心对齐
    .config("spark.sql.shuffle.partitions", "140")
    .getOrCreate()

# 生成range时直接指定分区数,让数据生成阶段就用满集群资源
df = spark.range(num_of_ints, numPartitions=140)\
    .withColumn("value", (rand(seed=42) * max_int).cast("integer"))

# 按10GB分片写入Parquet
out_file = os.path.join(out_folder, 'random_list.parquet')
partitions = math.ceil(size_in_mb / 10000)
df.repartition(partitions).write.mode("overwrite").parquet(out_file)

spark.stop()

3. EMR集群内存参数调优(针对m5.4xlarge)

单节点16核64GB内存,调整以下参数避免YARN内存阈值触发:

  • YARN配置(yarn-site.xml):
    • yarn.nodemanager.resource.memory-mb = 61440(预留3GB给系统进程)
    • yarn.scheduler.maximum-allocation-mb = 61440
  • Spark配置(spark-defaults.conf):
    • spark.executor.memory = 16g(每节点分配4个executor,16*4=64GB)
    • spark.executor.cores = 4(每个executor占用4核,16核节点刚好满负载)
    • spark.driver.memory = 8g

4. 替代方案:本地生成后上传至S3

如果Spark集群仍有内存问题,可采用本地批量生成后上传的方式(适合小集群或本地环境):

import boto3
import pandas as pd
import numpy as np

max_int = 2147483647
chunk_size = 10**8  # 每个chunk含1亿个int,约400MB
total_chunks = 125  # 125*400MB=50GB
s3_bucket = "your-bucket-name"
s3_prefix = "data/random_list/"

# 初始化S3客户端
s3 = boto3.client("s3")

for chunk_idx in range(total_chunks):
    # 生成随机整数DataFrame
    df = pd.DataFrame({
        "value": np.random.randint(0, max_int, size=chunk_size, dtype=np.int32)
    })
    # 临时保存Parquet文件
    temp_path = f"/tmp/random_chunk_{chunk_idx}.parquet"
    df.to_parquet(temp_path, compression="snappy")
    # 上传至S3
    s3.upload_file(temp_path, s3_bucket, f"{s3_prefix}chunk_{chunk_idx}.parquet")

内容的提问来源于stack exchange,提问作者Daniel

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.29 15:35:44