【问题标题】:Is there an alternative to do iterative join in spark - scala是否有替代方法可以在 spark - scala 中进行迭代加入
【发布时间】:2020-03-01 23:31:57
【问题描述】:

用例是在给定列中找到最多 n 行(这些可以是 n 列),一旦你有了 n 个键,你就可以将它加入到原始数据集中以获取你需要的所有行

val df = Seq(("12", "Tom", "Hanks"), (“13”,“梅丽尔”,“斯特里普”), (“12”,“汤姆”,“哈代”), (“12”、“约翰”、“斯特里普”) ).toDF("年龄", "名字", "姓氏")

假设我想将每一列单独加入一个更大的演员数据集,该数据集包含上述所有三列。

val v1 = actors.join(df, Seq("id"), "inner")
val v2 =actors.join(df, Seq("firstname"), "inner")
val v3 =actors.join(df, Seq("lastname"), "inner")
val output = v1.union(v2).union(v3)

有什么方法可以不重复执行此操作吗?还因为要连接的列可以是动态的。例如,有时它只能是 id,或者只能是 id 和名字。

【问题讨论】:

  • 在以下 2 个答案中阅读我的 cmets

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


【解决方案1】:

你可以尝试不同的方法,这样你就可以实现它:

actors.join(df).where(
actors("id") === df("id") || 
actors("firstname") === df("firstname") || 
actors("lastname") === df("lastname")
)

对于 n 列,您可以尝试一下:

  val joinCols = Seq("id", "firstname", "lastname") // or actors.columns
  val condition = joinCols
    .map(s => (actors(s) === df(s)))
    .reduce((a, b) => a || b)

你会得到以下条件:

condition.explain(true)
(((a#7 = a#7) || (b#8 = b#8)) || (c#9 = c#9))

最后使用它:

   actors.join(df).where(condition)

【讨论】:

  • 啊,但在这种情况下,列名是动态的。但这与静态列完美搭配!
  • 那么就用val joinCols = actors.columns
  • 啊哈。我认为这会解决我的问题。将确认并接受答案。谢谢!
  • 我认为这确实解决了我的问题。运行测试以检查性能@chlebek 对广播和 udf 的性能有任何担忧吗?虽然广播 udf 方法我不知道如何解决基于动态列的过滤。
  • 具有非等连接条件将导致笛卡尔连接运行。这似乎不是最优的。 @OP 你的多线连接会执行得更好。
【解决方案2】:

我认为使用较小的数据集进行广播并使用 udf 来检查较大的数据集可以解决问题。我一直在考虑加入!

【讨论】:

    【解决方案3】:

    @chlebek 解决方案应该可以正常工作,如果您想重现初始逻辑,这是另一种方法:

    val cols = Seq("id", "firstname", "lastname")
    
    val final_df = cols.map{
         df.join(actors, Seq(_), "inner") 
    }
    .reduce(_ union _)
    
    

    首先我们为每列生成一个内连接,然后我们联合它们。

    【讨论】:

    • 所有连接都将单独运行。为了尽量减少数据混洗,actordf 应该是 persistedcheckpointed。 Union 是一种狭义的转换,不会产生任何洗牌。
    猜你喜欢
    • 2020-09-11
    • 2016-07-01
    • 1970-01-01
    • 2013-08-13
    • 1970-01-01
    • 2023-03-27
    • 1970-01-01
    • 2011-01-18
    • 1970-01-01
    相关资源
    最近更新 更多