【问题标题】:How can I use the literal value of a spark dataframe column?如何使用 spark 数据框列的文字值?
【发布时间】:2017-10-08 02:06:38
【问题描述】:

我有一个看起来像这样的简单数据框,

+---+---+---+---+
|nm | ca| cb| cc|
+---+---+---+---+
|  a|123|  0|  0|
|  b|  1|  2|  3|
|  c|  0|  1|  0|
+---+---+---+---+

我想做的是,

+---+---+---+---+---+
|nm |ca |cb |cc |p  |
+---+---+---+---+---+
|a  |123|0  |0  |1  |
|b  |1  |2  |3  |1  |
|c  |0  |1  |0  |0  |
+---+---+---+---+---+

基本上添加了一个新列p,这样,如果列nm 的值为'a',则检查列ca 是否>0,如果是,则为列p1 设置'1',否则为0。

我的代码,

        def purchaseCol: UserDefinedFunction =
    udf((brand: String) => s"c$brand")

val a = ss.createDataset(List(
        ("a", 123, 0, 0),
        ("b", 1, 2, 3),
        ("c", 0, 1, 0)))
    .toDF("nm", "ca", "cb", "cc")

a.show()
a.withColumn("p", when(lit(DataFrameUtils.purchaseCol($"nm")) > 0, 1).otherwise(0))
.show(false)

它似乎不起作用,并且为 col 'p' 中的所有行返回 0。

PS:列数超过100,是动态生成的。

【问题讨论】:

    标签: apache-spark


    【解决方案1】:

    映射rdd,计算并将p 添加到每一行:

    val a = sc.parallelize(
        List(("a", 123, 0, 0),
             ("b", 1, 2, 3),
             ("c", 0, 1, 0))
    ).toDF("nm", "ca", "cb", "cc")
    
    import org.apache.spark.sql.Row
    import org.apache.spark.sql.types._
    
    val b = a.rdd.map(r => {
        val s = r.getAs[String]("nm")
        val v = r.getAs[Int](s"c$s")
        val p = if(v > 0) 1 else 0
        Row.fromSeq(r.toSeq :+ p)
    })
    
    val new_schema = StructType(a.schema :+ StructField("p", IntegerType, true))
    
    val df_new = spark.createDataFrame(b, new_schema)
    
    df_new.show
    +---+---+---+---+---+
    | nm| ca| cb| cc|  p|
    +---+---+---+---+---+
    |  a|123|  0|  0|  1|
    |  b|  1|  2|  3|  1|
    |  c|  0|  1|  0|  0|
    +---+---+---+---+---+
    

    【讨论】:

    • 是的,这是我大部分时间所做的,并不是真的想使用 rdds 并重建数据帧。正在寻找更优雅的解决方案,但是,+1
    【解决方案2】:

    如果“c*”列数有限,可以使用所有值的UDF:

      val nameMatcherFunct = (nm: String, ca: Int, cb: Int, cc: Int) => {
      val value = nm match {
        case "a" => ca
        case "b" => cb
        case "c" => cc
      }
      if (value > 0) 1 else 0
    }
    
    def purchaseValueUDF = udf(nameMatcherFunct)
    
    val result = a.withColumn("p", purchaseValueUDF(col("nm"), col("ca"), col("cb"), col("cc")))
    

    如果您有许多“c*”列,可以使用以 Row 作为参数的函数: How to pass whole Row to UDF - Spark DataFrame filter

    【讨论】:

    • 这就是问题所在,列是动态的
    【解决方案3】:

    看看你的逻辑

    如果nm列的值为'a',则检查ca列是否>0,如果是,则为p1列设置'1',否则为0。

    你可以这样做

    import org.apache.spark.sql.functions._
    a.withColumn("p", when((col("nm") === lit("a")) && (col("ca") > 0), lit(1)).otherwise(lit(0)))
    

    但是查看您的输出dataframe,您需要|| 而不是&&

    import org.apache.spark.sql.functions._
    a.withColumn("p", when((col("nm") === lit("a")) || (col("ca") > 0), lit(1)).otherwise(lit(0)))
    

    【讨论】:

    • 挑战是以编程方式进行,a 是 ca 的一部分,b 是 cb 的一部分,等等。
    【解决方案4】:
    val a1 = sc.parallelize(
        List(("a", 123, 0, 0),
             ("b", 1, 2, 3),
             ("c", 0, 1, 0))
    ).toDF("nm", "ca", "cb", "cc")
    
    a1.show()
    
    
    +---+---+---+---+
    | nm| ca| cb| cc|
    +---+---+---+---+
    |  a|123|  0|  0|
    |  b|  1|  2|  3|
    |  c|  0|  1|  0|
    +---+---+---+---+
    
    
    val newDf = a1.withColumn("P", when($"ca" > 0, 1).otherwise(0))
    newDf.show()
    
    +---+---+---+---+---+
    | nm| ca| cb| cc|  P|
    +---+---+---+---+---+
    |  a|123|  0|  0|  1|
    |  b|  1|  2|  3|  1|
    |  c|  0|  1|  0|  0|
    +---+---+---+---+---+
    

    【讨论】:

    • 不,这是错误的,您必须为每一行考虑不同的列。假设如果在 cc 之后有另一列 cd 对于另一行 d,1,4,5,0 ,则对应于此的 P 的值将为 0 但您的逻辑会将其标记为 1
    猜你喜欢
    • 2021-12-12
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2023-02-02
    • 2021-01-30
    • 1970-01-01
    • 1970-01-01
    • 2021-09-24
    相关资源
    最近更新 更多