【发布时间】: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 一起使用。任何建议将不胜感激!
【问题讨论】:
-
当然,我在 GitHub 上做了一个要点来详细说明我看到的问题。行为是确定性的。 Spark 和 Scala 版本也包含在 gist 中。 gist.github.com/igozali/d327a85646abe7ab10c2ae479bed431f
-
看来这实际上可能是一个错误。我已经提交了issues.apache.org/jira/browse/SPARK-19352,目前正在这个 GitHub PR 中进行处理:github.com/apache/spark/pull/16724
-
对于以后的人来说,似乎可以通过执行以下操作来实现所需的行为:``` ds.repartition("userId") .sortWithinPartitions("userId", "timestamp") .write .partitionBy ("userId") .option("header", "true") .csv(output_path) ``` 参见例如github.com/apache/spark/pull/16724#issuecomment-279190560 和 github.com/apache/spark/pull/16898
标签: python sorting apache-spark pyspark