哪些操作保留了RDD顺序?

sds*_*sds 48 apache-spark rdd

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恢复分区.

  • 据我所知,我可以确认重新分配确实*不*保持秩序.如果我运行`x = sc.textFile('somefile'); y = x.repartition(100); a = x.collect(); b = y.collect()`,然后`a == b`返回`False`. (6认同)
  • @mustachio:哎呀,谢谢!你是对的。`repartition` 使用 `shuffle=true` 调用 `coalesce`,所以很明显它会_shuffle_ RDD。我已经确定了名单。 (2认同)

Avs*_*riy 5

在 Spark 2.0.0+ 中,coalesce不保证合并期间的分区顺序。DefaultPartitionCoalescer具有基于分区局部性的优化算法。当分区包含有关其位置的信息时,DefaultPartitionCoalescer尝试合并同一主机上的分区。仅当没有位置信息时,它才会根据索引分割分区并保留分区顺序。

更新:

如果您从文件(例如 parquet)加载 DataFrame,Spark 在计划文件拆分时会破坏顺序。如果您使用它,您可以在DataSourceScanExec.scala#L629或新的 Spark 3.x FileScan#L152中看到它。spark.sql.files.maxPartitionBytes它只是按大小和小于最后一个分区的分割对分区进行排序。

因此,如果您需要从文件加载排序的数据集,您需要实现自己的阅读器。