Spark Join时按条件裁剪大字符串以避免内存溢出(OOM)问题
解决Spark Join大字符串导致的内存溢出(OOM)问题
问题场景
有两个DataFrame:
- 第一个DataFrame(约100万条记录)存储位置信息:
id start end 1 4 8 2 2 6 2 5 7
- 第二个DataFrame(仅10条记录)存储id与超大字符串(单条约100MB):
id string 1 my beautiful data 2 lorem ipsum
直接执行Join后裁剪字符串会触发OOM,期望在处理过程中就完成字符串裁剪,控制内存占用,最终得到结果:
id start end seq 1 | 4 | 8 | beaut 2 | 2 | 6 | orem i 2 | 5 | 7 | m i
当前实现方式为先Join再截取,触发OOM:
df1 = df1.join(df2, "id") df1.withColumn("substring", df1['string'].substr(df1.start, df1.end))).show()
尝试过broadcast、repartition优化,但中间结果体积仍超出内存限制,单独处理两个DataFrame无异常。
原执行计划
执行exonsDF = exons_raw.join(dnaDF, "seq_region_id").explain("cost")得到的计划:
== Optimized Logical Plan == Project [seq_region_id#63L, exon_id#62L, seq_region_start#64L, seq_region_end#65L, seq_region_strand#66, phase#67, end_phase#68, is_current#69, is_constitutive#70, stable_id#71, version#72, created_date#73, modified_date#74, transcript_id#75L, rank#76, sequence#148], Statistics(sizeInBytes=12.3 PiB) +- Join Inner, (seq_region_id#63L = cast(seq_region_id#147 as bigint)), Statistics(sizeInBytes=14.4 PiB) :- Repartition 10, true, Statistics(sizeInBytes=13.6 MiB) : +- Filter isnotnull(seq_region_id#63L), Statistics(sizeInBytes=13.6 MiB) : +- Relation [exon_id#62L,seq_region_id#63L,seq_region_start#64L,seq_region_end#65L,seq_region_strand#66,phase#67,end_phase#68,is_current#69,is_constitutive#70,stable_id#71,version#72,created_date#73,modified_date#74,transcript_id#75L,rank#76] orc, Statistics(sizeInBytes=13.6 MiB) +- Repartition 10000, true, Statistics(sizeInBytes=1083.8 MiB) +- Filter isnotnull(seq_region_id#147), Statistics(sizeInBytes=1083.8 MiB) +- Relation [seq_region_id#147,sequence#148] orc, Statistics(sizeInBytes=1083.8 MiB) == Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- Project [seq_region_id#63L, exon_id#62L, seq_region_start#64L, seq_region_end#65L, seq_region_strand#66, phase#67, end_phase#68, is_current#69, is_constitutive#70, stable_id#71, version#72, created_date#73, modified_date#74, transcript_id#75L, rank#76, sequence#148] +- SortMergeJoin [seq_region_id#63L], [cast(seq_region_id#147 as bigint)], Inner :- Sort [seq_region_id#63L ASC NULLS FIRST], false, 0 : +- Exchange hashpartitioning(seq_region_id#63L, 200), ENSURE_REQUIREMENTS, [plan_id=335] : +- Exchange RoundRobinPartitioning(10), REPARTITION_BY_NUM, [plan_id=330] : +- Filter isnotnull(seq_region_id#63L) : +- FileScan orc [exon_id#62L,seq_region_id#63L,seq_region_start#64L,seq_region_end#65L,seq_region_strand#66,phase#67,end_phase#68,is_current#69,is_constitutive#70,stable_id#71,version#72,created_date#73,modified_date#74,transcript_id#75L,rank#76] Batched: true, DataFilters: [isnotnull(seq_region_id#63L)], Format: ORC, Location: InMemoryFileIndex(1 paths)[file:/nfs/production/flicek/ensembl/infrastructure/mira/19tmp/exons], PartitionFilters: [], PushedFilters: [IsNotNull(seq_region_id)], ReadSchema: struct<exon_id:bigint,seq_region_id:bigint,seq_region_start:bigint,seq_region_end:bigint,seq_regi... +- Sort [cast(seq_region_id#147 as bigint) ASC NULLS FIRST], false, 0 +- Exchange hashpartitioning(cast(seq_region_id#147 as bigint), 200), ENSURE_REQUIREMENTS, [plan_id=336] +- Exchange RoundRobinPartitioning(10000), REPARTITION_BY_NUM, [plan_id=331] +- Filter isnotnull(seq_region_id#147) +- FileScan orc [seq_region_id#147,sequence#148] Batched: true, DataFilters: [isnotnull(seq_region_id#147)], Format: ORC, Location: InMemoryFileIndex(1 paths)[file:/nfs/production/flicek/ensembl/infrastructure/mira/6tmp/sequence], PartitionFilters: [], PushedFilters: [IsNotNull(seq_region_id)], ReadSchema: struct<seq_region_id:string,sequence:string>
更新后的执行计划
== Optimized Logical Plan == Project [seq_region_id#63L, exon_id#62L, seq_region_start#64L, seq_region_end#65L, seq_region_strand#66, phase#67, end_phase#68, is_current#69, is_constitutive#70, stable_id#71, version#72, created_date#73, modified_date#74, transcript_id#75L, rank#76, sequence#148], Statistics(sizeInBytes=12.3 PiB) +- Join Inner, (seq_region_id#63L = cast(seq_region_id#147 as bigint)), Statistics(sizeInBytes=14.4 PiB) :- RepartitionByExpression [seq_region_id#63L], Statistics(sizeInBytes=13.6 MiB) : +- Filter isnotnull(seq_region_id#63L), Statistics(sizeInBytes=13.6 MiB) : +- Relation [exon_id#62L,seq_region_id#63L,seq_region_start#64L,seq_region_end#65L,seq_region_strand#66,phase#67,end_phase#68,is_current#69,is_constitutive#70,stable_id#71,version#72,created_date#73,modified_date#74,transcript_id#75L,rank#76] orc, Statistics(sizeInBytes=13.6 MiB) +- RepartitionByExpression [cast(seq_region_id#147 as bigint)], Statistics(sizeInBytes=1083.8 MiB) +- Filter isnotnull(seq_region_id#147), Statistics(sizeInBytes=1083.8 MiB) +- Relation [seq_region_id#147,sequence#148] orc, Statistics(sizeInBytes=1083.8 MiB) == Physical Plan == AdaptiveSparkPlan isFinalPlan=false +- Project [seq_region_id#63L, exon_id#62L, seq_region_start#64L, seq_region_end#65L, seq_region_strand#66, phase#67, end_phase#68, is_current#69, is_constitutive#70, stable_id#71, version#72, created_date#73, modified_date#74, transcript_id#75L, rank#76, sequence#148] +- SortMergeJoin [seq_region_id#63L], [cast(seq_region_id#147 as bigint)], Inner :- Sort [seq_region_id#63L ASC NULLS FIRST], false, 0 : +- Exchange hashpartitioning(seq_region_id#63L, 200), REPARTITION_BY_COL, [plan_id=330] : +- Filter isnotnull(seq_region_id#63L) : +- FileScan orc [exon_id#62L,seq_region_id#63L,seq_region_start#64L,seq_region_end#65L,seq_region_strand#66,phase#67,end_phase#68,is_current#69,is_constitutive#70,stable_id#71,version#72,created_date#73,modified_date#74,transcript_id#75L,rank#76] Batched: true, DataFilters: [isnotnull(seq_region_id#63L)], Format: ORC, Location: InMemoryFileIndex(1 paths)[file:/nfs/production/flicek/ensembl/infrastructure/mira/0tmp/exons], PartitionFilters: [], PushedFilters: [IsNotNull(seq_region_id)], ReadSchema: struct<exon_id:bigint,seq_region_id:bigint,seq_region_start:bigint,seq_region_end:bigint,seq_regi... +- Sort [cast(seq_region_id#147 as bigint) ASC NULLS FIRST], false, 0 +- Exchange hashpartitioning(cast(seq_region_id#147 as bigint), 200), REPARTITION_BY_COL, [plan_id=331] +- Filter isnotnull(seq_region_id#147) +- FileScan orc [seq_region_id#147,sequence#148] Batched: true, DataFilters: [isnotnull(seq_region_id#147)], Format: ORC, Location: InMemoryFileIndex(1 paths)[file:/nfs/production/flicek/ensembl/infrastructure/mira/0tmp/sequence], PartitionFilters: [], PushedFilters: [IsNotNull(seq_region_id)], ReadSchema: struct<seq_region_id:string,sequence:string>
优化方案
核心思路是避免将完整大字符串带入Join后的中间结果,提供两种可行实现:
方案1:广播小表+UDF按需截取
利用小表数据量极小的特点,将其广播到所有节点,直接在大表上通过UDF匹配id并截取子串,完全跳过Join操作:
from pyspark.sql import functions as F from pyspark.sql.types import StringType # 将小表转为字典并广播 string_map = {row.id: row.string for row in df2.collect()} broadcast_strings = spark.sparkContext.broadcast(string_map) # 自定义UDF:根据id、start、end直接截取子串 def extract_substring(id_val, start, end): full_str = broadcast_strings.value.get(id_val, "") length = end - start + 1 return full_str[start-1 : start-1 + length] if len(full_str) >= start-1 + length else "" substr_udf = F.udf(extract_substring, StringType()) # 直接在大表生成结果 result_df = df1.withColumn("seq", substr_udf(F.col("id"), F.col("start"), F.col("end"))) result_df.show()
方案2:分组聚合+Join后展开截取
先将大表按id分组,收集同一id的所有位置信息,再与小表Join(每个id仅关联一次大字符串),最后展开并截取子串:
# 大表按id分组,聚合位置信息 grouped_df = df1.groupBy("id").agg( F.collect_list(F.struct("start", "end")).alias("position_ranges") ) # Join后展开每个位置并截取子串 result_df = grouped_df.join(df2, on="id").select( "id", F.explode("position_ranges").alias("pos"), "string" ).select( "id", F.col("pos.start").alias("start"), F.col("pos.end").alias("end"), F.expr("substring(string, start, end - start + 1)").alias("seq") ) result_df.show()
计划优化说明
原执行计划中,SortMergeJoin会将完整的sequence字段带入中间结果,导致内存占用激增(统计显示达12.3PiB)。优化方案:
- 方案1完全规避Join,内存仅消耗大表数据量+广播的小表字典(约1GB)。
- 方案2通过分组减少Join后的记录数,每个id仅存储一次大字符串,避免百万条记录重复存储100MB字符串的问题。
内容的提问来源于stack exchange,提问作者user453575457
相关产品推荐
相关产品推荐

