Spark Scala调用函数报Type Mismatch: 预期DataFrame却得到Dataset[DataFrame]
问题根因
- 直接触发类型不匹配报错的原因:
kvrdd1_toDF属于DataFrame/Dataset类型,调用map算子时会对每一行数据执行传入的匿名函数,你定义的function1返回值为DataFrame类型,因此map执行完成后的结果是Dataset[DataFrame]类型,和要赋值的df1的DataFrame类型完全不匹配,所以触发了该错误。 - 除此之外你的实现还有多个不符合Spark运行逻辑的问题,即便解决类型报错也无法正常运行:
- Spark的转换算子(如
map)内的逻辑会分发到Executor节点分布式执行,而SparkSession相关操作、DataFrame的创建、union、take等都属于仅能在Driver端运行的API,直接写在map内会抛出序列化错误、空指针异常等运行时错误。 - 你在
function1中调用的take(1)是Action算子,每执行一次都会触发一次Spark作业,反复调用会导致性能指数级下降。 - 你想通过逐行处理数据迭代更新DataFrame的思路不符合Spark的分布式计算模型,DataFrame是不可变数据结构,每次
union都会生成新的DataFrame,迭代大量数据时会产生严重的性能问题甚至OOM。
- Spark的转换算子(如
修正方案
你的需求看起来是要将kvrdd1_toDF的对应字段和原有DataFrame做合并,可直接按如下逻辑实现,避免在算子内操作DataFrame:
// 先从kvrdd1_toDF提取你需要的字段转成待合并的DataFrame val appendDf = kvrdd1_toDF.select(col("_c2").alias("g"), lit("nn").alias("tg")) // 如果原有df1是空的就用appendDf,否则合并两个df val finalDf = if(df1.take(1).isEmpty) { appendDf } else { df1.union(appendDf) }
如果确实需要逐行自定义逻辑处理,先把所有处理逻辑作用在Row上,统一处理完所有行之后再一次性转成DataFrame,不要逐行生成DataFrame再合并。
内容的提问来源于stack exchange,提问作者CloudSparkie
相关产品推荐
相关产品推荐

