【问题标题】:Spark sql Dataframe joins what's going on?Spark sql Dataframe 加入是怎么回事?
【发布时间】:2017-04-07 19:04:12
【问题描述】:

我有两个数据框,为简单起见,让我们左右调用它们,我将只显示示例结构。

数据框“左”:(这个数据框很大)

源代码 |夏令时 ------------ 乙 |一种 c | b 一个 | C

数据框“正确”(这个数据框很小)

位置 |姓名 ------------ 一个 |伦敦 乙 |巴黎

这两个数据框都是使用 hive 上下文和 sql 语句创建的。

如果我按如下方式在左侧数据帧上运行连接,一切正常:

left.join(right, left("src") === right("loc"), "left_outer")

这会按预期返回一个带有联接的数据框

我实际上想要做的是在 col1 和 col2 上进行匹配,实际上试图返回以下内容

源代码 | dst | src_loc | src_name | dst_loc | dst_name -------------------------------------------------- - 乙 |一个 |乙 |巴黎 |一个 |伦敦 c |乙 |空 |空 |乙 |巴黎 一个 | c |一个 |伦敦 |空 |空值

如果我尝试按如下方式在数据帧上执行此操作,整个 Spark 作业就会失败,它不会出错,但它要么花费的时间太长,要么发生了我不明白的事情。

val dfjoin1 = left.join(right, left("src") === right("loc"), "left_outer")
dfjoin1.join(right, dfjoin1("dst") === right("loc"), "left_outer")

出于沮丧,我尝试从第二个相同的 hive 查询中创建一个新的数据帧,而不是重复使用正确的数据帧

以下工作,但对我来说似乎非常错误(不应该为相同的数据调用 hive 两次)

val right = hiveContext.sql(FROM .....)
val right2 = hiveContext.sql(FROM .....)

val dfjoin1 = left.join(right, left("src") === right("loc"), "left_outer")
dfjoin1.join(right2, dfjoin1("dst") === right2("loc"), "left_outer")

我遇到的 ext 问题是我想过滤已添加的列,为了论证,假设我想获取所有 src loc 名称为 Paris 的列。

dfjoin1.filter($"name" === "Paris")

由于列名不明确,此操作失败。我该如何解决这个问题?作为连接的一部分,我可以轻松地为列添加名称前缀吗?

【问题讨论】:

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


    【解决方案1】:

    不确定 - 但我认为失败的原因也是列不明确 - 当您比较 dfjoin1("dst") === right("loc") 时,您实际上可能是在比较 dst 与由先前连接操作连接的 loc 列。

    换句话说,我相信您的两个问题都可以通过更准确的列命名来解决,这将确保没有歧义。实现这一点(并获得所需的输出模式)的更简单方法是在每次连接后重命名列:

    val result = left
      .join(right, $"src" === $"loc", "left_outer")
      .withColumnRenamed("loc", "src_loc")
      .withColumnRenamed("name", "src_name")
      .join(right, $"dst" === $"loc", "left_outer") // "loc" is now non-ambiguous, because we renamed left's "loc"
      .withColumnRenamed("loc", "dst_loc")
      .withColumnRenamed("name", "dst_name")
    
    result.show()
    // +---+---+-------+--------+-------+--------+
    // |src|dst|src_loc|src_name|dst_loc|dst_name|
    // +---+---+-------+--------+-------+--------+
    // |  b|  a|      b|   Paris|      a|  London|
    // |  c|  b|   null|    null|      b|   Paris|
    // |  a|  c|      a|  London|   null|    null|
    // +---+---+-------+--------+-------+--------+
    

    另一种方法可以在使用之前使用DataFrame.as(String) 来命名正确的数据框,每次使用不同的名称。结果略有不同,但仍然可用:

    left
      .join(right.as("src"), $"src" === $"src.loc", "left_outer")
      .join(right.as("dst"), $"dst" === $"dst.loc", "left_outer")
      .show()
    
    // +---+---+----+------+----+------+
    // |src|dst| loc|  name| loc|  name|
    // +---+---+----+------+----+------+
    // |  b|  a|   b| Paris|   a|London|
    // |  c|  b|null|  null|   b| Paris|
    // |  a|  c|   a|London|null|  null|
    // +---+---+----+------+----+------+
    

    架构显示了locname 具有相同名称的两列,但实际上可以使用相关前缀来引用它们,例如src.namedst.loc

    【讨论】:

    • 是否有任何方法可以为连接中的所有列添加前缀,或者我需要在每列上执行 withColumnRenamed 吗?我在右表中有大约 30 列,这意味着有大约 60 条重命名语句,似乎有点矫枉过正,但我​​想可能是必要的。
    • 我想知道这是否真的可以解决“整个 Spark 工作失败,它不会出错,但它需要的时间太长或发生了什么”问题 - 如果没有,请评论!
    • 会的,我可以在星期一对完整的数据集进行测试。左表大约有 10 亿个结果,但有时可能更多。右表最多几千条记录。
    • 我在下面添加了一个答案,并对 withColumn 变体稍作修改。我已经在本地数据集上进行了测试,似乎工作正常,但还没有机会在大型数据集上进行测试。
    【解决方案2】:

    进一步提到 Tzach Zohar,如 cmets 中所述,如果您有很多列,重命名它们会变得非常难看。为了解决这个问题,您可以使用表模式来获取列的名称并为所有列添加一个名称,如下所示:

    var tmp = left.join(right,$"src" === $"loc", "left_outer")
    
    right.schema.fields.foreach { x => tmp = tmp.withColumnRenamed(x.name, "src_" + x.name) }
    
    tmp = tmp.join(right,$"dst" === $"loc", "left_outer")
    
    right.schema.fields.foreach { x => tmp = tmp.withColumnRenamed(x.name, "dst_" + x.name) }
    
    // +---+---+-------+--------+-------+--------+
    // |src|dst|src_loc|src_name|dst_loc|dst_name|
    // +---+---+-------+--------+-------+--------+
    // |  b|  a|      b|   Paris|      a|  London|
    // |  c|  b|   null|    null|      b|   Paris|
    // |  a|  c|      a|  London|   null|    null|
    // +---+---+-------+--------+-------+--------+
    

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 2011-07-21
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-12-05
      • 2017-10-20
      • 2013-07-06
      • 2012-03-09
      相关资源
      最近更新 更多