【问题标题】:Calculating duration by subtracting two datetime columns in string format通过减去两个字符串格式的日期时间列来计算持续时间
【发布时间】:2015-07-28 18:57:04
【问题描述】:

我有一个 Spark Dataframe,其中包含一系列日期:

from pyspark.sql import SQLContext
from pyspark.sql import Row
from pyspark.sql.types import *
sqlContext = SQLContext(sc)
import pandas as pd

rdd = sc.parallelizesc.parallelize([('X01','2014-02-13T12:36:14.899','2014-02-13T12:31:56.876','sip:4534454450'),
                                    ('X02','2014-02-13T12:35:37.405','2014-02-13T12:32:13.321','sip:6413445440'),
                                    ('X03','2014-02-13T12:36:03.825','2014-02-13T12:32:15.229','sip:4534437492'),
                                    ('XO4','2014-02-13T12:37:05.460','2014-02-13T12:32:36.881','sip:6474454453'),
                                    ('XO5','2014-02-13T12:36:52.721','2014-02-13T12:33:30.323','sip:8874458555')])
schema = StructType([StructField('ID', StringType(), True),
                     StructField('EndDateTime', StringType(), True),
                     StructField('StartDateTime', StringType(), True)])
df = sqlContext.createDataFrame(rdd, schema)

我想做的是通过减去EndDateTimeStartDateTime来找到duration。我想我会尝试使用一个函数来做到这一点:

# Function to calculate time delta
def time_delta(y,x): 
    end = pd.to_datetime(y)
    start = pd.to_datetime(x)
    delta = (end-start)
    return delta

# create new RDD and add new column 'Duration' by applying time_delta function
df2 = df.withColumn('Duration', time_delta(df.EndDateTime, df.StartDateTime)) 

但这只是给了我:

>>> df2.show()
ID  EndDateTime          StartDateTime        ANI            Duration
X01 2014-02-13T12:36:... 2014-02-13T12:31:... sip:4534454450 null    
X02 2014-02-13T12:35:... 2014-02-13T12:32:... sip:6413445440 null    
X03 2014-02-13T12:36:... 2014-02-13T12:32:... sip:4534437492 null    
XO4 2014-02-13T12:37:... 2014-02-13T12:32:... sip:6474454453 null    
XO5 2014-02-13T12:36:... 2014-02-13T12:33:... sip:8874458555 null  

我不确定我的方法是否正确。如果没有,我很乐意接受另一种建议的方法来实现这一目标。

【问题讨论】:

  • 你试过在 REPL 中调试吗?
  • @dskrvk 因为我不是开发人员,所以我没有太多调试经验。但是,我怀疑问题在于 Spark 如何将数据传递给函数。例如, time_delta() 在纯 Python 中工作。出于某种原因,某些 Python/Pandas 函数不能很好地发挥作用。例如。 import re def extract_ani(x): extract = x.str.extract(r'(\d{10})') return extract Dates = Dates.withColumn('Cell', extract_ani(Dates.ANI)) 也出现错误Spark DataFrames,但当我将数据帧转换为 RDD 并将该函数用作 sc.map 的一部分时有效
  • 在 Scala 中,我将使用 TimestampType 而不是 StringType 来保存日期,然后创建一个 UDF 来计算两列之间的差异。我在任何地方都没有看到您将 time_delta 声明为用户定义的函数,但这是 Scala 中的必要步骤,以使其执行您想做的事情。
  • 是的,看看 pyspark.sql.functions.udf 下的spark.apache.org/docs/latest/api/python/…。您需要将 time_delta 创建为 UDF
  • @David Griffin 你是对的 :) 我最初忽略了注册 UDF,因为我认为你必须注册 UDF,只有你想使用 select 表达式

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


【解决方案1】:

这可以在 spark-sql 中通过将字符串日期转换为时间戳然后获取差异来完成。

1:转换为时间戳:

CAST(UNIX_TIMESTAMP(MY_COL_NAME,'dd-MMM-yy') as TIMESTAMP)

2:使用datediff函数获取日期之间的差异。

这将组合在一个嵌套函数中,例如:

spark.sql("select COL_1, COL_2, datediff( CAST( UNIX_TIMESTAMP( COL_1,'dd-MMM-yy') as TIMESTAMP), CAST( UNIX_TIMESTAMP( COL_2,'dd-MMM-yy') as TIMESTAMP) ) as LAG_in_days from MyTable")

结果如下:

+---------+---------+-----------+
|    COL_1|    COL_2|LAG_in_days|
+---------+---------+-----------+
|24-JAN-17|16-JAN-17|          8|
|19-JAN-05|18-JAN-05|          1|
|23-MAY-06|23-MAY-06|          0|
|18-AUG-06|17-AUG-06|          1|
+---------+---------+-----------+

参考:https://docs-snaplogic.atlassian.net/wiki/spaces/SD/pages/2458071/Date+Functions+and+Properties+Spark+SQL

【讨论】:

    【解决方案2】:

    使用 DoubleType 而不是 IntegerType

    from pyspark.sql import SQLContext, Row
    sqlContext = SQLContext(sc)
    from pyspark.sql.types import StringType, IntegerType, StructType, StructField
    from pyspark.sql.functions import udf
    
    
    # Build sample data
    rdd = sc.parallelize([('X01','2014-02-13T12:36:14.899','2014-02-13T12:31:56.876'),
                          ('X02','2014-02-13T12:35:37.405','2014-02-13T12:32:13.321'),
                          ('X03','2014-02-13T12:36:03.825','2014-02-13T12:32:15.229'),
                          ('XO4','2014-02-13T12:37:05.460','2014-02-13T12:32:36.881'),
                          ('XO5','2014-02-13T12:36:52.721','2014-02-13T12:33:30.323')])
    schema = StructType([StructField('ID', StringType(), True),
                         StructField('EndDateTime', StringType(), True),
                         StructField('StartDateTime', StringType(), True)])
    df = sqlContext.createDataFrame(rdd, schema)
    
    # define timedelta function (obtain duration in seconds)
    def time_delta(y,x): 
        from datetime import datetime
        end = datetime.strptime(y, '%Y-%m-%dT%H:%M:%S.%f')
        start = datetime.strptime(x, '%Y-%m-%dT%H:%M:%S.%f')
        delta = (end-start).total_seconds()
        return delta
    
    # register as a UDF 
    f = udf(time_delta, DoubleType())
    
    # Apply function
    df2 = df.withColumn('Duration', f(df.EndDateTime, df.StartDateTime))
    

    【讨论】:

      【解决方案3】:

      这是来自 jason 的 answer 的 spark 2.x 的工作版本

      from pyspark import SparkContext, SparkConf
      from pyspark.sql import SparkSession,SQLContext
      from pyspark.sql.types import StringType, StructType, StructField
      
      sc = SparkContext()
      sqlContext = SQLContext(sc)
      spark = SparkSession.builder.appName("Python Spark SQL basic example").getOrCreate()
      
      rdd = sc.parallelize([('X01','2014-02-13T12:36:14.899','2014-02-13T12:31:56.876'),
                            ('X02','2014-02-13T12:35:37.405','2014-02-13T12:32:13.321'),
                            ('X03','2014-02-13T12:36:03.825','2014-02-13T12:32:15.229'),
                            ('XO4','2014-02-13T12:37:05.460','2014-02-13T12:32:36.881'),
                            ('XO5','2014-02-13T12:36:52.721','2014-02-13T12:33:30.323')])
      schema = StructType([StructField('ID', StringType(), True),
                           StructField('EndDateTime', StringType(), True),
                           StructField('StartDateTime', StringType(), True)])
      df = sqlContext.createDataFrame(rdd, schema)
      
      # register as a UDF 
      from datetime import datetime
      sqlContext.registerFunction("time_delta", lambda y,x:(datetime.strptime(y, '%Y-%m-%dT%H:%M:%S.%f')-datetime.strptime(x, '%Y-%m-%dT%H:%M:%S.%f')).total_seconds())
      
      df.createOrReplaceTempView("Test_table")
      
      spark.sql("SELECT ID,EndDateTime,StartDateTime,time_delta(EndDateTime,StartDateTime) as time_delta FROM Test_table").show()
      
      sc.stop()
      

      【讨论】:

        【解决方案4】:
        datediff(Column end, Column start)
        

        返回从开始到结束的天数。

        https://spark.apache.org/docs/1.6.2/api/java/org/apache/spark/sql/functions.html

        【讨论】:

          【解决方案5】:

          从 Spark 1.5 开始,您可以使用 unix_timestamp:

          from pyspark.sql import functions as F
          timeFmt = "yyyy-MM-dd'T'HH:mm:ss.SSS"
          timeDiff = (F.unix_timestamp('EndDateTime', format=timeFmt)
                      - F.unix_timestamp('StartDateTime', format=timeFmt))
          df = df.withColumn("Duration", timeDiff)
          

          注意 Java 风格的时间格式。

          >>> df.show()
          +---+--------------------+--------------------+--------+
          | ID|         EndDateTime|       StartDateTime|Duration|
          +---+--------------------+--------------------+--------+
          |X01|2014-02-13T12:36:...|2014-02-13T12:31:...|     258|
          |X02|2014-02-13T12:35:...|2014-02-13T12:32:...|     204|
          |X03|2014-02-13T12:36:...|2014-02-13T12:32:...|     228|
          |XO4|2014-02-13T12:37:...|2014-02-13T12:32:...|     269|
          |XO5|2014-02-13T12:36:...|2014-02-13T12:33:...|     202|
          +---+--------------------+--------------------+--------+
          

          【讨论】:

          • 您可以除以 3600.0 转换为小时 df.withColumn("Duration_hours", df.Duration / 3600.0)
          【解决方案6】:

          感谢大卫格里芬。以下是如何执行此操作以供将来参考。

          from pyspark.sql import SQLContext, Row
          sqlContext = SQLContext(sc)
          from pyspark.sql.types import StringType, IntegerType, StructType, StructField
          from pyspark.sql.functions import udf
          
          # Build sample data
          rdd = sc.parallelize([('X01','2014-02-13T12:36:14.899','2014-02-13T12:31:56.876'),
                                ('X02','2014-02-13T12:35:37.405','2014-02-13T12:32:13.321'),
                                ('X03','2014-02-13T12:36:03.825','2014-02-13T12:32:15.229'),
                                ('XO4','2014-02-13T12:37:05.460','2014-02-13T12:32:36.881'),
                                ('XO5','2014-02-13T12:36:52.721','2014-02-13T12:33:30.323')])
          schema = StructType([StructField('ID', StringType(), True),
                               StructField('EndDateTime', StringType(), True),
                               StructField('StartDateTime', StringType(), True)])
          df = sqlContext.createDataFrame(rdd, schema)
          
          # define timedelta function (obtain duration in seconds)
          def time_delta(y,x): 
              from datetime import datetime
              end = datetime.strptime(y, '%Y-%m-%dT%H:%M:%S.%f')
              start = datetime.strptime(x, '%Y-%m-%dT%H:%M:%S.%f')
              delta = (end-start).total_seconds()
              return delta
          
          # register as a UDF 
          f = udf(time_delta, IntegerType())
          
          # Apply function
          df2 = df.withColumn('Duration', f(df.EndDateTime, df.StartDateTime)) 
          

          应用time_delta() 将给你以秒为单位的持续时间:

          >>> df2.show()
          ID  EndDateTime          StartDateTime        Duration
          X01 2014-02-13T12:36:... 2014-02-13T12:31:... 258     
          X02 2014-02-13T12:35:... 2014-02-13T12:32:... 204     
          X03 2014-02-13T12:36:... 2014-02-13T12:32:... 228     
          XO4 2014-02-13T12:37:... 2014-02-13T12:32:... 268     
          XO5 2014-02-13T12:36:... 2014-02-13T12:33:... 202 
          

          【讨论】:

          • 请使用 (end-start).total_seconds() 。否则你会得到这样令人讨厌的惊喜: time_delta('2014-02-13T12:36:14.000', '2014-02-13T12:36:15.900') 返回 86398 而不是 -1.9
          • 此代码不再起作用。持续时间显示为空。使用 zeppelin,火花 1.6
          猜你喜欢
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 1970-01-01
          • 2022-01-17
          • 2020-06-28
          • 1970-01-01
          • 1970-01-01
          相关资源
          最近更新 更多