代码之家  ›  专栏  ›  技术社区  ›  Tasos

拆分文本并在Spark数据框中查找常用词

  •  0
  • Tasos  · 技术社区  · 7 年前

    我正在用spark处理scala,我有一个包含两列文本的数据框架。

    这些列的格式是“term1,term2,term3,…”,我想创建第三列,其中两个列的通用术语。

    例如

    Col1 
    orange, apple, melon
    party, clouds, beach
    
    Col2
    apple, apricot, watermelon
    black, yellow, white
    

    结果是

    Col3
    1
    0
    

    到目前为止,我所做的是创建一个UDF,它分割文本并得到两列的交集。

    val common_terms = udf((a: String, b: String) => if (a.isEmpty || b.isEmpty) {
          0
        } else {
          split(a, ",").intersect(split(b, ",")).length
        })
    

    然后在我的数据框架上

    val results = termsDF.withColumn("col3", common_terms(col("col1"), col("col2"))
    

    但我有以下错误

    Error:(96, 13) type mismatch;
     found   : String
     required: org.apache.spark.sql.Column
          split(a, ",").intersect(split(b, ",")).length
    

    我很感谢你的帮助,因为我是斯卡拉的新手,只是想从在线教程中学习。

    编辑:

    val common_authors = udf((a: String, b: String) => if (a != null || b != null) {
          0
        } else {
          val tempA = a.split( ",")
          val tempB = b.split(",")
          if ( tempA.isEmpty || tempB.isEmpty ) {
            0
          } else {
            tempA.intersect(tempB).length
          }
        })
    

    编辑之后,如果我尝试 termsDF.show() 它运行。但是如果我做那样的事 termsDF.orderBy(desc("col3")) 然后我得到了一个 java.lang.NullPointerException

    2 回复  |  直到 7 年前
        1
  •  2
  •   M. Alexandru    7 年前

    尝试

    val common_terms = udf((a: String, b: String) => if (a.isEmpty || b.isEmpty) {
          0
        } else {
            var tmp1 = a.split(",")
            var tmp2 = b.split(",")
          tmp1.intersect(tmp2).length
        })
    
    val results = termsDF.withColumn("col3", common_terms($"a", $"b")).show
    

    split(a,”,)它是一个火花柱函数。 您使用的是UDF,因此需要使用string.split()wich是一个scala函数

    编辑后:将空验证更改为==not!=

        2
  •  0
  •   stack0114106    7 年前

    在Spark2.4SQL中,您可以在不使用UDF的情况下获得相同的结果。看看这个:

    scala> val df = Seq(("orange,apple,melon","apple,apricot,watermelon"),("party,clouds,beach","black,yellow,white"), ("orange,apple,melon","apple,orange,watermelon")).toDF("col1","col2")
    df: org.apache.spark.sql.DataFrame = [col1: string, col2: string]
    
    scala>
    
    scala> df.createOrReplaceTempView("tasos")
    
    scala> spark.sql(""" select col1,col2, filter(split(col1,','), x -> array_contains(split(col2,','),x) ) a1 from tasos """).show(false)
    +------------------+------------------------+---------------+
    |col1              |col2                    |a1             |
    +------------------+------------------------+---------------+
    |orange,apple,melon|apple,apricot,watermelon|[apple]        |
    |party,clouds,beach|black,yellow,white      |[]             |
    |orange,apple,melon|apple,orange,watermelon |[orange, apple]|
    +------------------+------------------------+---------------+
    

    如果你想要尺寸,那么

    scala> spark.sql(""" select col1,col2, filter(split(col1,','), x -> array_contains(split(col2,','),x) ) a1 from tasos """).withColumn("a1_size",size('a1)).show(false)
    +------------------+------------------------+---------------+-------+
    |col1              |col2                    |a1             |a1_size|
    +------------------+------------------------+---------------+-------+
    |orange,apple,melon|apple,apricot,watermelon|[apple]        |1      |
    |party,clouds,beach|black,yellow,white      |[]             |0      |
    |orange,apple,melon|apple,orange,watermelon |[orange, apple]|2      |
    +------------------+------------------------+---------------+-------+
    
    
    scala>
    
    推荐文章