【问题标题】:PySpark / Spark Window Function First/ Last IssuePySpark / Spark 窗口功能第一期/最后一期
【发布时间】:2019-02-15 19:07:36
【问题描述】:

据我了解,Spark 中的 first/last 函数将检索每个分区的第一行/最后一行/我无法理解为什么 LAST 函数给出了不正确的结果。

这是我的代码。

AgeWindow = Window.partitionBy('Dept').orderBy('Age')
df1 = df1.withColumn('first(ID)', first('ID').over(AgeWindow))\
        .withColumn('last(ID)', last('ID').over(AgeWindow))           
df1.show()
+---+----------+---+--------+--------------------------+-------------------------+
|Age|      Dept| ID|    Name|first(ID)                 |last(ID)                |
+---+----------+---+--------+--------------------------+-------------------------+
| 38|  medicine|  4|   harry|                         4|                        4|
| 41|  medicine|  5|hermione|                         4|                        5|
| 55|  medicine|  7| gandalf|                         4|                        7|
| 15|technology|  6|  sirius|                         6|                        6|
| 49|technology|  9|     sam|                         6|                        9|
| 88|technology|  1|     sam|                         6|                        2|
| 88|technology|  2|     nik|                         6|                        2|
| 75|       mba|  8|   ginny|                         8|                       11|
| 75|       mba| 10|     sam|                         8|                       11|
| 75|       mba|  3|     ron|                         8|                       11|
| 75|       mba| 11|     ron|                         8|                       11|
+---+----------+---+--------+--------------------------+-------------------------+

【问题讨论】:

    标签: sql apache-spark pyspark apache-spark-sql window-functions


    【解决方案1】:

    没有错。你的窗口定义不是你想象的那样。

    如果您提供ORDER BY 子句,则默认框架为RANGE BETWEEN UNBOUNDED PRECEDING AND CURRENT ROW:

    from pyspark.sql.window import Window
    from pyspark.sql.functions import first, last
    
    w = Window.partitionBy('Dept').orderBy('Age')
    
    df = spark.createDataFrame(
        [(38, "medicine", 4), (41, "medicine", 5), (55, "medicine", 7)],
        ("Age", "Dept", "ID")
    )
    
    df.select(
        "*",
        first('ID').over(w).alias("first_id"), 
        last('ID').over(w).alias("last_id")
    ).explain()
    
    == Physical Plan ==
    Window [first(ID#24L, false) windowspecdefinition(Dept#23, Age#22L ASC NULLS FIRST, specifiedwindowframe(RangeFrame, unboundedpreceding$(), currentrow$())) AS first_id#38L, last(ID#24L, false) windowspecdefinition(Dept#23, Age#22L ASC NULLS FIRST, specifiedwindowframe(RangeFrame, unboundedpreceding$(), currentrow$())) AS last_id#40L], [Dept#23], [Age#22L ASC NULLS FIRST]
    +- *(1) Sort [Dept#23 ASC NULLS FIRST, Age#22L ASC NULLS FIRST], false, 0
       +- Exchange hashpartitioning(Dept#23, 200)
          +- Scan ExistingRDD[Age#22L,Dept#23,ID#24L]
    

    这意味着窗口函数永远不会向前看,帧中的最后一行是当前行。

    您应该将窗口重新定义为

    w_uf = (Window
       .partitionBy('Dept')
       .orderBy('Age')
       .rowsBetween(Window.unboundedPreceding, Window.unboundedFollowing))
    
    result = df.select(
        "*", 
        first('ID').over(w_uf).alias("first_id"),
        last('ID').over(w_uf).alias("last_id")
    )
    
    == Physical Plan ==
    Window [first(ID#24L, false) windowspecdefinition(Dept#23, Age#22L ASC NULLS FIRST, specifiedwindowframe(RowFrame, unboundedpreceding$(), unboundedfollowing$())) AS first_id#56L, last(ID#24L, false) windowspecdefinition(Dept#23, Age#22L ASC NULLS FIRST, specifiedwindowframe(RowFrame, unboundedpreceding$(), unboundedfollowing$())) AS last_id#58L], [Dept#23], [Age#22L ASC NULLS FIRST]
    +- *(1) Sort [Dept#23 ASC NULLS FIRST, Age#22L ASC NULLS FIRST], false, 0
       +- Exchange hashpartitioning(Dept#23, 200)
          +- Scan ExistingRDD[Age#22L,Dept#23,ID#24L]
    
    result.show()
    
    +---+--------+---+--------+-------+
    |Age|    Dept| ID|first_id|last_id|
    +---+--------+---+--------+-------+
    | 38|medicine|  4|       4|      7|
    | 41|medicine|  5|       4|      7|
    | 55|medicine|  7|       4|      7|
    +---+--------+---+--------+-------+
    

    【讨论】:

    • 谢谢。有效。我正在尝试其他功能,例如 lag、lead、cume_dist,但它们对同一窗口给出了错误。 Window Frame specifiedwindowframe(RowFrame, unboundedpreceding$(), unboundedfollowing$()) must match the required frame specifiedwindowframe(RowFrame, 1, 1);
    • 某些功能只能在非常特定的框架中应用。例如lead 和lag 由它们的偏移量严格定义。
    • 很高兴知道一个关于 pyspark 中窗口函数的已解决错误jira.apache.org/jira/browse/SPARK-24033。由于该错误,您应该事先检查您的 spark 版本。
    • 你能帮忙回答一下吗 - stackoverflow.com/questions/64004622/…
    猜你喜欢
    • 2021-11-17
    • 1970-01-01
    • 2021-07-08
    • 2021-05-06
    • 1970-01-01
    • 2012-05-29
    • 1970-01-01
    • 1970-01-01
    • 2016-01-17
    相关资源
    最近更新 更多