【问题标题】:Joining two dataframes in Spark is blowing up在 Spark 中加入两个数据帧正在爆炸
【发布时间】:2018-09-04 17:26:17
【问题描述】:

我正在尝试在一个字段上加入两个数据框。为此,我必须首先确保该字段是唯一的。所以我的事件顺序是:

  1. 读入第一个数据帧
  2. 选择我要加入的字段(例如,field1),以及我要加入的另一个字段(field2
  3. .distinct

那么,对于第二张桌子..

  1. 读入第二个数据帧
  2. field1 上与第一个表进行左连接
  3. .distinct

我尝试运行我的脚本,但它花费的时间比它应该的要长。 为了调试它,我在连接前后的第一个表上添加了println 记录计数,结果如下:

在加入之前,记录数为904,326。之后,2,658,632。 所以我认为它正在爆炸,但不知道为什么。我认为这与选择两个字段后尝试仅使用一个“不同”有关..?

请帮忙!

代码如下:

    val ticketProduct = Source.fromArg(args, "f1").read
   .select($"INSTRUMENT_SK", $"TICKET_CODES_SK")
   .distinct

  val instrumentD = Source.fromArg(args, "f2").read
//    println("instrumentD count before join is " + instrumentD.count)
   .join(ticketProduct, Seq("INSTRUMENT_SK"), "leftouter")
//   .select($"SERIAL_NBR", $"TICKET_CODES_SK")
   .distinct        
    println("instrumentD count after join is " + instrumentD.count)

【问题讨论】:

  • 我们能看到代码,如果可能的话,还能看到数据集吗?
  • 另外,在加入它们之前尝试对每个将加入的数据帧进行区分
  • 这是我的代码主体:
  • val ticketProduct = Source.fromArg(args, "f1").read .select($"INSTRUMENT_SK", $"TICKET_CODES_SK") .distinct val instrumentD = Source.fromArg(args, "f2" ).read // println("instrumentD 加入前的计数是 " + instrumentD.count) .join(ticketProduct, Seq("INSTRUMENT_SK"), "leftouter") // .select($"SERIAL_NBR", $"TICKET_CODES_SK") .distinct println("连接后的instrumentD计数为" + instrumentD.count) // .writeToSource(Source.fromArg(args, "output")) } } }
  • 请在添加更多信息(代码等)时更新您的问题

标签: scala apache-spark dataframe join


【解决方案1】:

您遇到的问题是,通过调用 distinct 您只能删除 field1 field2 的值相同的行。 由于您加入 field1,您可能希望 field1 的值是唯一的。

您可以尝试以下方式,而不是调用distinct

dataframe1.groupBy($"field1").agg(org.apache.spark.sql.functions.array($"field2"))

这将产生一个数据框,其中 field1 列是唯一的,field2 的多个值聚合到一个数组中。

这同样适用于第二个数据帧。

举个例子:假设您有包含以下内容的数据框。

field1, field2
 1,       1
 1,       2

field1, field3
 1,      1
 1,      3

然后 distinct 对它们没有任何作用,因为行是不同的。 现在,如果您要加入 field1,您将获得以下信息。

 1, 1, 1
 1, 2, 1
 1, 1, 3
 1, 2, 3

相反的聚合会给出

field1, array_of_field2
 1,      [1,2]

field1, array_of field3
 1,      [1,3]

然后连接将产生以下数据帧。

1, [1,2], [1,3]

【讨论】:

  • 最终,我希望 field2 或 ticket_codes_sk 的值也不同,因为我将使用它执行另一个连接。那么我也应该对这个字段使用相同的方法吗?
  • 是的,第二个你也必须这样做。
  • 好吧好吧好吧..我想我还是不完全明白。您的解决方案是否要一起查看两个字段,获取 field1 的唯一值并将 field2 中的相应值分组到一个数组中?因此,当我在下一次连接中使用 field2 中的值时,它会具有由 agg 产生的这些值?我不认为这就是我想要的......我只想让第一个表中每个字段的值都是唯一的,这样以后的两个连接就不会爆炸......
  • 换句话说,我不想在获取唯一值时同时查看这两个字段,而是让它们分别不同..
  • 还有人建议,也许我可能想考虑一个联合,而不是一个连接,并且只取一个数据框的第一行?我不确定这是什么意思..
猜你喜欢
  • 1970-01-01
  • 2022-01-16
  • 1970-01-01
  • 2020-04-03
  • 2021-06-06
  • 1970-01-01
  • 2017-01-09
  • 2015-06-11
  • 2018-09-26
相关资源
最近更新 更多