python - 如何将向量拆分为列 - 使用 PySpark

标签 python apache-spark pyspark apache-spark-sql apache-spark-ml

上下文: 我有一个 DataFrame 有 2 列:单词和向量。其中“vector”的列类型为VectorUDT

一个例子:

word    |  vector
assert  | [435,323,324,212...]

我想得到这个:

word   |  v1 | v2  | v3 | v4 | v5 | v6 ......
assert | 435 | 5435| 698| 356|....

问题:

如何使用 PySpark 将包含向量的列拆分为每个维度的多个列?

提前致谢

最佳答案

Spark >= 3.0.0

从 Spark 3.0.0 开始,这可以在不使用 UDF 的情况下完成。

from pyspark.ml.functions import vector_to_array

(df
    .withColumn("xs", vector_to_array("vector")))
    .select(["word"] + [col("xs")[i] for i in range(3)]))

## +-------+-----+-----+-----+
## |   word|xs[0]|xs[1]|xs[2]|
## +-------+-----+-----+-----+
## | assert|  1.0|  2.0|  3.0|
## |require|  0.0|  2.0|  0.0|
## +-------+-----+-----+-----+

Spark <3.0.0

一种可能的方法是与 RDD 相互转换:

from pyspark.ml.linalg import Vectors

df = sc.parallelize([
    ("assert", Vectors.dense([1, 2, 3])),
    ("require", Vectors.sparse(3, {1: 2}))
]).toDF(["word", "vector"])

def extract(row):
    return (row.word, ) + tuple(row.vector.toArray().tolist())

df.rdd.map(extract).toDF(["word"])  # Vector values will be named _2, _3, ...

## +-------+---+---+---+
## |   word| _2| _3| _4|
## +-------+---+---+---+
## | assert|1.0|2.0|3.0|
## |require|0.0|2.0|0.0|
## +-------+---+---+---+

另一种解决方案是创建 UDF:

from pyspark.sql.functions import udf, col
from pyspark.sql.types import ArrayType, DoubleType

def to_array(col):
    def to_array_(v):
        return v.toArray().tolist()
    # Important: asNondeterministic requires Spark 2.3 or later
    # It can be safely removed i.e.
    # return udf(to_array_, ArrayType(DoubleType()))(col)
    # but at the cost of decreased performance
    return udf(to_array_, ArrayType(DoubleType())).asNondeterministic()(col)

(df
    .withColumn("xs", to_array(col("vector")))
    .select(["word"] + [col("xs")[i] for i in range(3)]))

## +-------+-----+-----+-----+
## |   word|xs[0]|xs[1]|xs[2]|
## +-------+-----+-----+-----+
## | assert|  1.0|  2.0|  3.0|
## |require|  0.0|  2.0|  0.0|
## +-------+-----+-----+-----+

对于 Scala 等价物,请参阅 Spark Scala: How to convert Dataframe[vector] to DataFrame[f1:Double, ..., fn: Double)] .

关于python - 如何将向量拆分为列 - 使用 PySpark,我们在Stack Overflow上找到一个类似的问题: https://stackoverflow.com/questions/38384347/

相关文章:

python - doc2vec:性能测量和 'workers' 参数

java - spark-class : line 71. ..没有那个文件或目录

apache-spark - 如何在Spark中计算过去时间值 `percent_rank`?

python - OSError 38 [Errno 38] 多处理

python - Pandas 将列从字符串转换为 float

python - 为什么第 13 行 a = input() 影响第 15 行 b=input(int() 上的变量?

apache-spark - pyspark.sql.DataFrameWriter.saveAsTable() 的格式

python - Spark : Pyspark: how to monitor python worker processes

apache-spark - 使用 registerTempTable 找不到表或 View

apache-spark - 使用 HDFS 存储的 Spark 作业