【问题标题】:PySpark Broadcast Variable JoinPySpark 广播变量连接
【发布时间】:2015-06-27 13:44:34
【问题描述】:

我正在执行联接,并且我的数据跨 100 多个节点。所以我有一个键/值的小列表,我正在与另一个键/值对加入。

我的列表如下所示:

[[1, 0], [2, 0], [3, 0], [4, 0], [5, 0], [6, 0], [7, 0], [8, 0], [9, 0], [10, 0], [11, 0], [16, 0], [18, 0], [19, 0], [20, 0], [21, 0], [22, 0], [23, 0], [24, 0], [25, 0], [26, 0], [27, 0], [28, 0], [29, 0], [36, 0], [37, 0], [38, 0], [39, 0], [40, 0], [41, 0], [42, 0], [44, 0], [46, 0]]

我有广播变量:

numB = sc.broadcast(numValuesKV)

当我加入时:

numRDD = columnRDD.join(numB.value)

我收到以下错误:

AttributeError: 'list' object has no attribute 'map'

【问题讨论】:

    标签: python apache-spark pyspark


    【解决方案1】:

    你正在广播一个列表,这绝对没问题。

    你需要做的是

    b=sc.broadcast(lst)
    rdd.map(lambda t: t if t[0] in b.value)
    

    这里的 t[0] 应该看起来像 [1,0] 等等。但我希望你明白了....

    【讨论】:

      【解决方案2】:

      你能不能试着把 numValuesKV 做成一个字典,看看它是否有效。

      【讨论】:

      • +1 我有同样的问题.. 我尝试将广播值转换为字典导致 TypeError: 'Broadcast' object is not iterable
      【解决方案3】:

      rdd.join(other) 意味着加入两个 RDD,因此它期望 other 是一个 RDD。要使用高效的“小表广播”连接技巧,您需要“手动”进行连接。在 Scala 中,它看起来像这样:

      rdd.mapPartitions{iter =>
          val valueMap = numB.value.toMap
          iter.map{case (k,v) => (k,(v,map(v))}
      }
      

      这会将使用广播值的连接以分布式方式应用于 RDD 的每个分区。

      PySpark 代码应该非常相似。

      【讨论】:

      • 谢谢,我会在 python 中尝试一下。所以我一直使用 join 作为一个非常低效的过滤器。我本质上做的是试图只保留该列表中的密钥。加入很昂贵,但我一直在尝试如何过滤掉键!=该列表没有成功。
      • 如果打算过滤,而不是 iter.map 使用 iter.filter(cond) 就完成了。
      猜你喜欢
      • 2016-03-07
      • 1970-01-01
      • 1970-01-01
      • 2020-06-28
      • 2015-01-13
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2021-06-10
      相关资源
      最近更新 更多