【发布时间】:2021-01-19 22:00:29
【问题描述】:
我想从 Spark 读取一个 CSV 文件(小于 50MB)并执行一些连接和过滤操作。 CSV 文件中的行按某些标准排序(在本例中为Score)。我想将结果保存在保留原始行顺序的单个 CSV 文件中。
输入 CSV 文件:
Id, Score
5, 100
3, 99
6, 98
7, 95
经过一些join&filter操作后:
val data = spark.read.option("header", "true").csv("s3://some-bucket/some-dir/123.csv")
val results = data
.dropDuplicates($"some_col")
.filter(x => ...)
.join(anotherDataset, Seq("some_col"), "left_anti")
results.repartition(1).write.option("header", "true").csv("...")
预期输出:
Id, Score
5, 100
6, 98
(过滤掉ID 3和7)
由于 Spark 可能会将数据加载到多个分区中,我该如何保持原始顺序?
【问题讨论】:
-
我的理解是如果输入的CSV文件小到可以加载到一个分区(默认分区大小为64MB),所有操作都在一个分区内完成,并保持顺序。
-
“所有操作都在单个分区内完成,并保持顺序”并非所有运算符都保持顺序。特别是大多数连接、分组、窗口函数和 obv 排序即使没有随机播放也不保留顺序。
-
dropDuplicates,join,repartition他们都是wide transformation。由于它们需要在多个 Spark 节点中打乱数据,因此初始顺序不会保留。为确保您在wide transformation之后的初始订单,您需要根据您的案例中的分数重新排序您的数据集。这是一个相关的discussion。
标签: csv apache-spark