【发布时间】:2018-10-08 23:15:17
【问题描述】:
我正在开发 Spark 2.3、Python 3.6 和 pyspark 2.3.1
我有一个 Spark DataFrame,其中每个条目都是一个工作步骤,我想将一些行放在一个工作会话中。这应该在下面的函数getSessions 中完成。我相信它有效。
我进一步创建了一个 RDD,其中包含我想要的所有信息 - 每个条目都是具有所需列的 Row 对象,看起来类型很好(一些数据伪装):
rddSessions_flattened.take(1)
# [Row(counter=1, end=datetime.datetime(2017, 11, 6, 9, 15, 20), end_id=2758327, no_of_entries=5, shortID=u'DISGUISED', start=datetime.datetime(2017, 11, 6, 9, 13, 59), start_id=INTEGERNUMBER, strDuration='0:01:21', tNumber=u'DISGUISED', timeDuration=datetime.timedelta(0, 81))]
如果我现在想将一个 DataFrame 作为 RDD 的我们,我会得到一个 TypeError。
df = rddSessions_flattened.toDF()
df.show()
# TypeError: not supported type: type 'datetime.timedelta'
(最后发布的完整错误消息)
任何想法出了什么问题以及如何解决这个问题?
- 我基本上依靠 Spark 来推断模式;我认为这是因为 spark.sql.types 模块有一个类 TimestampType
- 如何以编程方式定义它?我不清楚 Apache Spark 编程指南。
欣赏你的想法!
def getSessions(values, threshold=threshold):
"""
Create sessions for one person on one case
Arguments:
values: sorted list of tuples (datetime, id)
threshold: time delta object; max time of a session
Return:
sessions: list of sessions (SData)
"""
SData = Row(
'counter'
, 'start'
, 'start_id'
, 'end'
, 'end_id'
, 'strDuration'
, 'timeDuration'
, 'no_of_entries'
)
counter = 1
no_of_rows = 1
sessions = [] # list of sessions
session_start = values[0][0] # first entry of the first tuple in the list
start_row_id = values[0][1]
session_end = session_start
end_row_id = start_row_id
for row_time, row_id in values[1:]:
if row_time - session_start > threshold:
# new session found, so append previous session
sessions.append(SData(
counter
, session_start
, start_row_id
, session_end
, end_row_id
, str(session_end - session_start)
, session_end - session_start
, no_of_rows
)
)
# get the information for the next session
counter += 1
no_of_rows = 1
session_start = row_time
start_row_id = row_id
else:
no_of_rows +=1
# regardless if new session or not: session_end reset to current entry
session_end = row_time
end_row_id = row_id
# very last session has to be recorded as there is no "next" row
sessions.append(SData(
counter
, session_start
, start_row_id
, session_end
, end_row_id
, str(session_end - session_start)
, session_end - session_start
, no_of_rows
)
)
return sessions
完整的错误信息:
not supported type: <type 'datetime.timedelta'>
Traceback (most recent call last):
File "/app/cdh/lib/parcels/SPARK2-2.3.0.cloudera2-1.cdh5.13.3.p0.316101/lib/spark2/python/lib/pyspark.zip/pyspark/sql/session.py", line 58, in toDF
return sparkSession.createDataFrame(self, schema, sampleRatio)
File "/app/cdh/lib/parcels/SPARK2-2.3.0.cloudera2-1.cdh5.13.3.p0.316101/lib/spark2/python/lib/pyspark.zip/pyspark/sql/session.py", line 687, in createDataFrame
rdd, schema = self._createFromRDD(data.map(prepare), schema, samplingRatio)
File "/app/cdh/lib/parcels/SPARK2-2.3.0.cloudera2-1.cdh5.13.3.p0.316101/lib/spark2/python/lib/pyspark.zip/pyspark/sql/session.py", line 384, in _createFromRDD
struct = self._inferSchema(rdd, samplingRatio, names=schema)
File "/app/cdh/lib/parcels/SPARK2-2.3.0.cloudera2-1.cdh5.13.3.p0.316101/lib/spark2/python/lib/pyspark.zip/pyspark/sql/session.py", line 364, in _inferSchema
schema = _infer_schema(first, names=names)
File "/app/cdh/lib/parcels/SPARK2-2.3.0.cloudera2-1.cdh5.13.3.p0.316101/lib/spark2/python/lib/pyspark.zip/pyspark/sql/types.py", line 1096, in _infer_schema
fields = [StructField(k, _infer_type(v), True) for k, v in items]
File "/app/cdh/lib/parcels/SPARK2-2.3.0.cloudera2-1.cdh5.13.3.p0.316101/lib/spark2/python/lib/pyspark.zip/pyspark/sql/types.py", line 1070, in _infer_type
raise TypeError("not supported type: %s" % type(obj))
TypeError: not supported type: <type 'datetime.timedelta'>
【问题讨论】:
标签: apache-spark pyspark apache-spark-sql pyspark-sql