如何在PySpark DataFrame中explode展开特征向量为单独数据列
PySpark实现特征向量与文本数组配对展开
核心思路
Spark MLlib的Vector类型(包含稠密向量DenseVector、稀疏向量SparseVector)不支持直接与数组配对展开,需要先将Vector转换为普通数值数组,再与文本数组按位置配对后炸开即可得到目标结构。
完整实现代码
首先导入依赖包:
from pyspark.sql import SparkSession from pyspark.ml.linalg import Vectors, VectorUDT from pyspark.sql.functions import udf, col, explode, arrays_zip from pyspark.sql.types import ArrayType, FloatType, StructType, StructField, IntegerType, StringType
构造示例源数据(和给出的原始表结构一致):
spark = SparkSession.builder.appName("vector_expand").getOrCreate() # 原始数据 source_data = [ (0, ["a", "b", "c"], Vectors.sparse(3, [0,1,2], [1.0,1.0,1.0])), (1, ["a", "b", "c"], Vectors.sparse(3, [0,1,2], [2.0,2.0,1.0])) ] source_schema = StructType([ StructField("id", IntegerType(), True), StructField("texts", ArrayType(StringType()), True), StructField("vector", VectorUDT(), True) ]) source_df = spark.createDataFrame(source_data, schema=source_schema)
执行转换逻辑:
# 定义UDF将Vector类型转为普通浮点数组 vec2arr_udf = udf(lambda vec: vec.toArray().tolist(), ArrayType(FloatType())) result_df = source_df.withColumn("vec_arr", vec2arr_udf(col("vector"))) \ # 按位置将文本数组和数值数组配对后炸开 .withColumn("pair", explode(arrays_zip(col("texts"), col("vec_arr")))) \ # 提取配对后的字段 .select( col("id"), col("pair.texts").alias("texts"), col("pair.vec_arr").alias("list_2") ) # 打印结果验证 result_df.show()
运行后输出和目标表完全一致:
+---+-----+------+ | id|texts|list_2| +---+-----+------+ | 0| a| 1.0| | 0| b| 1.0| | 0| c| 1.0| | 1| a| 2.0| | 1| b| 2.0| | 1| c| 1.0| +---+-----+------+
注意事项
- 上述写法同时兼容稀疏向量和稠密向量,
toArray()方法会自动将稀疏向量补全为对应维度的完整数组,不会出现索引缺失错位问题 - 请保证
texts数组长度和vector的维度一致,否则会出现位置匹配错误 - 如果使用Spark 2.3及以下版本(无
arrays_zip函数),可以用posexplode按索引匹配实现,兼容写法如下:
from pyspark.sql.functions import posexplode result_df = source_df.withColumn("vec_arr", vec2arr_udf(col("vector"))) \ .select( col("id"), posexplode(col("texts")).alias("pos", "texts"), col("vec_arr") ) \ .withColumn("list_2", col("vec_arr")[col("pos")]) \ .select("id", "texts", "list_2")
内容的提问来源于stack exchange,提问作者Devansh Popat
相关产品推荐
相关产品推荐

