【问题标题】:How to flatten tuple created using zip transformation in PySpark如何展平在 PySpark 中使用 zip 转换创建的元组
【发布时间】:2015-11-10 13:09:23
【问题描述】:

我有两个 RDD - RDD1 和 RDD2,结构如下:

RDD1:

[(u'abc', 1.0), (u'cde', 1.0),....]

RDD2:

[3.0, 0.0,....]

现在我想形成第三个 RDD,它的值来自上述两个 RDD 的每个索引。所以上面的输出应该变成:

RDD3:

[(u'abc', 1.0,3.0), (u'cde', 1.0,0.0),....]

如您所见,来自 RDD2 的值已添加到 RDD1 的元组中。我怎样才能做到这一点?我试图做RDD3 = RDD1.map(lambda x:x).zip(RDD2),但它产生了这个输出 - [((u'abc', 1.0),3.0), ((u'cde', 1.0),0.0),....] 这不是我想要的,因为你可以看到 () 的 RDD1 和 RDD2 的值之间存在分隔。

注意:我的 RDD1 是使用 - RDD1 = data.map(lambda x:(x[0])).zip(val)

【问题讨论】:

  • 使用后续映射创建所需的元组。
  • @Marcin 我做了RDD3 = RDD1.map(lambda x:x).zip(RDD2) 就像我在上面的帖子中提到的那样,但这并没有产生所需的输出
  • 后续;您的地图也是无操作的,因为它应用了身份转换。

标签: python apache-spark ipython pyspark rdd


【解决方案1】:

您可以在压缩后简单地重塑数据:

rdd1 = sc.parallelize([(u'abc', 1.0), (u'cde', 1.0)])
rdd2 = sc.parallelize([3.0, 0.0])

rdd1.zip(rdd2).map(lambda t: (t[0][0], t[0][1], t[1]))

在 Python 2 中可以使用:

rdd1.zip(rdd2).map(lambda ((x1, x2), y): (x1, x2, y))

但 Python 3 不再支持它。

如果您要使用索引提取更多值可能会很乏味

lambda t: (t[0][0], t[0][1], t[0][2], ..., t[1]))

所以你可以尝试这样的事情:

lambda t: tuple(list(t[0]) + [t[1]])

或实施更复杂的解决方案,例如:Flatten (an irregular) list of lists

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2013-09-01
    • 1970-01-01
    • 2020-05-15
    • 1970-01-01
    • 2023-01-21
    • 2016-05-02
    • 1970-01-01
    相关资源
    最近更新 更多