【发布时间】: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)
我想做的是通过减去EndDateTime和StartDateTime来找到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