【发布时间】:2020-04-08 22:39:29
【问题描述】:
我正在尝试尽可能优化地执行 isin 过滤器。有没有办法使用 Scala API 广播 collList?
编辑:我不是在寻找替代方案,我知道它们,但我需要 isin 以便我的 RelationProviders 会下推这些值。
val collList = collectedDf.map(_.getAs[String]("col1")).sortWith(_ < _)
//collList.size == 200.000
val retTable = df.filter(col("col1").isin(collList: _*))
我传递给“isin”方法的列表有多达 200.000 个独特的元素。
我知道这看起来不是最好的选择,并且加入听起来更好,但我需要将这些元素推入过滤器,在阅读 (我的存储是 Kudu,但它也适用于 HDFS+Parquet,基础数据太大,查询只能处理大约 1% 的数据),我已经测量了所有内容,它为我节省了大约 30 分钟的执行时间:)。另外,如果 isin 大于 200.000,我的方法已经很小心了。
我的问题是,我收到了一些 Spark“任务太大”(每个任务约 8mb)警告,一切正常,所以没什么大不了的,但我希望删除它们并进行优化。
我已经尝试过,它什么也没做,因为我仍然收到警告(因为广播的 var 在 Scala 中得到解析并传递给我猜的 vargargs):
val collList = collectedDf.map(_.getAs[String]("col1")).sortWith(_ < _)
val retTable = df.filter(col("col1").isin(sc.broadcast(collList).value: _*))
还有这个不能编译:
val collList = collectedDf.map(_.getAs[String]("col1")).sortWith(_ < _)
val retTable = df.filter(col("col1").isin(sc.broadcast(collList: _*).value))
而这个不起作用(任务太大仍然出现)
val broadcastedList=df.sparkSession.sparkContext.broadcast(collList.map(lit(_).expr))
val filterBroadcasted=In(col("col1").expr, collList.value)
val retTable = df.filter(new Column(filterBroadcasted))
关于如何广播这个变量的任何想法? (允许黑客攻击)。任何允许过滤器下推的 isin 替代方案也是有效的我见过有人在 PySpark 上这样做,但 API 不一样。
PS:不可能对存储进行更改,我知道分区(已经分区,但不是按该字段)等可能会有所帮助,但用户输入完全是随机的,并且数据被访问并更改了我的许多客户端。
【问题讨论】:
标签: scala apache-spark