如何在PySpark Schema中直接定义TimestampType类型列?
解决PySpark直接通过Schema定义TimestampType列的问题
报错核心原因:TimestampType对应的Python类型是**datetime.datetime对象**,而非字符串。直接传递字符串给定义为TimestampType的列时,Spark无法自动完成类型转换,因此抛出错误。以下是两种无需先定义StringType再cast的直接实现方式:
方法1:将数据中的时间字符串转为Python datetime对象
直接把数据列表里的时间字符串转换成Python标准库的datetime对象,让其匹配TimestampType的Schema要求:
from pyspark.sql import SparkSession from pyspark.sql.types import StructType, StructField, TimestampType, StringType, DoubleType from datetime import datetime spark = (SparkSession .builder .master("local[*]") .appName("Unit-tests") .getOrCreate()) # 把时间字符串转为datetime对象 data = [ (datetime.strptime("2023-04-25 00:00:00", "%Y-%m-%d %H:%M:%S"), "A", 100.5), (datetime.strptime("2023-04-26 00:00:00", "%Y-%m-%d %H:%M:%S"), "A", 110.0), (datetime.strptime("2023-04-28 00:00:00", "%Y-%m-%d %H:%M:%S"), "A", 105.0), (datetime.strptime("2023-04-27 00:00:00", "%Y-%m-%d %H:%M:%S"), "B", 50.5), (datetime.strptime("2023-04-29 00:00:00", "%Y-%m-%d %H:%M:%S"), "B", 55.5), ] schema = StructType([ StructField("time", TimestampType(), True), StructField("id", StringType(), True), StructField("value", DoubleType(), True) ]) df = spark.createDataFrame(data=data, schema=schema) df.printSchema() # 输出Schema: # root # |-- time: timestamp (nullable = true) # |-- id: string (nullable = true) # |-- value: double (nullable = true)
方法2:结合Row与to_timestamp函数生成Timestamp类型列(Spark 3.0+适用)
通过to_timestamp函数在数据构造阶段直接生成Timestamp类型的列,再配合预设的Schema定义:
from pyspark.sql import SparkSession, Row from pyspark.sql.types import StructType, StructField, TimestampType, StringType, DoubleType from pyspark.sql.functions import to_timestamp spark = (SparkSession .builder .master("local[*]") .appName("Unit-tests") .getOrCreate()) # 使用Row和to_timestamp生成Timestamp类型值 data = [ Row(time=to_timestamp("2023-04-25 00:00:00"), id="A", value=100.5), Row(time=to_timestamp("2023-04-26 00:00:00"), id="A", value=110.0), Row(time=to_timestamp("2023-04-28 00:00:00"), id="A", value=105.0), Row(time=to_timestamp("2023-04-27 00:00:00"), id="B", value=50.5), Row(time=to_timestamp("2023-04-29 00:00:00"), id="B", value=55.5), ] schema = StructType([ StructField("time", TimestampType(), True), StructField("id", StringType(), True), StructField("value", DoubleType(), True) ]) df = spark.createDataFrame(data=data, schema=schema) df.printSchema()
补充:全局配置自动识别时间字符串(不推荐用于单元测试)
如果是非单元测试场景,可通过设置Spark配置让其自动将符合格式的字符串转为Timestamp,但单元测试建议保持显式转换以避免环境依赖:
spark = (SparkSession .builder .master("local[*]") .appName("Unit-tests") .config("spark.sql.legacy.timeParserPolicy", "LEGACY") # 可根据Spark版本调整为CORRECTED .getOrCreate())
内容的提问来源于stack exchange,提问作者andKaae
相关产品推荐
相关产品推荐

