【问题标题】:How to filter Spark dataframe by array column containing any of the values of some other dataframe/set如何按包含其他数据帧/集的任何值的数组列过滤 Spark 数据帧
【发布时间】:2020-10-15 18:33:22
【问题描述】:

我有一个包含一列数组字符串的数据框 A。

...
 |-- browse: array (nullable = true)
 |    |-- element: string (containsNull = true)
...

例如三个样本行将是

+---------+--------+---------+
| column 1|  browse| column n|
+---------+--------+---------+
|     foo1| [X,Y,Z]|     bar1|
|     foo2|   [K,L]|     bar2|
|     foo3|     [M]|     bar3|

另一个包含一列字符串的Dataframe B

|-- browsenodeid: string (nullable = true)

它的一些示例行将是

+------------+
|browsenodeid|
+------------+
|           A|
|           Z|
|           M|

如何过滤 A 以便保留 browse 包含来自 B 的任何 browsenodeid 值的所有行?根据上述示例,结果将是:

+---------+--=-----+---------+
| column 1|  browse| column n|
+---------+--------+---------+
|     foo1| [X,Y,Z]|     bar1| <- because Z is a value of B.browsenodeid
|     foo3|     [M]|     bar3| <- because M is a value of B.browsenodeid

如果我只有一个值,那么我会使用类似的东西

A.filter(array_contains(A("browse"), single_value))

但是我该如何处理值的列表或 DataFrame?

【问题讨论】:

    标签: apache-spark apache-spark-sql spark-dataframe


    【解决方案1】:

    我为此找到了一个优雅的解决方案,无需将 DataFrames/Datasets 转换为 RDDs。

    假设你有一个 DataFrame dataDF:

    +---------+--------+---------+
    | column 1|  browse| column n|
    +---------+--------+---------+
    |     foo1| [X,Y,Z]|     bar1|
    |     foo2|   [K,L]|     bar2|
    |     foo3|     [M]|     bar3|
    

    还有一个数组b,其中包含您要在browse 中匹配的值

    val b: Array[String] = Array(M,Z)
    

    实现udf:

    import org.apache.spark.sql.expressions.UserDefinedFunction
    import scala.collection.mutable.WrappedArray
    
    def array_contains_any(s:Seq[String]): UserDefinedFunction = {
    udf((c: WrappedArray[String]) =>
      c.toList.intersect(s).nonEmpty)
    }
    

    然后简单地使用filterwhere 函数(带有一点花哨的柯里化:P)进行过滤,如下所示:

    dataDF.where(array_contains_any(b)($"browse"))
    

    【讨论】:

      【解决方案2】:

      在 Spark >= 2.4.0 中,您可以使用arrays_overlap

      import org.apache.spark.sql.functions.{array, arrays_overlap, lit}
      
      val df = Seq(
        ("foo1", Seq("X", "Y", "Z"), "bar1"),
        ("foo2", Seq("K", "L"), "bar2"),
        ("foo3", Seq("M"), "bar3")
      ).toDF("col1", "browse", "coln")
      
      val b = Seq("M" ,"Z") 
      val searchArray = array(b.map{lit}:_*) // cast to lit(i) then create Spark array
      
      df.where(arrays_overlap($"browse", searchArray)).show()
      
      // +----+---------+----+
      // |col1|   browse|coln|
      // +----+---------+----+
      // |foo1|[X, Y, Z]|bar1|
      // |foo3|      [M]|bar3|
      // +----+---------+----+
      

      【讨论】:

        【解决方案3】:

        假设输入数据:Dataframe A

        browse
        200,300,889,767,9908,7768,9090
        300,400,223,4456,3214,6675,333
        234,567,890
        123,445,667,887
        

        你必须将它与 Dataframe B 匹配

        browsenodeid:(我把browsenodeid列展平)123,200,300

        val matchSet = "123,200,300".split(",").toSet
        val rawrdd = sc.textFile("D:\\Dataframe_A.txt")
        rawrdd.map(_.split("|"))
              .map(arr => arr(0).split(",").toSet.intersect(matchSet).mkString(","))
              .foreach(println)
        

        你的输出:

        300,200
        300
        123
        

        更新

        val matchSet = "A,Z,M".split(",").toSet
        
        val rawrdd = sc.textFile("/FileStore/tables/mvv45x9f1494518792828/input_A.txt")
        
        rawrdd.map(_.split("|"))
              .map(r => if (! r(1).split(",").toSet.intersect(matchSet).isEmpty) org.apache.spark.sql.Row(r(0),r(1), r(2))).collect.foreach(println)
        

        输出是

        foo1,X,Y,Z,bar1
        foo3,M,bar3
        

        【讨论】:

        • Manish 感谢您的回复,但这并不是我想要的。我在描述中添加了一些示例,以阐明我所拥有的以及我想要实现的目标
        • Manish 感谢您的回答。我想要一个可以插入到Datasetfilter/where 函数的解决方案,以便它更具可读性并且更容易集成到现有代码库中(主要围绕DataFrames 而不是@ 987654330@s)。检查我上面的答案,如果你喜欢它,请为我投票!
        猜你喜欢
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 2022-01-23
        • 2018-02-20
        • 2015-12-13
        • 2011-08-28
        相关资源
        最近更新 更多