【问题标题】:Does writing a dataframe to HDFS affect its sorting将数据帧写入 HDFS 是否会影响其排序
【发布时间】:2017-07-03 22:37:41
【问题描述】:

我在多节点环境(一个主节点和两个从节点)上的 apache spark 上运行代码,在该环境中我正在操作数据帧,然后对其执行逻辑回归。在这两者之间,我还写出了临时转换的文件。我目睹了一个特殊的观察结果(是的,我已经仔细检查和三次检查了),我无法解释并想确认这可能是因为我的代码还是可能有其他因素在起作用。

我有一个类似

的数据框

df

uid rank text
a   1    najn
b   2    dak
c   1    kksa
c   3    alkw
b   1    bdsj
c   2    asma

我用下面的代码排序

sdf = df.orderBy("uid", "rank")
sdf.show()

uid rank text
a   1    najn
b   1    bdsj
b   2    dak
c   1    kksa
c   2    asma
c   3    alkw

并使用将转换后的df写入HDFS

sdf.repartition(1)
  .write.format("com.databricks.spark.csv")
  .option("header", "true")
  .save("/someLocation")

现在当我再次尝试查看数据时,它似乎失去了排序

sdf.show() 
uid rank text
a   1    najn
c   2    asma
b   2    dak
c   1    kksa
c   3    alkw
b   1    bdsj

当我跳过编写代码时,它工作正常。

如果这可能是一个有效的案例,任何人都有任何指针,我们可以做一些事情来解决它。

附:我尝试了编写代码的各种变体,增加分区数量,完全删除分区并将其保存为其他格式。

【问题讨论】:

  • repartition 打乱所有数据删除和先前的顺序。否则,应以像这样的简单输出格式保留顺序。

标签: apache-spark dataframe hdfs apache-spark-sql spark-dataframe


【解决方案1】:

问题不在于写入 HDFS,而是在 cmets by zero323 中所述的重新分区。

如果您打算将所有内容写到一个文件中,您应该这样做:

sdf.coalesce(1).orderBy("uid", "rank").write...

coalesce 避免了重新分区(它只是一个接一个地复制分区,而不是通过哈希对所有内容进行混洗),这意味着您的数据仍将在原始分区中排序,因此排序速度更快(当然,您总是会丢失原始订单,因为它在这里没有多大帮助)。

请注意,这是不可扩展的,因为您将所有内容都拉到一个分区中。如果没有任何重新分区就出错了,您将根据 sdf 的原始分区数获得许多文件。每个文件都将在内部排序,以便您轻松组合它们。

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2017-11-01
    • 2017-10-30
    • 1970-01-01
    • 2018-12-25
    • 1970-01-01
    • 2021-04-24
    • 2021-04-01
    • 1970-01-01
    相关资源
    最近更新 更多