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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.31 00:57:28