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

