【发布时间】:2021-06-23 21:36:09
【问题描述】:
情况如下:我有一个时间序列数据框,它由一个对序列进行排序的索引列组成;和一列这样的离散值:
id value
0 A
1 A
2 B
3 C
4 A
5 A
6 A
7 B
我现在想减少所有连续重复,使其看起来像这样:
id value
0 A
2 B
3 C
4 A
7 B
我想出了一个窗口并使用lag()、when() 并在之后进行过滤。问题是窗口需要特定的分区列。然而,我想要的是首先删除每个分区中的连续行,然后检查分区边界(因为窗口每个分区都有效,所以分区边界上的连续行仍然存在)。
df_with_block = df.withColumn(
"block", (col("id") / df.rdd.getNumPartitions()).cast("int"))
window = Window.partitionBy("block").orderBy("id")
get_last = when(lag("value", 1).over(window) == col("value"), False).otherwise(True)
reduced_df = unificated_with_block.withColumn("reduced",get_last)
.where(col("reduced")).drop("reduced")
在第一行中,我通过整数除以 id 创建了一个具有均匀分布分区的新数据帧。 get_last 然后包含布尔值,具体取决于当前行是否等于前面的行。 reduce_df 然后过滤掉重复项。
现在的问题是分区边界:
id value
0 A
2 B
3 C
4 A
6 A
7 B
如您所见,id=6 的行没有被删除,因为它是在不同的分区中处理的。我正在考虑不同的想法来解决这个问题:
- 使用
coalesce()合并分区并再次过滤? - 想办法从下一个分区访问第一个值
- 使用 RDD 而不是 Dataframe 来完成所有这些操作
- 更改我的分区函数,使其不会切入重复项所在的位置(如何?)
我很好奇这是怎么解决的。
【问题讨论】:
-
在一定程度上是这样,但我希望它并行执行,因此 Spark 不必将所有数据移动到单个分区。
标签: pyspark apache-spark-sql time-series partitioning pyspark-sql