【问题标题】:Spark: Faster way to join two dataframe?Spark:加入两个数据框的更快方法?
【发布时间】:2018-02-09 01:34:30
【问题描述】:

我有两个数据框 df1ip2Countrydf1 包含 IP 地址,我正在尝试将 IP 地址映射到 geolocation 信息,例如 longitudelatitude,它们是 @987654325 中的列@。

我将它作为 Spark 提交作业运行,但即使 df1 只有不到 2500 行,操作也需要很长时间。

我的代码:

val agg =df1.join(ip2Country, ip2Country("network_start_int")=df1("sint") , “内在”) .select($"src_ip" ,$"country_name".alias("scountry") ,$"iso_3".alias("scode") ,$"经度".alias("slong") ,$"纬度".alias("slat") ,$"dst_ip",$"dint",$"count") .filter($"slong".isNotNull) val agg1 =agg.join(ip2Country, ip2Country("network_start_int")=agg("dint") , “内在”) .select($"src_ip",$"scountry" ,$"scode",$"slong" ,$"slat",$"dst_ip" ,$"country_name".alias("dcountry") ,$"iso_3".alias("dcode") ,$"经度".alias("dlong") ,$"纬度".alias("dlat"),$"count") .filter($"dlong".isNotNull)

有没有其他方法可以加入这两个表?还是我做错了?

【问题讨论】:

  • agg 或 agg1 哪个花费更多时间?
  • 实际上两者都需要很长时间。当我打印 sys.time 时,agg1 需要更长的时间
  • Largedf.join(broadcast(smalldf)) 将在广播提示框架的情况下工作
  • 所以就像 ip2Country.join(broadcast(df1),...) ?
  • 是的,请看我的回答here,它将清楚地解释为什么它会更好。以很好的方式解释了更多细节。如果你喜欢,请投票。谢谢!

标签: scala apache-spark


【解决方案1】:

如果您有一个大数据框需要与一个小数据框连接 - 广播连接非常有效。在这里阅读:Broadcast Joins (aka Map-Side Joins)

bigdf.join(broadcast(smalldf))

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-10-20
    • 2021-12-03
    • 2020-10-28
    • 2017-07-16
    • 2018-09-13
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多