【问题标题】:Creating a large dictionary in pyspark在 pyspark 中创建一个大字典
【发布时间】:2014-08-22 04:49:12
【问题描述】:

我正在尝试使用 pyspark 解决以下问题。 我在 hdfs 上有一个文件,格式是查找表的转储。

key1, value1
key2, value2
...

我想将它加载到 pyspark 中的 python 字典中,并将其用于其他目的。所以我尝试这样做:

table = {}
def populateDict(line):
    (k,v) = line.split(",", 1)
    table[k] = v

kvfile = sc.textFile("pathtofile")
kvfile.foreach(populateDict)

我发现表变量没有被修改。那么,有没有办法在 spark 中创建一个大的内存哈希表?

【问题讨论】:

    标签: python apache-spark


    【解决方案1】:

    foreach 是一种分布式计算,因此您不能期望它修改仅在驱动程序中可见的数据结构。你想要的是。

    kv.map(line => { line.split(" ") match { 
        case Array(k,v) => (k,v)
        case _ => ("","")
    }.collectAsMap()
    

    这是在 scala 中,但你明白了,重要的函数是 collectAsMap(),它将地图返回给驱动程序。

    如果您的数据非常大,您可以使用 PairRDD 作为地图。首先映射到对

        kv.map(line => { line.split(" ") match { 
            case Array(k,v) => (k,v)
            case _ => ("","")
        }
    

    然后您可以使用rdd.lookup("key") 进行访问,它返回与键关联的一系列值,尽管这肯定不会像其他分布式 KV 存储那样高效,因为 spark 并不是为此而构建的。

    【讨论】:

    • 非常感谢。这是否意味着地图必须适合驱动程序的内存?还是仍在分发?
    • @Kamal 是的,它必须适合内存。你可以使用 pair rdd 作为查找表。也想到了一个可累积的解决方案,很快就会发布
    • 好的。我正在寻找火花中的分布式地图。看来不可能!
    • 谢谢!我会试一试
    • 你是不是漏掉了一个}?
    【解决方案2】:

    效率见:sortByKey() and lookup()

    查找(键):

    返回 RDD 中键 key 的值列表。如果 RDD 有一个已知的分区器,则此操作会有效地完成,只需搜索键映射到的分区。

    RDD 将通过 sortByKey() (see: OrderedRDD) 重新分区,并在lookup() 调用期间有效地搜索。在代码中,类似,

    kvfile = sc.textFile("pathtofile")
    sorted_kv = kvfile.flatMap(lambda x: x.split("," , 1)).sortByKey()
    
    sorted_kv.lookup('key1').take(10)
    

    既能作为 RDD 又能高效地完成任务。

    【讨论】:

      猜你喜欢
      • 2021-03-09
      • 1970-01-01
      • 2021-05-15
      • 1970-01-01
      • 2012-01-15
      • 2019-03-15
      • 1970-01-01
      • 2016-06-14
      • 1970-01-01
      相关资源
      最近更新 更多