如何将LDA.transform生成的含Vector类型的Spark DataFrame转为CSV?
解决Spark DataFrame中Vector类型字段写入CSV的问题
嘿,这个场景我之前也遇到过!问题出在Spark的CSV写入器默认不支持直接序列化Vector这种复杂类型,所以咱们得先把向量转换成字符串格式,就能用你熟悉的方式保存成CSV了,还能完美保留show()里显示的那种[x1,x2,...]格式。
下面是Scala环境下的具体操作步骤:
1. 导入必要的类
首先得导入UDF相关工具和Vector的处理类:
import org.apache.spark.sql.functions.udf import org.apache.spark.ml.linalg.Vector
2. 写个小UDF把Vector转成字符串
咱们自定义一个简单的函数,把Vector对象直接转成它的字符串表示:
val vectorToString = udf((vec: Vector) => vec.toString)
3. 转换DataFrame里的Vector列
你可以选择添加一个新的字符串列,或者直接替换原有的topicDistribution列:
// 方案一:新增字符串列,保留原Vector列(方便后续可能的其他操作) val dfWithStringVector = df_new.withColumn("topicDistribution_str", vectorToString($"topicDistribution")) // 方案二:直接替换原列(如果不需要再用原Vector列的话) val dfWithStringVector = df_new.withColumn("topicDistribution", vectorToString($"topicDistribution"))
4. 正常写入CSV就完事了
现在就可以用你习惯的代码写入CSV啦:
// 用方案一的话,记得选择要输出的列 dfWithStringVector.select("label", "topicDistribution_str") .coalesce(1) .write .option("header", true) .csv("<你的输出路径>") // 用方案二的话,直接写原列名就行 dfWithStringVector .coalesce(1) .write .option("header", true) .csv("<你的输出路径>")
额外提一句
要是你用的是Spark 3.0及以上版本,其实也可以用内置函数vector_to_array把向量转成数组,再用array_join拼接成字符串,但这种方式出来的格式是x1,x2,x3而不是你想要的[x1,x2,x3]。所以自定义UDF的方式更贴合你的需求,能完全复现show()里的格式。
内容的提问来源于stack exchange,提问作者BlueIvy
相关产品推荐
相关产品推荐

