【发布时间】:2021-07-19 21:23:47
【问题描述】:
我正在尝试获取出现在此列表中的前两个计数,按它们出现的最早 log_date。
state count log_date
GU 7402 2021-07-19
GU 7402 2021-07-18
GU 7402 2021-07-17
GU 7402 2021-07-16
GU 7397 2021-07-15
GU 7397 2021-07-14
GU 7397 2021-07-13
GU 7402 2021-07-12
GU 7402 2021-07-11
GU 7225 2021-07-10
GU 7225 2021-07-10
在这种情况下,我的预期输出是:
GU 7402 2021-07-16
GU 7397 2021-07-13
这就是我所做的工作,但在一些边缘情况下,计数可能会先下降然后再上升,如上面的示例所示。此代码返回 2021-07-11 作为 count=7402 的最早日期。
df = df.withColumnRenamed("count", "case_count")
df2 = df.groupBy(
"state", "case_count"
).agg(
F.min("log_time").alias("earliest_date")
)
df2 = df2.select("state", "case_count", "earliest_date").distinct()
df = df2.withColumn(
"last_import_date",
F.max("earliest_date").over(Window.partitionBy("state"))
).withColumn(
"max_count",
F.min(
F.when(
F.col("earliest_date") == F.col("last_import_date"),
F.col("case_count")
)
).over(Window.partitionBy("state"))
)
df = df.select("state", "max_count", "last_import_date").distinct()
我认为我需要做的是根据 state 和 log_date(desc) 排序选择前两个计数,然后获取每个计数的最小 log_date。我认为 rank() 可能会在这里通过为每个计数获取最高排名来工作,但我对如何将它应用于这种情况感到困惑。无论我尝试什么,我都无法摆脱最后两条 count=7402 记录。也许我忽略了一种更简单的方法?
df = df.withColumnRenamed(
"count", "case_count"
)
df = df.withColumn(
"rank",
F.rank().over(
Window.partitionBy(
"state", "case_count"
).orderBy(
F.col("state").asc(),
F.col("log_date").desc()
)
)
).orderBy(
F.col("log_date").desc(),
F.col("state").asc(),
F.col("rank").desc()
)
# output
state count log_date rank
GU 7402 2021-07-19 1
GU 7402 2021-07-18 2
GU 7402 2021-07-17 3
GU 7402 2021-07-16 4
GU 7397 2021-07-15 1
GU 7397 2021-07-14 2
GU 7397 2021-07-13 3
GU 7402 2021-07-12 5
GU 7402 2021-07-11 6
【问题讨论】:
标签: python dataframe apache-spark pyspark