【问题标题】:How to pass dataset column value to a function while using spark filter with scala?如何在使用带有scala的火花过滤器时将数据集列值传递给函数?
【发布时间】:2018-04-16 21:12:37
【问题描述】:

我有一个包含用户 ID 和操作类型的操作数组

+-------+-------+
|user_id|   type|
+-------+-------+
|     11| SEARCH|
+-------+-------+
|     11| DETAIL|
+-------+-------+
|     12| SEARCH|
+-------+-------+

我想过滤属于至少有一个搜索操作的用户的操作。

所以我创建了一个具有搜索操作的用户 ID 的布隆过滤器。

然后我尝试根据布隆过滤器的用户状态过滤所有操作

val df = spark.read...
val searchers = df.filter($"type" === "SEARCH").select("user_id").distinct.as[String].collect
val bloomFilter = BloomFilter.create(100)
searchers.foreach(bloomFilter.putString(_))
df.filter(bloomFilter.mightContainString($"user_id"))

但是代码给出了异常

type mismatch;
found   : org.apache.spark.sql.ColumnName
required: String

请告诉我如何将列值传递给 BloomFilter.mightContainString 方法?

【问题讨论】:

    标签: scala apache-spark bloom-filter


    【解决方案1】:

    创建过滤器:

    val expectedNumItems: Long = ???
    val fpp: Double = ???
    val f = df.stat.bloomFilter("user_id", expectedNumItems, fpp)
    

    使用udf 进行过滤:

    import org.apache.spark.sql.functions.udf
    
    val mightContain = udf((s: String) => f.mightContain(s))
    df.filter(mightContain($"user_id"))
    

    如果您当前的 Bloom 过滤器实现是可序列化的,您应该能够以相同的方式使用它,但如果数据大到足以证明 Bloom 过滤器的合理性,您应该避免收集。

    【讨论】:

      【解决方案2】:

      你可以这样做,

      val sparkSession = ???
      val sc = sparkSession.sparkContext
      
      val bloomFilter = BloomFilter.create(100)
      
      val df = ???
      
      val searchers = df.filter($"type" === "SEARCH").select("user_id").distinct.as[String].collect
      

      在这一点上,我会提到collect 不是一个好主意的事实。接下来你可以做类似的事情。

      import org.apache.spark.sql.functions.udf
      val bbFilter = sc.broadcast(bloomFilter)
      
      val filterUDF = udf((s: String) => bbFilter.value.mightContainString(s))
      
      df.filter(filterUDF($"user_id"))
      

      如果bloomFilter实例是可序列化的,你可以移除广播。

      希望这会有所帮助,干杯。

      【讨论】:

      • 不太确定为什么这被否决了。 cmets 将有助于提高答案的质量。
      • 谢谢@Chitral,它有效。关于为什么我需要一个包装 udf 函数而不是直接调用 mightContainString 的任何评论?
      • 是的,像 filter、select、groupBy、where 等 spark sql DSL 接收 Column 对象,这就是需要 UDF 包装器的原因。如果有帮助,请点赞。
      猜你喜欢
      • 2018-10-27
      • 1970-01-01
      • 2016-05-01
      • 1970-01-01
      • 2017-03-09
      • 1970-01-01
      • 1970-01-01
      • 2022-10-14
      • 2021-04-04
      相关资源
      最近更新 更多