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).
由于 Spark 是分布式的,因此通常假设确定性结果是不安全的。您的示例采用 DataFrame 的“第一个”10,000 行。在这里,“第一”的含义是模糊的(因此是非确定性的)。这将取决于 Spark 的内部结构。例如,它可能是响应驱动程序的第一个分区。该分区可能会随着网络、数据位置等而改变。
即使你缓存了数据,我仍然不会依赖于每次都取回相同的数据,尽管我当然希望它比从磁盘读取更一致。
来自Spark 文档:
该
LIMIT子句用于约束语句返回的行数SELECT。一般来说,该子句与 结合使用ORDER BY以确保结果是确定性的。
因此,如果您希望调用.limit()具有确定性,则需要事先对行进行排序。但是有一个问题!如果您按每行不具有唯一值的列进行排序,则所谓的“捆绑”行(具有相同排序键值的行)将不会被确定性地排序,因此.limit()仍然是不确定的。
您有两种选择来解决此问题:
df.orderBy('someCol', 'rowId').limit(n)。您可以这样rowIddf = df.withColumn('rowId', func.monotonically_increasing_id())df.limit(n).cache(),以便至少该 limit 的结果不会因连续的操作调用而改变,否则会重新计算结果limit并弄乱结果。| 归档时间: |
|
| 查看次数: |
8732 次 |
| 最近记录: |