Spark与PostgreSQL交互触发OOM问题求助
解决方案:Spark Pivot导致PostgreSQL OOM问题处理
一、优化Spark Pivot操作,降低内存压力
- 限制Pivot展开的列数:如果
attr_id的不同取值过多,全量Pivot会生成海量列,直接触发内存溢出。指定业务需要的列,避免无差别展开:// 只保留需要的attr_id值,减少生成的列数 List<String> targetAttrIds = Arrays.asList("attr_1", "attr_2", "attr_3"); Dataset<Row> iter4 = iter_3 .groupBy( col("super_parent"), col("parent_id").cast("string").as("parent_id") ) .pivot("attr_id", targetAttrIds) // 指定要Pivot的列 .agg(first(col("value")).cast("string").as("value")) .as("iter_4"); - 调整Spark Shuffle配置:默认Shuffle分区数(20)会导致单分区数据量过大,调大分区数并优化Shuffle内存占比:
// SparkSession初始化时设置 spark.conf.set("spark.sql.shuffle.partitions", "200") spark.conf.set("spark.shuffle.memoryFraction", "0.4") // 分配40%的Executor内存给Shuffle - 提前过滤冗余数据:在Pivot前过滤空值和无用
attr_id,减少分组后的数据量:Dataset<Row> filteredIter3 = iter_3 .filter(col("attr_id").isNotNull()) .filter(col("value").isNotNull()); // 基于过滤后的数据集执行Pivot Dataset<Row> iter4 = filteredIter3.groupBy(...).pivot(...).agg(...);
二、优化PostgreSQL端资源与查询
- 调整PostgreSQL内存参数:修改
postgresql.conf中的关键配置,避免查询时内存不足:shared_buffers = 1GB # 设为系统内存的1/4(假设服务器总内存4GB) work_mem = 64MB # 单个排序/分组操作的内存,减少磁盘落盘 maintenance_work_mem = 256MB # 维护操作(如建索引)的内存 temp_buffers = 128MB # 临时表使用的内存 - 给临时表添加索引:针对OOM前的两个长查询,给临时表的过滤字段加索引,加速查询并减少内存占用:
CREATE INDEX idx_temp_parent ON temp_2166008c_4a87_4f77_a267_077adde08453_v2 (parent_id); CREATE INDEX idx_temp_object ON temp_2166008c_4a87_4f77_a267_077adde08453_v2 (object_id); - 避免并发执行大查询:两个长查询同时扫描临时表会加倍消耗内存,改为串行执行或合并查询减少重复扫描:
-- 合并查询,一次扫描获取所有所需字段 SELECT parent_id, root_id, object_id, attr_id FROM temp_2166008c_4a87_4f77_a267_077adde08453_v2 WHERE parent_id IS NOT NULL OR object_id IS NOT NULL;
三、优化数据处理流程
- 分批次处理数据:将3GB数据按
super_parent或parent_id分片,分批次处理后合并结果,避免一次性加载全量数据:// 获取所有distinct的super_parent值 List<String> superParents = iter_3.select("super_parent").distinct().as(Encoders.STRING()).collectAsList(); // 分批次处理 List<Dataset<Row>> batchResults = new ArrayList<>(); for (String sp : superParents) { Dataset<Row> batch = iter_3.filter(col("super_parent").equalTo(sp)) .groupBy(...).pivot(...).agg(...); batchResults.add(batch); } // 合并所有批次结果 Dataset<Row> finalIter4 = spark.createDataFrame(spark.emptyDataFrame().rdd(), batchResults.get(0).schema()); for (Dataset<Row> batch : batchResults) { finalIter4 = finalIter4.union(batch); } - 中间结果落地到文件系统:将Spark中间结果保存为Parquet格式到本地/HDFS,而非依赖PostgreSQL临时表,减少数据库压力:
// 保存中间结果到Parquet iter_3.write().mode(SaveMode.Overwrite).parquet("/path/to/iter3_data"); // 从Parquet加载数据继续处理 Dataset<Row> iter3FromParquet = spark.read().parquet("/path/to/iter3_data"); - 排查数据倾斜:检查
super_parent或parent_id是否存在热点值,若存在,给热点值添加随机后缀拆分分组:// 给热点super_parent添加随机后缀,拆分分组 Dataset<Row> skewedFixedIter3 = iter_3.withColumn("super_parent_skew", when(col("super_parent").equalTo("hot_value"), concat(col("super_parent"), lit("_"), floor(rand() * 10)) // 拆成10个小分组 ).otherwise(col("super_parent")) ); // 基于拆分后的字段分组Pivot Dataset<Row> iter4 = skewedFixedIter3.groupBy("super_parent_skew", "parent_id") .pivot(...).agg(...);
内容的提问来源于stack exchange,提问作者Maxim
相关产品推荐
相关产品推荐

