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

Spark中如何将多个小DataFrame存入合并DataFrame的列并实现读取?

解决Spark中存储并还原嵌套DataFrame数据的问题

看起来你是想把多个小型DataFrame打包到一个统一的DataFrame里存储,但目前的写法是把DataFrame对象本身序列化到字段中,这不是正确的做法——Spark的DataFrame本质是分布式数据集的逻辑引用,序列化这个对象后反序列化会因为依赖上下文失效,自然没法获取到实际数据。下面给你两种可行的解决方案:

方案1:将DataFrame数据转为嵌套结构存储(推荐)

这种方式直接存储DataFrame里的实际数据,不需要序列化整个DataFrame对象,适合你的小型DataFrame场景:

case class AccountsData(empId: String, deptId: String, fName: String, lName: String)
case class PositionsData(positionId: String, date: String)
import sparkSession.implicits._
import org.apache.spark.sql.functions._
import org.apache.spark.sql.Row

// 构建原始小DataFrame
val accountdata = List(AccountsData("1", "100", "FN1", "LN1"),AccountsData("2", "100", "FN2", "LN2"))
val accountsDf = accountdata.toDF()
val positionsData = List(PositionsData("10011001", "01-Jan- 2022"),PositionsData("20012001", "02-Jan-2022"))
val positionsDf = positionsData.toDF()

// 将每个小DataFrame的数据转为数组,再构建合并后的DataFrame
val combinedDf = List(
  ("accounts-data", accountsDf.collect()),
  ("positions-data", positionsDf.collect())
).toDF("key", "actual-dataframe")

// 还原accounts-data对应的DataFrame
val acctRetainedRows = combinedDf
  .filter(col("key") === "accounts-data")
  .select(col("actual-dataframe"))
  .as[Array[Row]]
  .first()

// 用原始Schema重新构建DataFrame
val acctDfRestored = sparkSession.createDataFrame(
  sparkSession.sparkContext.parallelize(acctRetainedRows),
  accountsDf.schema
)

// 验证结果
acctDfRestored.show()

说明

  • collect()会把小型DataFrame的数据拉到Driver节点,完全适配你的场景;
  • 还原时通过原始Schema重新创建DataFrame,能完整保留数据结构和内容;
  • 不需要自定义Kryo编码器,Spark原生支持Array[Row]的序列化逻辑。

方案2:序列化DataFrame数据为二进制(适合需同时存储Schema的场景)

如果需要同时保存数据和Schema信息,可以把整个DataFrame序列化为二进制(比如Parquet格式),读取时再反序列化:

import org.apache.spark.sql.DataFrame
import java.io.{ByteArrayInputStream, ByteArrayOutputStream}
import org.apache.spark.sql.execution.datasources.parquet.ParquetStreams

// 自定义DataFrame序列化/反序列化工具方法
def serializeDf(df: DataFrame): Array[Byte] = {
  val baos = new ByteArrayOutputStream()
  ParquetStreams.writeDataFrame(df).to(baos).close()
  baos.toByteArray
}

def deserializeDf(bytes: Array[Byte], spark: org.apache.spark.sql.SparkSession): DataFrame = {
  val bais = new ByteArrayInputStream(bytes)
  ParquetStreams.readDataFrame(spark).from(bais)
}

// 构建合并后的DataFrame
val combinedDf = List(
  ("accounts-data", serializeDf(accountsDf)),
  ("positions-data", serializeDf(positionsDf))
).toDF("key", "actual-dataframe")

// 还原accounts-data对应的DataFrame
val acctBytes = combinedDf
  .filter(col("key") === "accounts-data")
  .select(col("actual-dataframe"))
  .as[Array[Byte]]
  .first()

val acctDfRestored = deserializeDf(acctBytes, sparkSession)

// 验证结果
acctDfRestored.show()

说明

  • 用Parquet序列化能同时保存数据和Schema,还原时不需要依赖原始Schema定义;
  • 序列化效率高,适合结构化数据存储;
  • 同样只适合小型DataFrame,避免单个Row存储过大数据导致内存溢出。

为什么你的原始代码无法读取数据?

你之前用Kryo序列化的是DataFrame对象本身,这个对象包含了SparkSession、执行计划等上下文引用,反序列化后这些引用已经失效,自然无法访问到实际数据。你需要序列化的是DataFrame中的数据内容,而不是DataFrame这个容器对象。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.28 13:23:13