RDD具有一个有意义的(与存储模型强加的一些随机顺序相反),如果它被处理sortBy(),则如本回复中所解释的那样.
现在,哪些操作保留了该订单?
例如,是否保证(之后a.sortBy())
a.map(f).zip(a) ===
a.map(x => (f(x),x))
Run Code Online (Sandbox Code Playgroud)
怎么样
a.filter(f).map(g) ===
a.map(x => (x,g(x))).filter(f(_._1)).map(_._2)
Run Code Online (Sandbox Code Playgroud)
关于什么
a.filter(f).flatMap(g) ===
a.flatMap(x => g(x).map((x,_))).filter(f(_._1)).map(_._2)
Run Code Online (Sandbox Code Playgroud)
这里"平等" ===被理解为"功能等同",即,没有办法使用用户级操作来区分结果(即,没有读取日志和c).
Dan*_*bos 55
除了那些明确没有的操作外,所有操作都会保留订单.订购总是"有意义的",而不仅仅是在一个之后sortBy.例如,如果您读取文件(sc.textFile),则RDD的行将按它们在文件中的顺序排列.
如果没有试图给一个完整的清单,map,filter,flatMap,和coalesce(同shuffle=false)做维护秩序.sortBy,partitionBy,join不保留订单.
原因是大多数RDD操作Iterator在分区内部工作.所以map或者filter只是没有办法搞砸订单.您可以查看代码以便自己查看.
现在你可能会问:如果我有一个的RDD HashPartitioner.当我map用来更换钥匙时会发生什么?好吧,它们将保持原位,现在RDD没有被密钥分区.您可以使用partitionByshuffle恢复分区.
在 Spark 2.0.0+ 中,coalesce不保证合并期间的分区顺序。DefaultPartitionCoalescer具有基于分区局部性的优化算法。当分区包含有关其位置的信息时,DefaultPartitionCoalescer尝试合并同一主机上的分区。仅当没有位置信息时,它才会根据索引分割分区并保留分区顺序。
更新:
如果您从文件(例如 parquet)加载 DataFrame,Spark 在计划文件拆分时会破坏顺序。如果您使用它,您可以在DataSourceScanExec.scala#L629或新的 Spark 3.x FileScan#L152中看到它。spark.sql.files.maxPartitionBytes它只是按大小和小于最后一个分区的分割对分区进行排序。
因此,如果您需要从文件加载排序的数据集,您需要实现自己的阅读器。
| 归档时间: |
|
| 查看次数: |
17109 次 |
| 最近记录: |