使用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
相关产品推荐
相关产品推荐

