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

Spark读取JSON数据时如何将日期时间字符串转为timestamp[us]?

问题:Spark合并JSON数据到Hudi表时时间戳异常

场景复现

JSON数据示例:

{
  "id":1,
  "time":"2023-01-01 12:34:56"
}

目标Hudi表通过PyArrow读取的Schema:

id: int64
time: timestamp[us, tz=UTC]

使用Spark合并的代码:

with SparkSession.builder.getOrCreate() as spark:
    df = spark.read.json('path/to/json/files')
    df.createOrReplaceTempView('upsert_data')
    spark.sql(f"CREATE TABLE snapshot USING hudi LOCATION 'path/to/hudi/table'")
    sql_upsert='''MERGE INTO snapshot AS target
                    USING upsert_data AS source
                    ON source.id = target.id
                  WHEN MATCHED THEN UPDATE SET target.time=source.time
                  WHEN NOT MATCHED THEN INSERT (id, time) values (source.id, source.time)'''
    spark.sql(sql_upsert)

合并后Hudi表中time字段变为56019-01-08 12:34:56,核心原因是Spark默认将时间字符串转为timestamp[ns],与Hudi表的timestamp[us]精度不匹配,导致时间数值被放大1000倍,最终出现时间异常。

解决方法

核心思路是统一时间戳精度,让Spark读取的JSON数据时间精度与Hudi表保持一致(微秒级),具体有两种可行方案:

方案1:读取JSON时显式指定Schema并转换精度

定义与Hudi表匹配的Schema,同时将Spark默认的纳秒级时间戳转换为微秒级:

from pyspark.sql.types import StructType, StructField, IntegerType, TimestampType

# 定义匹配Hudi表的Schema
custom_schema = StructType([
    StructField("id", IntegerType(), nullable=False),
    StructField("time", TimestampType(), nullable=True)
])

with SparkSession.builder.getOrCreate() as spark:
    # 读取JSON时指定Schema
    df = spark.read.schema(custom_schema).json('path/to/json/files')
    # 将纳秒级timestamp转换为微秒级(6位小数对应微秒精度)
    df = df.withColumn("time", df["time"].cast("timestamp(6)"))
    df.createOrReplaceTempView('upsert_data')
    
    spark.sql(f"CREATE TABLE snapshot USING hudi LOCATION 'path/to/hudi/table'")
    sql_upsert='''MERGE INTO snapshot AS target
                    USING upsert_data AS source
                    ON source.id = target.id
                  WHEN MATCHED THEN UPDATE SET target.time=source.time
                  WHEN NOT MATCHED THEN INSERT (id, time) values (source.id, source.time)'''
    spark.sql(sql_upsert)

方案2:在MERGE语句中直接转换时间精度

如果无法修改读取逻辑,可在MERGE的SQL语句中直接转换源数据的时间精度:

with SparkSession.builder.getOrCreate() as spark:
    df = spark.read.json('path/to/json/files')
    df.createOrReplaceTempView('upsert_data')
    spark.sql(f"CREATE TABLE snapshot USING hudi LOCATION 'path/to/hudi/table'")
    # 修改MERGE语句,转换时间精度为微秒级
    sql_upsert='''MERGE INTO snapshot AS target
                    USING upsert_data AS source
                    ON source.id = target.id
                  WHEN MATCHED THEN UPDATE SET target.time=CAST(source.time AS timestamp(6))
                  WHEN NOT MATCHED THEN INSERT (id, time) values (source.id, CAST(source.time AS timestamp(6)))'''
    spark.sql(sql_upsert)

内容的提问来源于stack exchange,提问作者Rinze

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 18:18:30