【问题标题】:How to preserve order of a DataFrame when writing it as CSV with partitioning by columns?将 DataFrame 写为按列分区的 CSV 时如何保留它的顺序?
【发布时间】:2019-01-28 21:40:57
【问题描述】:

我对@9​​87654321@ 的行进行排序并将其写入磁盘,如下所示:

df.
  orderBy("foo").
  write.
  partitionBy("bar", "moo").
  option("compression", "gzip").
  csv(outDir)

当我查看生成的 .csv.gz 文件时,它们的顺序没有保留。这是 Spark 的做法吗?将 DF 写入带有分区的磁盘时,有没有办法保持顺序?

编辑:更准确地说:不是 CSV 的顺序关闭,而是它们内部的顺序。假设我在df.orderBy 之后有如下内容(为简单起见,我现在只按一列分区):

foo | bar | baz
===============
  1 |   1 |   1
  1 |   2 |   2
  1 |   1 |   3
  2 |   3 |   4
  2 |   1 |   5
  3 |   2 |   6
  3 |   3 |   7
  4 |   2 |   9
  4 |   1 |  10

我希望它是这样的,例如对于文件夹bar=1中的文件:

part-00000-NNN.csv.gz:

1,1
1,3
2,5

part-00001-NNN.csv.gz:

3,8
4,10

但它是什么样的:

part-00000-NNN.csv.gz:

1,1
2,5
1,3

part-00001-NNN.csv.gz:

4,10
3,8

【问题讨论】:

  • 你使用的是哪个 spark 版本
  • 我使用的是 2.3.1
  • 这个周末我也要试试这个。有“桶”。但通常在阅读时你会得到一个拆分,并且需要始终重新分区以确保。哈希与 RangeBy ... ?有趣的东西会引起混乱,然后 Hive 过去兼容。
  • 任务的结果用于.net 应用程序,该应用程序将数据读回并使用 FILESTREAM 列将其放入 MSSQL 数据库(这是第三方的要求)。由于每个文件夹都是一个接一个地读取的,而且所有 CSV 中的条目总和不会超过 600000 个,因此我目前将所有内容读入内存并重新排序。对于这种特殊情况,这是可以容忍的,但如果能更多地了解正在发生的事情,那就太好了。

标签: scala apache-spark apache-spark-sql


【解决方案1】:

已经有一段时间了,但我再次目睹了这一点。我终于找到了一个解决方法。

假设,你的架构是这样的:

  • 时间:大整数
  • 频道:字符串
  • 值:双倍

如果你这样做:

df.sortBy("time").write.partitionBy("channel").csv("hdfs:///foo")

单个 part-* 文件中的时间戳会被扔掉。

如果你这样做:

df.sortBy("channel", "time").write.partitionBy("channel").csv("hdfs:///foo")

顺序正确。

我认为这与洗牌有关。因此,作为一种解决方法,我现在按我希望我的数据首先分区的列进行排序,然后按我希望在单个文件中对其进行排序的列。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2020-12-21
    • 1970-01-01
    • 1970-01-01
    • 2013-03-17
    • 2018-12-10
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多