【问题标题】:Spark: Joining with arraySpark:加入数组
【发布时间】:2018-01-14 17:12:52
【问题描述】:

我需要将带有字符串列的数据框与带有字符串数组的数据框连接起来,这样如果数组中的一个值匹配,行就会连接起来。

我试过了,但我猜它不支持。 还有其他方法吗?

import org.apache.spark.SparkConf
import org.apache.spark.sql.SparkSession

val sparkConf = new SparkConf().setMaster("local[*]").setAppName("test")
val spark = SparkSession.builder().config(sparkConf).getOrCreate()

import spark.implicits._

val left = spark.sparkContext.parallelize(Seq(1, 2, 3)).toDF("col1")
val right = spark.sparkContext.parallelize(Seq((Array(1, 2), "Yes"),(Array(3),"No"))).toDF("col1", "col2")

left.join(right,"col1")

投掷:

org.apache.spark.sql.AnalysisException: 无法解析 '(col1 =col1)' 由于数据 类型不匹配:'(col1 =

中的不同类型

col1)'(整数和数组)。;;

【问题讨论】:

    标签: scala apache-spark apache-spark-sql


    【解决方案1】:

    一种选择是创建一个 UDF 来构建您的连接条件:

    import org.apache.spark.sql.functions._
    import scala.collection.mutable.WrappedArray
    
    val left = spark.sparkContext.parallelize(Seq(1, 2, 3)).toDF("col1")
    val right = spark.sparkContext.parallelize(Seq((Array(1, 2), "Yes"),(Array(3),"No"))).toDF("col1", "col2")
    
    val checkValue = udf { 
      (array: WrappedArray[Int], value: Int) => array.contains(value) 
    }
    val result = left.join(right, checkValue(right("col1"), left("col1")), "inner")
    
    result.show
    
    +----+------+----+
    |col1|  col1|col2|
    +----+------+----+
    |   1|[1, 2]| Yes|
    |   2|[1, 2]| Yes|
    |   3|   [3]|  No|
    +----+------+----+
    

    【讨论】:

    • 好吧,我猜这个解决方案稍微好一点,但我不确定,因为 UDF 可能很慢。然而,在内存方面,这个解决方案肯定会比另一个解决方案使用更少的空间。
    • 我想我会更好地知道如果我会尝试两种方式。谢谢。
    • @aclokay 你可以尝试对你的数据帧进行二元化,然后通过两个条件加入,bucketId 和 UDF。然后 Spark 将能够排序合并连接而不是笛卡尔连接
    • @T.Gawęda 你说的 bickitize 是什么意思?
    • @aclokay Bucketize - 抱歉名称错误;)这意味着,例如,您可以添加名称为 bucketId 且值 = 数组中元素数的列。然后加入 a.bucketId == b.bucketId 和 udf - 它应该改变查询计划:)
    【解决方案2】:

    最简洁的方法是使用 array_contains spark sql 表达式,如下所示,这表示我已经将其性能与执行爆炸和连接的性能进行了比较,如上一个答案和爆炸似乎性能更高。

    import org.apache.spark.sql.functions.expr
    import spark.implicits._
    
    val left = Seq(1, 2, 3).toDF("col1")
    
    val right = Seq((Array(1, 2), "Yes"),(Array(3),"No")).toDF("col1", "col2").withColumnRenamed("col1", "col1_array")
    
    val joined = left.join(right, expr("array_contains(col1_array, col1)")).show
    
    +----+----------+----+
    |col1|col1_array|col2|
    +----+----------+----+
    |   1|    [1, 2]| Yes|
    |   2|    [1, 2]| Yes|
    |   3|       [3]|  No|
    +----+----------+----+
    

    请注意,您不能直接使用 org.apache.spark.sql.functions.array_contains 函数,因为它要求第二个参数是文字而不是列表达式。

    【讨论】:

      【解决方案3】:

      您可以在加入之前在 Array 列上使用 explode。 Explode 为数组中的每个元素创建一个新行:

      right = right.withColumn("exploded_col",explode(right("col1")))
      right.show()
      
      +------+----+--------------+
      |  col1|col2|exploded_col_1|
      +------+----+--------------+
      |[1, 2]| Yes|             1|
      |[1, 2]| Yes|             2|
      |   [3]|  No|             3|
      +------+----+--------------+
      

      然后您可以轻松加入您的第一个数据集。

      【讨论】:

      • 谢谢!这与下面的答案相比如何......性能方面?
      • @aclokay 我认为我的回答比较慢,尤其是当您的数组变大时,因为您必须为数组中的每个元素创建(即复制)一行。
      • 我也想过这种方式,但我希望 spark 能做一些魔法来以某种方式优化它。不过谢谢!
      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2016-07-27
      • 2017-08-19
      • 2020-08-23
      • 2022-11-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      相关资源
      最近更新 更多