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
相关产品推荐
相关产品推荐

