You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何在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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.23 17:32:06