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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 03:27:09