【发布时间】:2018-02-09 01:34:30
【问题描述】:
我有两个数据框 df1 和 ip2Country。
df1 包含 IP 地址,我正在尝试将 IP 地址映射到 geolocation 信息,例如 longitude 和 latitude,它们是 @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