【问题标题】:Why do Window functions fail with "Window function X does not take a frame specification"?为什么窗口函数会因“窗口函数 X 不采用帧规范”而失败?
【发布时间】:2015-11-29 08:27:24
【问题描述】:

我正在尝试在 pyspark 1.4.1 中使用 Spark 1.4 window functions

但大多是错误或意外结果。 这是一个我认为应该可行的非常简单的示例:

from pyspark.sql.window import Window
import pyspark.sql.functions as func

l = [(1,101),(2,202),(3,303),(4,404),(5,505)]
df = sqlContext.createDataFrame(l,["a","b"])

wSpec = Window.orderBy(df.a).rowsBetween(-1,1)

df.select(df.a, func.rank().over(wSpec).alias("rank"))  
    ==> Failure org.apache.spark.sql.AnalysisException: Window function rank does not take a frame specification.

df.select(df.a, func.lag(df.b,1).over(wSpec).alias("prev"), df.b, func.lead(df.b,1).over(wSpec).alias("next"))  
    ===>  org.apache.spark.sql.AnalysisException: Window function lag does not take a frame specification.;


wSpec = Window.orderBy(df.a)

df.select(df.a, func.rank().over(wSpec).alias("rank"))
    ===> org.apache.hadoop.hive.ql.exec.UDFArgumentTypeException: One or more arguments are expected.

df.select(df.a, func.lag(df.b,1).over(wSpec).alias("prev"), df.b, func.lead(df.b,1).over(wSpec).alias("next")).collect()

    [Row(a=1, prev=None, b=101, next=None), Row(a=2, prev=None, b=202, next=None), Row(a=3, prev=None, b=303, next=None)]

如您所见,如果我添加 rowsBetween 框架规范,rank()lag/lead() 窗口函数都不会识别它:“窗口函数不接受框架规范”。

如果我在 leas lag/lead() 处省略了 rowsBetween 框架规范,请不要抛出异常但返回意外的(对我而言)结果:总是 None。并且 rank() 仍然无法使用不同的异常。

谁能帮我把窗口功能弄好?

更新

好吧,这开始看起来像一个 pyspark 错误。 我已经在纯 Spark(Scala,spark-shell)中准备了相同的测试:

import sqlContext.implicits._
import org.apache.spark.sql._
import org.apache.spark.sql.types._

val l: List[Tuple2[Int,Int]] = List((1,101),(2,202),(3,303),(4,404),(5,505))
val rdd = sc.parallelize(l).map(i => Row(i._1,i._2))
val schemaString = "a b"
val schema = StructType(schemaString.split(" ").map(fieldName => StructField(fieldName, IntegerType, true)))
val df = sqlContext.createDataFrame(rdd, schema)

import org.apache.spark.sql.expressions.Window
import org.apache.spark.sql.functions._

val wSpec = Window.orderBy("a").rowsBetween(-1,1)
df.select(df("a"), rank().over(wSpec).alias("rank"))
    ==> org.apache.spark.sql.AnalysisException: Window function rank does not take a frame specification.;

df.select(df("a"), lag(df("b"),1).over(wSpec).alias("prev"), df("b"), lead(df("b"),1).over(wSpec).alias("next"))
    ===> org.apache.spark.sql.AnalysisException: Window function lag does not take a frame specification.;


val wSpec = Window.orderBy("a")
df.select(df("a"), rank().over(wSpec).alias("rank")).collect()
    ====> res10: Array[org.apache.spark.sql.Row] = Array([1,1], [2,2], [3,3], [4,4], [5,5])

df.select(df("a"), lag(df("b"),1).over(wSpec).alias("prev"), df("b"), lead(df("b"),1).over(wSpec).alias("next"))
    ====> res12: Array[org.apache.spark.sql.Row] = Array([1,null,101,202], [2,101,202,303], [3,202,303,404], [4,303,404,505], [5,404,505,null])

即使 rowsBetween 不能在 Scala 中应用,rank()lag()/lead() 在省略 rowsBetween 时都可以正常工作。

【问题讨论】:

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


    【解决方案1】:

    据我所知,有两个不同的问题。 Hive GenericUDAFRankGenericUDAFLagGenericUDAFLead 根本不支持窗口框架定义,因此您看到的错误是预期的行为。

    关于以下 PySpark 代码的问题

    wSpec = Window.orderBy(df.a)
    df.select(df.a, func.rank().over(wSpec).alias("rank"))
    

    它看起来与我的问题https://stackoverflow.com/q/31948194/1560062 有关,应该由SPARK-9978 解决。到目前为止,您可以通过将窗口定义更改为:

    wSpec = Window.partitionBy().orderBy(df.a)
    

    【讨论】:

      猜你喜欢
      • 2014-07-22
      • 2014-05-07
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2023-03-21
      • 1970-01-01
      相关资源
      最近更新 更多