Spark写入Cassandra出现Map转(AnyRef,AnyRef)类型转换异常如何解决
问题根源
你当前报错和schema不符合预期的核心原因是collect_list(map($"key", $"value"))的返回结果是Array[Map[String,String]]类型(也就是多个单键值对小map组成的数组),但Cassandra表要求attributes字段是单个Map[String,String]类型,类型不匹配导致转换失败。
解决方案
适用Spark 2.4及以上版本
直接用Spark内置的map_from_entries函数,把收集到的键值对结构直接转成单个map,替换你原来的agg逻辑即可:
原来的错误行:
.agg(collect_list(map($"key", $"value")).as("attributes"))
替换为:
.agg(map_from_entries(collect_list(struct($"key", $"value"))).as("attributes"))
修正后完整代码
val maxItem = 65000 dataFrame.select($"partition_key", $"row_key", $"data_as_of_date", posexplode($"attributes")) .withColumn("group", $"pos".divide(maxItem).cast("int")) .groupBy($"partition_key", $"row_key", $"data_as_of_date", $"group") // 核心修改点:收集键值对结构后直接转成单个map .agg(map_from_entries(collect_list(struct($"key", $"value"))).as("attributes")) .select($"partition_key", $"row_key", $"group", $"attributes", $"data_as_of_date") .write .format("org.apache.spark.sql.cassandra") .mode("append") .options(Map( "keyspace" -> keySpace, "table" -> tableName )) .save()
低版本Spark兼容方案(2.4以下)
如果你的Spark版本低于2.4没有map_from_entries函数,可以自定义UDF实现数组合并为单个map:
import org.apache.spark.sql.functions.udf import scala.collection.mutable.Map val mergeMaps = udf { arr: Seq[Map[String, String]] => val merged = Map.empty[String, String] arr.foreach(merged ++= _) merged.toMap } // 调用时agg部分改为 .agg(mergeMaps(collect_list(map($"key", $"value"))).as("attributes"))
修改后重新执行,attributes字段的schema就会变成预期的Map[String,String]类型,即可正常写入Cassandra。
内容的提问来源于stack exchange,提问作者vijayinani
相关产品推荐
相关产品推荐

