【发布时间】:2017-04-24 10:41:35
【问题描述】:
我有一个巨大的 RDD(源),我需要从中创建一个 BloomFilter 数据,因此对用户数据的后续更新将只考虑真正的“差异”,没有重复。
看起来 BloomFilter 的大多数实现都是不可序列化的(虽然可以很容易地修复),但我想要稍微不同的工作流程:
- 处理每个分区并为每个分区创建一个适当的 BloomFilter 实例。对于每个 BloomFilter 对象 - 将其写入某个二进制文件。我实际上不知道如何处理整个分区 - RDD 上有
mapPartition函数可用,但这希望我返回一个迭代器。也许我可以使用传递的迭代器,创建一个 BloomFilter 的实例,将其写入某个位置并将指向创建文件的链接作为Iterator.singleton[PathToFile]返回? - 在主节点 -
consume该处理的结果(文件的路径列表),读取这些文件并在内存中聚合 BloomFilter。然后将响应写入二进制文件。
我不知道正确的方法:
- 在传递给
mapPartitions的函数中,在集群支持的FS中创建一个文件(可以是HDFS、S3N或本地文件) - 在第二阶段使用
consume读取文件的内容(当我有一个带有文件路径的RDD 时,我必须使用SparkContext来读取它们 - 不知道怎么可能) .
谢谢!
【问题讨论】: