【问题标题】:Apache Spark sort partition by user ID and write each partition to CSVApache Spark 按用户 ID 对分区进行排序并将每个分区写入 CSV
【发布时间】:2017-01-23 19:05:53
【问题描述】:

我有一个使用 Spark 似乎相对简单的用例,但似乎无法找到可靠的方法来解决这个问题。

我有一个数据集,其中包含各种用户的时间序列数据。我要做的就是:

  • 按用户 ID 分区此数据集
  • 对每个用户的时间序列数据进行排序,届时这些数据应该包含在各个分区中,
  • 将每个分区写入单个 CSV 文件。最后,我希望每个用户 ID 有 1 个 CSV 文件。

我尝试使用以下代码 sn-p,但最终得到了令人惊讶的结果。我最终每个用户 ID 有 1 个 csv 文件,并且一些用户的时间序列数据最终得到了排序,但很多其他用户没有排序。

# repr(ds) = DataFrame[userId: string, timestamp: string, c1: float, c2: float, c3: float, ...]
ds = load_dataset(user_dataset_path)
ds.repartition("userId")
    .sortWithinPartitions("timestamp")
    .write
    .partitionBy("userId")
    .option("header", "true")
    .csv(output_path)

我不清楚为什么会发生这种情况,也不完全确定该怎么做。我也不确定这是否是 Spark 中的潜在错误。

我将 Spark 2.0.2 与 Python 2.7.12 一起使用。任何建议将不胜感激!

【问题讨论】:

标签: python sorting apache-spark pyspark


【解决方案1】:

以下代码适用于我(此处以 Scala 显示,但在 Python 中类似)。

我得到每个用户名的一个文件,输出文件中的行按时间戳排序值。

testDF
  .select( $"username", $"timestamp", $"activity" )
  .repartition(col("username"))
  .sortWithinPartitions(col("username"),col("timestamp")) // <-- both here
  .write
  .partitionBy("username")
  .mode(SaveMode.Overwrite)
  .option("header", "true")
  .option("delimiter", ",")
  .csv(folder + "/useractivity")

导入的东西是both将用户名和时间戳列作为sortWithinPartitions的参数。

这是其中一个输出文件的外观(我使用一个简单的整数作为时间戳):

timestamp,activity
345,login
402,upload
515,download
600,logout

【讨论】:

    猜你喜欢
    • 2022-10-13
    • 2017-04-03
    • 2021-02-05
    • 2015-07-07
    • 2020-08-24
    • 1970-01-01
    • 2015-01-01
    • 2016-02-10
    • 1970-01-01
    相关资源
    最近更新 更多