【问题标题】:Can unix_timestamp() return unix time in milliseconds in Apache Spark?unix_timestamp() 可以在 Apache Spark 中以毫秒为单位返回 unix 时间吗?
【发布时间】:2021-03-10 18:56:17
【问题描述】:

我试图从时间戳字段中获取 unix 时间,以毫秒(13 位)为单位,但目前它以秒(10 位)返回。

scala> var df = Seq("2017-01-18 11:00:00.000", "2017-01-18 11:00:00.123", "2017-01-18 11:00:00.882", "2017-01-18 11:00:02.432").toDF()
df: org.apache.spark.sql.DataFrame = [value: string]

scala> df = df.selectExpr("value timeString", "cast(value as timestamp) time")
df: org.apache.spark.sql.DataFrame = [timeString: string, time: timestamp]


scala> df = df.withColumn("unix_time", unix_timestamp(df("time")))
df: org.apache.spark.sql.DataFrame = [timeString: string, time: timestamp ... 1 more field]

scala> df.take(4)
res63: Array[org.apache.spark.sql.Row] = Array(
[2017-01-18 11:00:00.000,2017-01-18 11:00:00.0,1484758800], 
[2017-01-18 11:00:00.123,2017-01-18 11:00:00.123,1484758800], 
[2017-01-18 11:00:00.882,2017-01-18 11:00:00.882,1484758800], 
[2017-01-18 11:00:02.432,2017-01-18 11:00:02.432,1484758802])

即使 2017-01-18 11:00:00.1232017-01-18 11:00:00.000 不同,我得到相同的 unix 时间回到 1484758800

我错过了什么?

【问题讨论】:

    标签: apache-spark apache-spark-sql unix-timestamp


    【解决方案1】:

    毫秒隐藏在小数部分时间戳格式中

    试试这个:

    df = df.withColumn("time_in_milliseconds", col("time").cast("double"))
    

    你会得到类似 1484758800.792,其中 792 是毫秒

    至少它对我有用(Scala、Spark、Hive)

    【讨论】:

      【解决方案2】:

      实现Dao Thi's answer中建议的方法

      import pyspark.sql.functions as F
      df = spark.createDataFrame([('22-Jul-2018 04:21:18.792 UTC', ),('23-Jul-2018 04:21:25.888 UTC',)], ['TIME'])
      df.show(2,False)
      df.printSchema()
      

      输出:

      +----------------------------+
      |TIME                        |
      +----------------------------+
      |22-Jul-2018 04:21:18.792 UTC|
      |23-Jul-2018 04:21:25.888 UTC|
      +----------------------------+
      root
      |-- TIME: string (nullable = true)
      

      字符串时间格式(包括毫秒)转换为unix_timestamp(double)。使用 substring 方法从字符串中提取毫秒(start_position = -7,length_of_substring=3)并将毫秒单独添加到 unix_timestamp。 (转换为子字符串以浮动添加)

      df1 = df.withColumn("unix_timestamp",F.unix_timestamp(df.TIME,'dd-MMM-yyyy HH:mm:ss.SSS z') + F.substring(df.TIME,-7,3).cast('float')/1000)
      

      在 Spark 中将 unix_timestamp(double) 转换为 timestamp 数据类型

      df2 = df1.withColumn("TimestampType",F.to_timestamp(df1["unix_timestamp"]))
      df2.show(n=2,truncate=False)
      

      这将为您提供以下输出

      +----------------------------+----------------+-----------------------+
      |TIME                        |unix_timestamp  |TimestampType          |
      +----------------------------+----------------+-----------------------+
      |22-Jul-2018 04:21:18.792 UTC|1.532233278792E9|2018-07-22 04:21:18.792|
      |23-Jul-2018 04:21:25.888 UTC|1.532319685888E9|2018-07-23 04:21:25.888|
      +----------------------------+----------------+-----------------------+
      

      检查架构:

      df2.printSchema()
      
      
      root
       |-- TIME: string (nullable = true)
       |-- unix_timestamp: double (nullable = true)
       |-- TimestampType: timestamp (nullable = true)
      

      【讨论】:

        【解决方案3】:

        unix_timestamp() 以秒为单位返回 unix 时间戳。

        时间戳的最后 3 位与毫秒字符串的最后 3 位相同(1.999sec = 1999 milliseconds),因此只需将时间戳字符串的最后 3 位附加到毫秒字符串的末尾即可。

        【讨论】:

          【解决方案4】:

          在 Spark 版本 3.0.1 之前,无法使用 SQL 内置函数 unix_timestamp 将时间戳转换为以毫秒为单位的 unix 时间。

          根据Spark的DateTimeUtils上的代码

          “时间戳对外暴露为java.sql.Timestamp,内部存储为longs,能够以微秒级精度存储时间戳。”

          因此,如果您定义具有 java.sql.Timestamp 作为输入的 UDF,您可以调用 getTime 以毫秒为单位的 Long。如果您申请 unix_timestamp,您将只能获得精确到秒的 unix 时间。

          val tsConversionToLongUdf = udf((ts: java.sql.Timestamp) => ts.getTime)
          

          将此应用于各种时间戳:

          val df = Seq("2017-01-18 11:00:00.000", "2017-01-18 11:00:00.111", "2017-01-18 11:00:00.110", "2017-01-18 11:00:00.100")
            .toDF("timestampString")
            .withColumn("timestamp", to_timestamp(col("timestampString")))
            .withColumn("timestampConversionToLong", tsConversionToLongUdf(col("timestamp")))
            .withColumn("timestampUnixTimestamp", unix_timestamp(col("timestamp")))
          
          df.printSchema()
          df.show(false)
          
          // returns
          root
           |-- timestampString: string (nullable = true)
           |-- timestamp: timestamp (nullable = true)
           |-- timestampConversionToLong: long (nullable = false)
           |-- timestampCastAsLong: long (nullable = true)
          
          +-----------------------+-----------------------+-------------------------+-------------------+
          |timestampString        |timestamp              |timestampConversionToLong|timestampUnixTimestamp|
          +-----------------------+-----------------------+-------------------------+-------------------+
          |2017-01-18 11:00:00.000|2017-01-18 11:00:00    |1484733600000            |1484733600         |
          |2017-01-18 11:00:00.111|2017-01-18 11:00:00.111|1484733600111            |1484733600         |
          |2017-01-18 11:00:00.110|2017-01-18 11:00:00.11 |1484733600110            |1484733600         |
          |2017-01-18 11:00:00.100|2017-01-18 11:00:00.1  |1484733600100            |1484733600         |
          +-----------------------+-----------------------+-------------------------+-------------------+
          

          【讨论】:

            【解决方案5】:

            无法使用 unix_timestamp() 完成,但从 Spark 3.1.0 开始,有一个名为 unix_millis() 的内置函数:

            unix_millis(timestamp) - 返回自 1970-01-01 00:00:00 UTC 以来的毫秒数。截断更高级别的精度。

            【讨论】:

              猜你喜欢
              • 1970-01-01
              • 2011-08-02
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              • 1970-01-01
              相关资源
              最近更新 更多