如何将PySpark DataFrame的timestamp列转换为Jalali日期
问题原因
你写的代码无法运行的核心原因是:jdatetime是Python原生第三方日期库,不支持直接调用PySpark的col()对象,PySpark的分布式执行引擎无法直接识别Python本地对象的链式方法,必须通过UDF(用户自定义函数)封装转换逻辑才能在集群上分布式执行。
解决方案
前置依赖
首先确保所有集群节点都安装了jdatetime依赖,执行安装命令:
pip install jdatetime
基础实现(普通UDF)
通过PySpark的Python UDF封装公历转Jalali日期的逻辑,转换时自动保留原时间的时分秒部分,代码如下:
import jdatetime from pyspark.sql import functions as F from pyspark.sql.types import StringType # 定义转换UDF @F.udf(returnType=StringType()) def to_jalali_ts(creation_ts): if not creation_ts: return None # 将公历timestamp转为Jalali日期时间 j_dt = jdatetime.datetime.fromgregorian(datetime=creation_ts) # 按指定格式返回,保留时分秒 return j_dt.strftime("%Y-%m-%d %H:%M:%S") # 应用转换 df_result = df_etl_test_piko1.withColumn( "CreationDate", to_jalali_ts(F.col("CreationDate")) ) # 查看结果 df_result.show()
执行后输出和预期完全一致:
+----+-------------------+ |Name|CreationDate | +----+-------------------+ |Sara|1400-10-12 10:49:43| |Mina|1399-10-13 12:30:21| +----+-------------------+
高性能实现(向量式UDF,适配大数据量)
如果你的数据量超过百万级,普通Python UDF的序列化开销较高,可以替换为Pandas向量式UDF提升执行效率,代码如下:
import jdatetime import pandas as pd from pyspark.sql import functions as F from pyspark.sql.types import StringType @F.pandas_udf(returnType=StringType()) def to_jalali_ts_vector(ts_series: pd.Series) -> pd.Series: return ts_series.apply( lambda x: jdatetime.datetime.fromgregorian(datetime=x).strftime("%Y-%m-%d %H:%M:%S") if pd.notnull(x) else None ) # 调用方式和普通UDF完全一致 df_result = df_etl_test_piko1.withColumn( "CreationDate", to_jalali_ts_vector(F.col("CreationDate")) )
原代码错误说明
你之前的写法jdatetime.datetime.col('creationdate')存在语法错误:col()是PySpark SQL模块下的列选取方法,归属于pyspark.sql.functions包,jdatetime库本身不存在col属性,运行时会直接抛出AttributeError。
内容的提问来源于stack exchange,提问作者sara moradi
相关产品推荐
相关产品推荐

