【发布时间】:2019-05-10 11:09:07
【问题描述】:
我每天都在使用 PySpark 处理一个文件,以获取有关设备通过网络导航的信息。在每个月底,我想使用窗口功能来为每个设备进行导航。这是一个非常缓慢的处理,即使有很多节点,所以我正在寻找加快它的方法。
我的想法是对数据进行分区,但我有 20 亿个不同的键,所以 partitionBy 似乎不合适。即使bucketBy 也可能不是一个好的选择,因为我每天都会创建n 存储桶,因此不会附加文件,但每天都会创建x 个文件。
有人有解决办法吗?
所以这里是每天导出的示例(在每个 parquet 文件中,我们找到 9 个分区):
这是我们在每个月初启动的 partitionBy 查询(compute_visit_number 和 compute_session_number 是我在笔记本上创建的两个 udf):
【问题讨论】:
-
您能否添加一些示例数据和您尝试过且需要优化的代码?这将有助于我们了解您想要做什么。
-
我添加了更多截图和示例
标签: apache-spark pyspark apache-spark-sql databricks azure-databricks