为什么df.limit在Pyspark中不断变化?

Ger*_*nuk 10 apache-spark pyspark spark-dataframe

我创建了一些数据帧数据样本df与

rdd = df.limit(10000).rdd
Run Code Online (Sandbox Code Playgroud)

这个操作需要相当长的时间(实际上为什么?在10000行之后它不能短路?),所以我假设我现在有了一个新的RDD.

但是,当我现在工作时rdd,每次访问它时都会有不同的行.好像它重新重新采样一样.缓存RDD有点帮助,但肯定不是保存?

它背后的原因是什么?

更新:这是Spark 1.5.2的再现

from operator import add
from pyspark.sql import Row
rdd=sc.parallelize([Row(i=i) for i in range(1000000)],100)
rdd1=rdd.toDF().limit(1000).rdd
for _ in range(3):
    print(rdd1.map(lambda row:row.i).reduce(add))
Run Code Online (Sandbox Code Playgroud)

输出是

499500
19955500
49651500
Run Code Online (Sandbox Code Playgroud)

我很惊讶.rdd没有修复数据.

编辑:为了表明它比重​​新执行问题更棘手,这里是一个单一的操作,在Spark 2.0.0.2.5.0上产生不正确的结果

from pyspark.sql import Row
rdd=sc.parallelize([Row(i=i) for i in range(1000000)],200)
rdd1=rdd.toDF().limit(12345).rdd
rdd2=rdd1.map(lambda x:(x,x))
rdd2.join(rdd2).count()
# result is 10240 despite doing a self-join
Run Code Online (Sandbox Code Playgroud)

基本上,每当您使用limit结果时可能会出错.我并不是指"只是众多样本中的一个",而是非常不正确(因为在这种情况下结果应始终为12345).

san*_*ton 9

由于 Spark 是分布式的,因此通常假设确定性结果是不安全的。您的示例采用 DataFrame 的“第一个”10,000 行。在这里,“第一”的含义是模糊的(因此是非确定性的)。这将取决于 Spark 的内部结构。例如,它可能是响应驱动程序的第一个分区。该分区可能会随着网络、数据位置等而改变。

即使你缓存了数据,我仍然不会依赖于每次都取回相同的数据,尽管我当然希望它比从磁盘读取更一致。

  • 您能更具体地说“拨打一次电话”吗?您可以定义一次 RDD(例如 `rdd = df.limit(10000).rdd`)。但 Spark 的计算是惰性的,所以直到你调用像“rdd.first()”这样的东西时,计算才会发生。此时,spark 或多或少将其解释为“给我来自‘df’的 10000 个随机行中的第一个元素”。当您稍后调用 rdd.first() 时,您将启动另一次计算,这可能会给出与第一次不同的结果。 (2认同)
  • 我觉得还是有一个误区。不要认为 RDD 实际上已经实现了。也就是说,不要将“rdd1”视为 12,345 个整数。将其视为保证最多包含 12,345 个整数的列表的计算描述,但_这些整数是什么_没有具体性_。何时、如何或多久引用 RDD 并不重要,它只是计算的描述。当您多次引用“rdd1”时,您会要求两次计算的输出,并且不能保证一致性。 (2认同)

dsa*_*laj 7

来自Spark 文档:

该LIMIT子句用于约束语句返回的行数SELECT。一般来说,该子句与 结合使用ORDER BY以确保结果是确定性的。

因此,如果您希望调用.limit()具有确定性,则需要事先对行进行排序。但是有一个问题!如果您按每行不具有唯一值的列进行排序,则所谓的“捆绑”行(具有相同排序键值的行)将不会被确定性地排序,因此.limit()仍然是不确定的。

您有两种选择来解决此问题:

  • 确保在排序调用中包含唯一的行 ID。
    例如df.orderBy('someCol', 'rowId').limit(n)。您可以这样
    定义:rowId
    df = df.withColumn('rowId', func.monotonically_increasing_id())
  • 如果您只需要在单次运行中获得确定性结果,则可以简单地缓存 limit 的结果df.limit(n).cache(),以便至少该 limit 的结果不会因连续的操作调用而改变,否则会重新计算结果limit并弄乱结果。


ale*_*lov 6

Spark 是惰性的,因此您执行的每个操作都会重新计算 limit() 返回的数据。如果底层数据分布在多个分区,那么每次评估它时,limit 可能会从不同的分区中提取(即如果您的数据存储在 10 个 Parquet 文件中,第一个 limit 调用可能从文件 1 中提取,第二个来自文件 7,等等)。