使用PySpark转换Spark DataFrame为字典列表时datetime字段返回空值
问题解决:Spark DataFrame Timestamp数组转换为字典列表时返回空值
问题描述
我有一个Spark DataFrame df,结构如下:
| ID | values | datetime | ID2 | Port | TotalTime |
|---|---|---|---|---|---|
| id1 | [0, 1, 2, 3, 4, 5] | [01/01/23, 02/01/23, 03/01/23, 04/01/23, 05/01/23] | ab1 | 7 | 16 |
| id2 | [3, 6, 9] | [01/01/23, 02/01/23, 03/01/23] | ab2 | 64 | 27 |
需要将其转换为字典列表all_parts,每个字典结构要求:
part['ID'] = df['ID'] part['form']['values'] = df['values'] part['form']['datetime'] = df['datetime'] part['ID2'] = df['ID2'] part['port'] = df['Port'] part['TotalTime'] = df['TotalTime']
使用以下代码转换后,其他字段正常,但datetime字段返回空值,该字段是timestamp类型的数组,元素间隔20毫秒:
all_parts = [] for row in df.collect(): part = {} part['ID'] = row['ID'] part['form'] = {} part['form']['values'] = row['values'] part['form']['datetime'] = row['datetime'] part['ID2'] = row['ID2'] part['port'] = row['Port'] part['TotalTime'] = row['TotalTime'] all_parts.append(part)
问题原因
Spark的Timestamp类型数组在collect到本地Python环境时,默认的序列化逻辑无法正确解析数组中的Timestamp元素,直接访问row['datetime']会返回空值或无法识别的对象。
解决方案
提供两种可行的解决方式:
方法1:Spark层面转换为字符串数组
先通过Spark内置函数,将datetime数组中的每个Timestamp元素转为指定格式的字符串,再collect到本地处理:
from pyspark.sql import functions as F # 转换datetime列:将数组内的每个timestamp转为字符串格式 df_processed = df.withColumn( "datetime", F.transform( "datetime", lambda t: F.date_format(t, "dd/MM/yy") # 可根据需求调整日期格式 ) ) # 执行原转换逻辑 all_parts = [] for row in df_processed.collect(): part = {} part['ID'] = row['ID'] part['form'] = {} part['form']['values'] = row['values'] part['form']['datetime'] = row['datetime'] part['ID2'] = row['ID2'] part['port'] = row['Port'] part['TotalTime'] = row['TotalTime'] all_parts.append(part)
方法2:Python层面解析Timestamp元素
如果需要保留Python datetime对象而非字符串,可在collect后手动遍历数组,将每个Spark Timestamp转换为Python datetime:
all_parts = [] for row in df.collect(): part = {} part['ID'] = row['ID'] part['form'] = {} part['form']['values'] = row['values'] # 手动转换每个Spark Timestamp为Python datetime对象 part['form']['datetime'] = [t.to_pydatetime() for t in row['datetime']] part['ID2'] = row['ID2'] part['port'] = row['Port'] part['TotalTime'] = row['TotalTime'] all_parts.append(part)
验证
两种方法均可正常获取datetime字段的非空值:方法1适合需要字符串格式的场景,方法2适合后续需要对datetime对象进行操作的场景。
内容的提问来源于stack exchange,提问作者JGW
相关产品推荐
相关产品推荐

