【问题标题】:Dynamic/Variable Offset in SparkSQL Lead/Lag functionSparkSQL Lead/Lag 函数中的动态/变量偏移量
【发布时间】:2021-03-02 20:03:28
【问题描述】:

我们能否以某种方式使用一个偏移值,该偏移值取决于 spark SQL 中超前/滞后函数中的列值?

示例:以下是可行的方法。

val sampleData = Seq( ("bob","Developer",125000),
  ("mark","Developer",108000),
  ("carl","Tester",70000),
  ("peter","Developer",185000),
  ("jon","Tester",65000),
  ("roman","Tester",82000),
  ("simon","Developer",98000),
  ("eric","Developer",144000),
  ("carlos","Tester",75000),
  ("henry","Developer",110000)).toDF("Name","Role","Salary")

val window = Window.orderBy("Role")

//Derive lag column for salary
val laggingCol = lag(col("Salary"), 1).over(window)

//Use derived column LastSalary to find difference between current and previous row
val salaryDifference = col("Salary") - col("LastSalary")

//Calculate trend based on the difference
//IF ELSE / CASE can be written using when.otherwise in spark
val trend = when(col("SalaryDiff").isNull || col("SalaryDiff").===(0), "SAME")
  .when(col("SalaryDiff").>(0), "UP")
  .otherwise("DOWN")

sampleData.withColumn("LastSalary", laggingCol)
  .withColumn("SalaryDiff",salaryDifference)
  .withColumn("Trend", trend).show()

现在,我的用例是,我们必须传递的偏移量取决于整数类型的特定列。这有点我想工作:

val sampleData = Seq( ("bob","Developer",125000,2),
  ("mark","Developer",108000,3),
  ("carl","Tester",70000,3),
  ("peter","Developer",185000,2),
  ("jon","Tester",65000,1),
  ("roman","Tester",82000,1),
  ("simon","Developer",98000,2),
  ("eric","Developer",144000,3),
  ("carlos","Tester",75000,2),
  ("henry","Developer",110000,2)).toDF("Name","Role","Salary","ColumnForOffset")

val window = Window.orderBy("Role")

//Derive lag column for salary
val laggingCol = lag(col("Salary"), col("ColumnForOffset")).over(window)

//Use derived column LastSalary to find difference between current and previous row
val salaryDifference = col("Salary") - col("LastSalary")

//Calculate trend based on the difference
//IF ELSE / CASE can be written using when.otherwise in spark
val trend = when(col("SalaryDiff").isNull || col("SalaryDiff").===(0), "SAME")
  .when(col("SalaryDiff").>(0), "UP")
  .otherwise("DOWN")

sampleData.withColumn("LastSalary", laggingCol)
  .withColumn("SalaryDiff",salaryDifference)
  .withColumn("Trend", trend).show()
   

这将按预期抛出异常,因为 offset 仅采用 Integer 值。 让我们讨论一下我们是否可以以某种方式为此实现逻辑。

【问题讨论】:

    标签: dataframe apache-spark apache-spark-sql


    【解决方案1】:

    您可以添加一个行号列,并根据行号和偏移量进行自连接,例如:

    val df = sampleData.withColumn("rn", row_number().over(window))
    
    val df2 = df.alias("t1").join(
        df.alias("t2"),
        expr("t1.rn = t2.rn + t1.ColumnForOffset"),
        "left"
    ).selectExpr("t1.*", "t2.Salary as LastSalary")
    
    df2.show
    +------+---------+------+---------------+---+----------+
    |  Name|     Role|Salary|ColumnForOffset| rn|LastSalary|
    +------+---------+------+---------------+---+----------+
    |   bob|Developer|125000|              2|  1|      null|
    |  mark|Developer|108000|              3|  2|      null|
    | peter|Developer|185000|              2|  3|    125000|
    | simon|Developer| 98000|              2|  4|    108000|
    |  eric|Developer|144000|              3|  5|    108000|
    | henry|Developer|110000|              2|  6|     98000|
    |  carl|   Tester| 70000|              3|  7|     98000|
    |   jon|   Tester| 65000|              1|  8|     70000|
    | roman|   Tester| 82000|              1|  9|     65000|
    |carlos|   Tester| 75000|              2| 10|     65000|
    +------+---------+------+---------------+---+----------+
    

    【讨论】:

    • 你好,如果我们还有一个列需要分组呢?
    • 然后您可以将该列添加到您的窗口定义中
    • 使用 .partitionBy
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2016-09-19
    • 1970-01-01
    • 2022-11-26
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多