【问题标题】:Iterative filter in spark doesn't seem to workspark中的迭代过滤器似乎不起作用
【发布时间】:2021-01-20 10:53:52
【问题描述】:

我正在尝试逐个删除 RDD 的元素,但这不起作用,因为元素重新出现。

这是我的代码的一部分:

rdd = spark.sparkContext.parallelize([0,1,2,3,4])
for i in range(5):
    rdd=rdd.filter(lambda x:x!=i)
print(rdd.collect())
[0, 1, 2, 3]

所以似乎只有最后一个过滤器是“记住”。我在想,在这个循环之后,rdd 会是空的。

但是,我不明白为什么,因为每次我将通过过滤器获得的新 rdd 保存在“rdd”中,所以它不应该保留所有转换吗?如果没有,我该怎么办?

感谢您指出我错在哪里!

【问题讨论】:

  • 因为rdd 变量在每个循环中都被新值替换。比如rdd = 1'filter 然后rdd=2 filter 等等。由于您在循环之外打印,因此您只会看到最新的值。\
  • @venky__ 是的,但是 rdd 不应该将它的元素一个接一个地删除吗?在第一个循环之后,应该只有 [1,2,3];在第二个之后不应该只是 [2,3] 吗?为什么/何时返回 0?
  • 你应该包含代码你的问题。仅使用图像来显示您获得的输出,或者您可以在问题中编写的任何内容,否则无法复制粘贴代码来测试错误。您还应该包括一个输入示例(可复制粘贴)和所需的输出。通过这种方式,您将有更高的机会获得相关答案。
  • 另见stackoverflow.com/questions/41666977/…,它解释了过滤器的工作原理

标签: python apache-spark pyspark rdd


【解决方案1】:

结果实际上是正确的 - 这不是 Spark 的错误。请注意,lambda 函数定义为x != ii 没有代入 lambda 函数。所以在 for 循环的每次迭代中,RDD 看起来像

rdd
rdd.filter(lambda x: x != i)
rdd.filter(lambda x: x != i).filter(lambda x: x != i)
rdd.filter(lambda x: x != i).filter(lambda x: x != i).filter(lambda x: x != i)

等等

由于过滤器都是相同的,并且它们将被替换为i 的最新值,因此在每次 for 循环迭代中只过滤掉一项。

为避免这种情况,您可以使用部分函数来确保将i 替换到函数中:

from functools import partial
 
rdd = spark.sparkContext.parallelize([0,1,2,3,4])
for i in range(5):
    rdd = rdd.filter(partial(lambda x, i: x != i, i))

print(rdd.collect())

或者你可以使用reduce:

from functools import reduce

rdd = spark.sparkContext.parallelize([0,1,2])
rdd = reduce(lambda r, i: r.filter(lambda x: x != i), range(3), rdd)
print(rdd.collect())

【讨论】:

    猜你喜欢
    • 2012-06-11
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-12-10
    • 2017-07-31
    • 1970-01-01
    • 2015-11-14
    • 1970-01-01
    相关资源
    最近更新 更多