【问题标题】:Filtering dataframe using hashmap使用 hashmap 过滤数据帧
【发布时间】:2018-12-30 23:32:51
【问题描述】:

我有一个存储值的哈希图

Map(862304021470656 -> List(0.0, 0.0, 0.0, 0.0, 1.540980096E9, 74.365111, 22.302669, 0.0),866561010400483 -> List(0.0, 1.0, 1.0, 2.0, 1.543622306E9, 78.0204, 10.005262, 56.0))

这是数据框

|             id|       lt|       ln|       evt|    lstevt|  s|  d|agl|chg| d1| d2| d3| d4|ebt|ibt|port| a1| a2| a3| a4|nos|dfrmd|
+---------------+---------+---------+----------+----------+---+---+---+---+---+---+---+---+---+---+----+---+---+---+---+---+-----+
|862304021470656|25.284158|82.435973|1540980095|1540980095|  0| 39|298|  0|  0|  1|  1|  2|  0|  5|  97| 12| -1| -1| 22|  0|    0|
|862304021470656|25.284158|82.435973|1540980105|1540980105|  0|  0|298|  0|  0|  1|  1|  2|  0|  5|  97| 12| -1| -1| 22|  0|    0|
|862304021470656|25.284724|82.434222|1540980155|1540980155| 14| 47|289|  0|  0|  1|  1|  2|  0|  5|  97| 11| -1| -1| 22|  0|    0|
|866561010400483|25.284858|82.433831|1544980165|1540980165| 12| 42|295|  0|  0|  1|  1|  2|  0|  5|  97| 12| -1| -1| 22|  0|    0|

我只想从数据帧中过滤这些值,比较 evt 列中列表的第 4 个索引,只选择 evt 值大于列表第 4 个索引值的行,映射中的键是数据帧的 id 列.

【问题讨论】:

  • 你的 hashmap 有多大?它是广播的好选择吗(例如,小于~300 MB)?那么在 hashmap 中没有匹配键的行呢?

标签: scala hashmap apache-spark-sql


【解决方案1】:

这是使用 UDF 获取 evt 值以进行比较的一种方法:

import org.apache.spark.sql.functions._

val df = Seq(
  (862304021470656L, 25.284158, 82.435973, 1540980095),
  (862304021470656L, 25.284158, 82.435973, 1540980105),
  (862304021470656L, 25.284724, 82.434222, 1540980155),
  (866561010400483L, 25.284858, 82.433831, 1544980165)
).toDF("id", "lt", "ln", "evt")

val listMap = Map(
  862304021470656L -> List(0.0, 0.0, 0.0, 0.0, 1.540980096E9, 74.365111, 22.302669, 0.0),
  866561010400483L -> List(0.0, 1.0, 1.0, 2.0, 1.543622306E9, 78.0204, 10.005262, 56.0)
)

def evtLimit(m: Map[Long, List[Double]], evtIdx: Int) = udf(
  (id: Long) => m.get(id) match {
      case Some(ls) => if (evtIdx < ls.size) ls(evtIdx) else Double.MaxValue
      case None => Double.MaxValue
    }
)

df.where($"evt" > evtLimit(listMap, 4)($"id")).show
// +---------------+---------+---------+----------+
// |             id|       lt|       ln|       evt|
// +---------------+---------+---------+----------+
// |862304021470656|25.284158|82.435973|1540980105|
// |862304021470656|25.284724|82.434222|1540980155|
// |866561010400483|25.284858|82.433831|1544980165|
// +---------------+---------+---------+----------+

请注意,如果提供的 Map 中的键不匹配或值无效,UDF 将返回 Double.MaxValue。这当然可以根据特定的业务需求进行修改。

【讨论】:

  • 您好,先生,请告诉我如何获取过滤后的数据帧,其中 hashmap 的键不在数据帧的 id 值中。
  • @experiment,如果我正确理解您的问题,您可以使用内置方法isin,如df.where(!$"id".isin(listMap.keys.toArray: _*))
【解决方案2】:

你可以用一个简单的 sql 得到这个:

import spark.implicits._
import org.apache.spark.sql.functions._
val df = ... //your main Dataframe
val map = Map(..your data here..).toDF("id", "list")
val join = df.join(map, "id").filter(length($"list") >= 5 /* <-- just in case */)
val res = join.filter($"evt" > $"list"(4))

【讨论】:

  • 感谢您的努力,但我想避免加入比较,因为这是一个繁重的操作。
  • 如果你的地图足够小,Spark会广播它,不会有shuffle。
  • 至少会有 10 万个数据,那可行吗?
  • @experiment 1lakh?你是说1kb?您可以通过spark.sql.autoBroadcastJoinThreshold 设置 Spark 广播阈值,默认为 10MB。
猜你喜欢
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 2023-01-09
  • 2017-03-16
  • 2021-04-10
  • 2018-02-20
  • 2017-03-08
  • 2017-10-21
相关资源
最近更新 更多