【问题标题】:Date and Interval Addition in SparkSQLSparkSQL 中的日期和间隔相加
【发布时间】:2017-07-28 19:20:13
【问题描述】:

我正在尝试对 spark-shell 中的某个数据帧执行一个简单的 SQL 查询,该查询将 1 周的间隔添加到某个日期,如下所示:

原始查询:

scala> spark.sql("select Cast(table1.date2 as Date) + interval 1 week from table1").show()

现在我做了一些测试:

scala> spark.sql("select Cast('1999-09-19' as Date) + interval 1 week from table1").show()

我得到了正确的结果

+----------------------------------------------------------------------------+
|CAST(CAST(CAST(1999-09-19 AS DATE) AS TIMESTAMP) + interval 1 weeks AS DATE)|
+----------------------------------------------------------------------------+
|                                                                  1999-09-26|
+----------------------------------------------------------------------------+

(只需将 7 天添加到 19 = 26)

但是当我将年份改为 1997 年而不是 1999 年时,结果发生了变化!

scala> spark.sql("select Cast('1997-09-19' as Date) + interval 1 week from table1").show()

+----------------------------------------------------------------------------+
|CAST(CAST(CAST(1997-09-19 AS DATE) AS TIMESTAMP) + interval 1 weeks AS DATE)|
+----------------------------------------------------------------------------+
|                                                                  1997-09-25|
+----------------------------------------------------------------------------+

为什么结果会改变?不应该是 26 而不是 25?

那么,这是 sparkSQL 中与某种迭代计算丢失相关的错误,还是我遗漏了什么?

【问题讨论】:

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


    【解决方案1】:

    这可能是转换为当地时间的问题。 INTERVAL 将数据转换为TIMESTAMP,然后返回到DATE

    scala> spark.sql("SELECT CAST('1997-09-19' AS DATE) + INTERVAL 1 weeks").explain
    == Physical Plan ==
    *Project [10130 AS CAST(CAST(CAST(1997-09-19 AS DATE) AS TIMESTAMP) + interval 1 weeks AS DATE)#19]
    +- Scan OneRowRelation[]
    

    (注意第二个和第三个CASTs),Spark 被称为inconsequent when handling timestamps

    DATE_ADD 应该表现出更稳定的行为:

    scala> spark.sql("SELECT DATE_ADD(CAST('1997-09-19' AS DATE), 7)").explain
    == Physical Plan ==
    *Project [10130 AS date_add(CAST(1997-09-19 AS DATE), 7)#27]
    +- Scan OneRowRelation[]
    

    【讨论】:

    • 也不一致:如果您有一个跨越两个时区的集群,时间戳到日期的转换会完全崩溃(除非您每次都使用具有明确时区的方法)。
    【解决方案2】:

    从 Spark 3 开始,此错误已得到修复。让我们使用您提到的日期创建一个 DataFrame 并添加一周间隔。创建 DataFrame。

    import java.sql.Date
    
    val df = Seq(
      (Date.valueOf("1999-09-19")),
      (Date.valueOf("1997-09-19"))
    ).toDF("some_date")
    

    添加一周间隔:

    df
      .withColumn("plus_one_week", expr("some_date + INTERVAL 1 week"))
      .show()
    
    +----------+-------------+
    | some_date|plus_one_week|
    +----------+-------------+
    |1999-09-19|   1999-09-26|
    |1997-09-19|   1997-09-26|
    +----------+-------------+
    

    您也可以使用make_interval() SQL 函数获得相同的结果:

    df
      .withColumn("plus_one_week", expr("some_date + make_interval(0, 0, 1, 0, 0, 0, 0)"))
      .show()
    

    我们正在开发getting make_interval() exposed as Scala/PySpark functions,因此不必使用expr 来访问该功能。

    date_add 仅适用于添加天数,因此它是有限的。 make_interval() 更强大,因为它允许您添加年/月/日/小时/分钟/秒的任意组合。

    【讨论】:

    • 我可以确认函数 makeInterval 可用并在 spark 3.0.1 makeInterval(years = 0, months = 0, weeks = 1, days = 0, hours = 0, mins = 0, secs = Decimal(0)) 中使用数据集@
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 2023-04-03
    • 1970-01-01
    • 2014-12-17
    • 2023-03-29
    • 2021-08-15
    • 2023-03-07
    • 1970-01-01
    相关资源
    最近更新 更多