【问题标题】:spark: save ordered data to parquetspark:将有序数据保存到镶木地板
【发布时间】:2020-03-11 20:58:48
【问题描述】:

我有 30TB 的数据按日期和小时划分,每小时分成 300 个文件。我做了一些数据转换,然后希望数据按排序顺序排序和保存,以便 C++ 程序轻松摄取。我了解当您进行序列化时,排序仅在文件中是正确的。我希望通过更好地划分数据来规避这个问题。

我想同时按 sessionID 和时间戳排序。我不希望 sessionID 在不同文件之间拆分。如果我在 SessionID 上进行分区,我将有太多,所以我做一个模 N 来生成 N 个桶,旨在获得 1 个约 100-200MB 的数据桶:

df = df.withColumn("bucket", F.abs(F.col("sessionId")) % F.lit(50))

然后我按日期、时间和存储桶遣返,然后再进行排序

df = df.repartition(50,"dt","hr","bucket")
df = df.sortWithinPartitions("sessionId","timestamp")
df.write.option("compression","gzip").partitionBy("dt","hr","bucket").parquet(SAVE_PATH)

这会将数据保存到 dt/hr/bucket,每个存储桶中有 1 个文件,但排序丢失。如果我不创建存储桶和重新分区,那么我最终会得到 200 个文件,数据是有序的,但是 sessionIds 被拆分到多个文件中。

编辑: 问题似乎出在使用partitionBy("dt","hr","bucket") 保存时,它会随机对数据进行重新分区,因此不再对其进行排序。如果我在没有partitionBy 的情况下保存,那么我得到的正是我所期望的——N 个存储桶/分区的 N 个文件和 sessionIds 跨越一个文件,所有这些都正确排序。所以我有一个非火花黑客手动迭代所有日期+小时目录

如果您按列分区,排序,然后使用 partitionBy 写入同一列,那么您希望直接转储已排序的分区,而不是对数据进行随机重新洗牌,这似乎是一个错误。

【问题讨论】:

  • 尝试使用sortWithinPartitions(),例如:df.write.option("compression","gzip").partitionBy("dt","hr","bucket").sortWithinPartitions("sessionId", "timestamp").parquet(SAVE_PATH)
  • 这是我一直在尝试的事情,但我得到了这个错误:AttributeError: 'DataFrameWriter' object has no attribute 'sortWithinPartitions' 似乎在写出分区时,您几乎没有选择来控制文件大小、存储桶和排序。作为一个初学者,您可以订购 30TB 的数据似乎很令人困惑,然后在保存时您放弃了昂贵的订购
  • @syadav。 quoteinvestigator.com/2012/04/28/shorter-letter我做了一些编辑并改进了文本,但我喜欢提供足够的上下文,以便人们知道我打算实现什么,因为可能会有更全面的解决方案

标签: apache-spark pyspark sql-order-by parquet partition-by


【解决方案1】:

将分区列放在已排序的列列表中可能会奏效。

完整描述在这里 - https://stackoverflow.com/a/59161488/3061686

【讨论】:

    猜你喜欢
    • 2016-07-04
    • 1970-01-01
    • 2020-03-18
    • 1970-01-01
    • 2015-10-23
    • 2019-01-08
    • 1970-01-01
    • 2017-04-25
    • 2020-10-23
    相关资源
    最近更新 更多