【问题标题】:Spark and BloomFilter sharingSpark 和 BloomFilter 共享
【发布时间】:2017-04-24 10:41:35
【问题描述】:

我有一个巨大的 RDD(源),我需要从中创建一个 BloomFilter 数据,因此对用户数据的后续更新将只考虑真正的“差异”,没有重复。

看起来 BloomFilter 的大多数实现都是不可序列化的(虽然可以很容易地修复),但我想要稍微不同的工作流程:

  1. 处理每个分区并为每个分区创建一个适当的 BloomFilter 实例。对于每个 BloomFilter 对象 - 将其写入某个二进制文件。我实际上不知道如何处理整个分区 - RDD 上有 mapPartition 函数可用,但这希望我返回一个迭代器。也许我可以使用传递的迭代器,创建一个 BloomFilter 的实例,将其写入某个位置并将指向创建文件的链接作为 Iterator.singleton[PathToFile] 返回?
  2. 在主节​​点 - consume 该处理的结果(文件的路径列表),读取这些文件并在内存中聚合 BloomFilter。然后将响应写入二进制文件。

我不知道正确的方法:

  • 在传递给mapPartitions的函数中,在集群支持的FS中创建一个文件(可以是HDFS、S3N或本地文件)
  • 在第二阶段使用consume 读取文件的内容(当我有一个带有文件路径的RDD 时,我必须使用SparkContext 来读取它们 - 不知道怎么可能) .

谢谢!

【问题讨论】:

    标签: apache-spark bloom-filter


    【解决方案1】:

    breeze 实现不是最快的,但它带有通常的 Spark 依赖项,可以与 simple aggregate 一起使用:

    import breeze.util.BloomFilter
    
    // Adjust values to fit your case
    val numBuckets: Int = 100
    val numHashFunctions: Int = 30
    
    val rdd = sc.parallelize(Seq("a", "d", "f", "e", "g", "j", "z", "k"), 4)
    val bf = rdd.aggregate(new BloomFilter[String](numBuckets, numHashFunctions))(
      _ += _, _ |= _
    )
    
    bf.contains("a")
    
    Boolean = true
    
    bf.contains("n")
    
    Boolean = false
    

    在 Spark 2.0+ 中你可以使用DataFrameStatFunctions.bloomFilter:

    val df = rdd.toDF
    
    val expectedNumItems: Long = 1000 
    val fpp: Double = 0.005
    
    val sbf = df.stat.bloomFilter($"value", expectedNumItems, fpp)
    
    sbf.mightContain("a")
    
    Boolean = true
    
    sbf.mightContain("n")
    
    Boolean = false
    

    Algebird 实现也可以工作,并且可以与breeze 实现类似地使用。

    【讨论】:

    • 什么是 $"value" in val sbf = df.stat.bloomFilter($"value", expectedNumItems, fpp)
    • $"value"Column 类的对象 - 换句话说,它是您要在其上创建布隆过滤器的列。例如$ 等价于org.apache.spark.sql.functions.col
    • 在 spark 中使用时,您必须将 Bloom Filter 设置为 BroadCast 变量,否则最终会导致每个任务的开销过多。每个执行器广播一次广播变量,否则,bloomfilter 将每个任务传输一次(或为每个任务执行构建)。
    • @YoYo 除非数据在阶段之间重复使用或者您担心反序列化成本,否则没有充分的理由广播过滤器。
    • @YoYo 不是。引用编程指南 Spark 会自动广播每个阶段内任务所需的公共数据。以这种方式广播的数据以序列化的形式缓存起来,并在运行每个任务之前进行反序列化.
    猜你喜欢
    • 2015-04-15
    • 1970-01-01
    • 2012-02-20
    • 1970-01-01
    • 2016-11-07
    • 2015-08-14
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多