【问题标题】:How to remove unchanged values from a timestamped dataframe in Spark?如何从 Spark 中的时间戳数据框中删除未更改的值?
【发布时间】:2020-01-19 13:13:37
【问题描述】:

所以我有以下 CSV 文件:

Timestamp,Point,Value
2019-09-01,A,1
2019-09-01,B,2
2019-09-02,A,1
2019-09-02,B,2
2019-09-03,A,3
2019-09-03,B,4
2019-09-04,A,3
2019-09-04,B,4
2019-09-05,A,1
2019-09-05,B,2

我正在使用以下代码在 Spark 2.4.3(Azure 上的 Databricks 5.4)上阅读它:

val df = spark
  .read
  .format("csv")
  .option("header", "true")
  .option("inferSchema", "true")
  .load("/test-data/data.csv")

我得到一个具有以下架构的数据框:

df:org.apache.spark.sql.DataFrame
  Timestamp:timestamp
  Point:string
  Value:double

此文件包含在不同时间点读取不同“点”的值。在此示例中,A 和 B 每 1 天有一次读数,但其中一些值与之前的读数相同。

我需要应用一个转换,该转换将只保留其 Value 列与先前读数相比已更改为同一点的行。

|Timestamp |Point|Value|
|----------|-----|-----|
|2019-09-01|A    |1    | // A = 1
|2019-09-01|B    |2    | // B = 2 
|2019-09-02|A    |1    | // A unchanged, should be removed
|2019-09-02|B    |2    | // B unchanged, should be removed
|2019-09-03|A    |3    | // A = 3
|2019-09-03|B    |4    | // B = 4
|2019-09-04|A    |3    | // A unchanged, should be removed
|2019-09-04|B    |4    | // B unchanged, should be removed
|2019-09-05|A    |1    | // A = 1
|2019-09-05|B    |2    | // B = 2

在这个简化的例子中,我想要一个如下的数据框:

|Timestamp |Point|Value|
|----------|-----|-----|
|2019-09-01|A    |1    |
|2019-09-01|B    |2    |
|2019-09-03|A    |3    |
|2019-09-03|B    |4    |
|2019-09-05|A    |1    |
|2019-09-05|B    |2    |

【问题讨论】:

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


    【解决方案1】:

    Spark 2.4.3 你可以使用Window函数来达到想要的效果。

    scala> var df_1= Seq(("2019-09-01","A",1),("2019-09-01","B",2),("2019-09-02","A",1),("2019-09-02","B",2),("2019-09-03","A",3),("2019-09-03","B",4),("2019-09-04","A",3),("2019-09-04","B",4),("2019-09-05","A",1),("2019-09-05","B",2)).toDF("Timestamp","Point","Value")
    
    scala> import org.apache.spark.sql.expressions.Window
    
    
    scala> df_1.show
    +----------+-----+-----+
    | Timestamp|Point|Value|
    +----------+-----+-----+
    |2019-09-01|    A|    1|
    |2019-09-01|    B|    2|
    |2019-09-02|    A|    1|
    |2019-09-02|    B|    2|
    |2019-09-03|    A|    3|
    |2019-09-03|    B|    4|
    |2019-09-04|    A|    3|
    |2019-09-04|    B|    4|
    |2019-09-05|    A|    1|
    |2019-09-05|    B|    2|
    +----------+-----+-----+
    
    scala> val win = Window.partitionBy("Point").orderBy("Timestamp","Point","Value")
    scala> val compareCols = List("Point", "Value")
    
    scala> val df2 = df_1. withColumn("compCols", struct(compareCols.map(col): _*)). withColumn("rowNum", row_number.over(win)). withColumn("toKeep", when($"rowNum" === 1 || $"compCols" =!= lag($"compCols", 1).over(win), true). otherwise(false) )
    scala> df2.filter(col("toKeep")===true).drop(col("compCols")).drop(col("rowNum")).drop(col("toKeep")).show
    +----------+-----+-----+
    | Timestamp|Point|Value|
    +----------+-----+-----+
    |2019-09-01|    B|    2|
    |2019-09-03|    B|    4|
    |2019-09-05|    B|    2|
    |2019-09-01|    A|    1|
    |2019-09-03|    A|    3|
    |2019-09-05|    A|    1|
    +----------+-----+-----+
    

    如果您有任何与此相关的疑问,请告诉我。

    【讨论】:

    • 这看起来很有希望。明天在工作中测试它并让你知道。谢谢!
    • @emzero 你确定
    • 谢谢,但它没有做我需要的。这似乎在做某种不同的事情,因为如果 Value 回到之前的任何值(但不是之前的值),它仍然会删除它。请参阅我更新的问题以更好地理解我在说什么。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2018-09-19
    • 2018-08-18
    相关资源
    最近更新 更多