【问题标题】:Spark - read single CSV file, process and write results to single CSV file while keep original row orderSpark - 读取单个 CSV 文件,处理并将结果写入单个 CSV 文件,同时保持原始行顺序
【发布时间】: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


【解决方案1】:

你需要做的是在你做任何改变记录顺序的操作之前附加一个带有 monotonically_increasing_id() 的列,比如 group-bys、join、distinct 等。这个函数可以帮助你重新创建记录的顺序一个分区。

"生成的ID保证单调递增且唯一,但不连续。当前实现将分区ID放在高31位,低33位表示每个分区内的记录号"

val data = spark.read.option("header", "true").csv("s3://some-bucket/some-dir/123.csv")
val results = data
  .withColumn("rowId",monotonically_increasing_id())
  .dropDuplicates($"some_col"). // this might need to be replaced with a window function.
  .filter(x => ...)
  .join(anotherDataset, Seq("some_col"), "left_anti")

results.repartition(1)
.orderBy("rowId")
.write.option("header", "true").csv("...")

请注意,由于某种原因,spark sql 不包含简单的内置函数来获取 spark 分区 id 或 spark 分区行号,但幸运的是 monotonically_increasing_id 做得很好。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2015-10-18
    • 2023-02-17
    • 1970-01-01
    • 2022-01-13
    相关资源
    最近更新 更多