如何让PySpark SQL API返回带时区的datetime对象?测试失败分析
问题
我编写了如下单元测试,创建带时区信息的datetime对象并尝试返回:
from datetime import datetime, timezone from pyspark.sql import SparkSession from pyspark.sql.types import StructField, StructType, TimestampType def test_timezones() -> None: session = SparkSession.builder.master("local[4]").getOrCreate() df = session.createDataFrame( [(datetime(2022, 10, 1, 10, 53, tzinfo=timezone.utc),)], StructType( [ StructField("my_field", TimestampType(), nullable=False), ] ), ) assert df.take(1)[0].my_field == datetime(2022, 10, 1, 10, 53, tzinfo=timezone.utc)
但测试失败,报错信息如下:
Expected :datetime.datetime(2022, 10, 1, 10, 53, tzinfo=datetime.timezone.utc) Actual :datetime.datetime(2022, 10, 1, 11, 53)
我使用的是PySpark 3.3.1版本,为何无法得到预期的datetime对象?如何让PySpark SQL API返回带时区信息的datetime对象?
原因分析
- TimestampType的时区特性:PySpark的
TimestampType对应SQL标准的TIMESTAMP类型,本身不存储时区信息。它会将传入的带时区datetime转换为UTC时间戳存储,但读取时会自动转换为Spark会话的本地时区(你的环境时区应为UTC+1),最终返回的datetime对象不带tzinfo,且时间值是本地时区的时间。 - 对象匹配逻辑问题:预期的是带UTC时区信息的datetime对象,而实际返回的是无时区的本地时间datetime,两者既无时区匹配,时间值也因时区转换存在差异,导致断言失败。
解决方案
方法1:将Spark会话时区设为UTC
修改SparkSession配置,强制会话使用UTC时区,这样读取时返回的datetime时间值与原始UTC时间一致(仍不带tzinfo):
def test_timezones() -> None: session = SparkSession.builder \ .master("local[4]") \ .config("spark.sql.session.timeZone", "UTC") \ .getOrCreate() df = session.createDataFrame( [(datetime(2022, 10, 1, 10, 53, tzinfo=timezone.utc),)], StructType( [ StructField("my_field", TimestampType(), nullable=False), ] ), ) # 此时返回的datetime不带tzinfo,但时间值为UTC的10:53 assert df.take(1)[0].my_field == datetime(2022, 10, 1, 10, 53)
方法2:手动为返回的datetime添加时区信息
如果需要严格匹配带tzinfo的datetime对象,在会话时区设为UTC后,为取出的datetime手动添加UTC时区:
def test_timezones() -> None: session = SparkSession.builder \ .master("local[4]") \ .config("spark.sql.session.timeZone", "UTC") \ .getOrCreate() df = session.createDataFrame( [(datetime(2022, 10, 1, 10, 53, tzinfo=timezone.utc),)], StructType( [ StructField("my_field", TimestampType(), nullable=False), ] ), ) retrieved_dt = df.take(1)[0].my_field # 手动注入UTC时区信息 assert retrieved_dt.replace(tzinfo=timezone.utc) == datetime(2022, 10, 1, 10, 53, tzinfo=timezone.utc)
方法3:使用带时区的字符串存储(可选)
若需完全保留时区元数据,可将datetime转换为ISO 8601格式的带时区字符串存储,读取时再解析为带tzinfo的datetime对象,但这种方式会失去Timestamp类型的计算优势。
内容的提问来源于stack exchange,提问作者Jon
相关产品推荐
相关产品推荐

