【问题标题】:How does Spark DataFrame distinguish between different VectorUDT objects?Spark DataFrame 如何区分不同的 VectorUDT 对象?
【发布时间】:2016-12-05 10:28:13
【问题描述】:

我正在尝试了解 DataFrame 列类型。当然,DataFrame 不是物化对象,它只是 Spark 的一组指令,将来要转换成代码。但我想象这个类型列表代表了在执行操作时可能在 JVM 中实现的对象类型。

import pyspark
import pyspark.sql.types as T
import pyspark.sql.functions as F
data = [0, 3, 0, 4]
d = {}
d['DenseVector'] = pyspark.ml.linalg.DenseVector(data)
d['old_DenseVector'] = pyspark.mllib.linalg.DenseVector(data)
d['SparseVector'] = pyspark.ml.linalg.SparseVector(4, dict(enumerate(data)))
d['old_SparseVector'] = pyspark.mllib.linalg.SparseVector(4, dict(enumerate(data)))
df = spark.createDataFrame([d])
df.printSchema()

四个向量值的列在printSchema()(或schema)中看起来相同:

root
 |-- DenseVector: vector (nullable = true)
 |-- SparseVector: vector (nullable = true)
 |-- old_DenseVector: vector (nullable = true)
 |-- old_SparseVector: vector (nullable = true)

但是当我逐行检索它们时,它们结果是不同的:

> for x in df.first().asDict().items():
  print(x[0], type(x[1]))
(2) Spark Jobs
old_SparseVector <class 'pyspark.mllib.linalg.SparseVector'>
SparseVector <class 'pyspark.ml.linalg.SparseVector'>
old_DenseVector <class 'pyspark.mllib.linalg.DenseVector'>
DenseVector <class 'pyspark.ml.linalg.DenseVector'>

我对@9​​87654327@ 类型的含义感到困惑(相当于VectorUDT 用于编写UDF)。 DataFrame 如何知道它在每个 vector 列中具有四种向量类型中的哪一种?这些向量列中的数据是否存储在 JVM 或 python VM 中?如果VectorUDT不是listed here的官方类型之一,为什么VectorUDT可以存储在DataFrame中?

(我知道mllib.linalg 中的四种向量类型中的两种最终将被弃用。)

【问题讨论】:

    标签: apache-spark dataframe pyspark apache-spark-mllib apache-spark-ml


    【解决方案1】:

    VectorUDT 怎么可以存储在 DataFrame 中

    UDT a.k.a 用户定义类型应该是这里的提示。 Spark 提供(现在是私有的)机制来将自定义对象存储在 DataFrame 中。您可以查看我对How to define schema for custom type in Spark SQL? 的回答或 Spark 源以了解详细信息,但长话短说,这一切都是关于解构对象并将它们编码为 Catalyst 类型。

    我对向量类型的含义感到困惑

    很可能是因为您看错了东西。简短的描述很有用,但它不能确定类型。相反,您应该检查架构。让我们创建另一个数据框:

    import pyspark.mllib.linalg as mllib
    import pyspark.ml.linalg as ml
    
    df = sc.parallelize([
        (mllib.DenseVector([1, ]), ml.DenseVector([1, ])),
        (mllib.SparseVector(1, [0, ], [1, ]), ml.SparseVector(1, [0, ], [1, ]))
    ]).toDF(["mllib_v", "ml_v"])
    
    df.show()
    
    ## +-------------+-------------+
    ## |      mllib_v|         ml_v|
    ## +-------------+-------------+
    ## |        [1.0]|        [1.0]|
    ## |(1,[0],[1.0])|(1,[0],[1.0])|
    ## +-------------+-------------+
    

    并检查数据类型:

    {s.name: type(s.dataType) for s in df.schema}
    
    ## {'ml_v': pyspark.ml.linalg.VectorUDT,
    ##  'mllib_v': pyspark.mllib.linalg.VectorUDT}
    

    正如您所见,UDT 类型是完全限定的,因此这里没有混淆。

    DataFrame如何知道它在每个向量列中有四种向量类型中的哪一种?

    如上所示,DataFrame 只知道其架构,可以区分ml / mllib 类型,但不关心向量变体(稀疏或密集)。

    向量类型由其type 字段(byte 字段,0 -> 稀疏,1 -> 密集)确定,但总体架构是相同的。 mlmllib 之间的内部表示也没有区别。

    这些向量列中的数据是存储在JVM还是Python中

    DataFrame 是一个纯 JVM 实体。 Python 互操作性是通过耦合的 UDT 类实现的:

    • Scala UDT 可以定义pyUDT 属性。
    • Python UDT 可以定义scalaUDT 属性。

    【讨论】:

    • 泰!我试过print(df.schema),但VectorUDT instances 的字符串表示不包括完整的限定符,我没想过直接检查他们的类。我只是对一些奇怪的魔法感到不舒服,但现在一切似乎都很合理。
    猜你喜欢
    • 2017-01-26
    • 1970-01-01
    • 1970-01-01
    • 2017-01-15
    • 2016-03-06
    • 2021-05-28
    • 2018-01-19
    • 1970-01-01
    • 2016-09-19
    相关资源
    最近更新 更多