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

如何从Kafka中正确获取CreationAt时间戳字段数据

问题

现有Kafka数据源传输的JSON格式如下:

{
    "CustomerId":606811,
    "Latitude":35.834896,
    "Longitude":50.019657,
    "Response":{
        "Stores":[
            {
                "Id":771,
                "LegacyStoreId":5497,
                "LegacyStoreTypeId":1,
                "PartnerId":3,
                "StoreName":"test",
                "StoreDisplayName":"test",
                "PartnerName":"test",
                "ServiceRadius":3,
                "Longitude":56.009797,
                "Latitude":35.829067,
                "StatusCode":1,
                "CityId":200,
                "CityName":null,
                "Description":"test",
                "IsOk24":false,
                "RouteDistanceInMeter":0,
                "IsRouteDistanceValid":false,
                "IsOutOfOrders":false,
                "AirDistanceInMeter":1100,
                "IsAirDistanceInValid":true,
                "IsDeliveryCoverage":true,
                "IsNonCoverageArea":false,
                "Rate":3.9,
                "Reviews":560,
                "IsHighPriorityStore":false,
                "StoreScore":0,
                "PartnerRank":1,
                "DeliveryCost":"80000",
                "FirstDeliveryTime":"test",
                "Labels":[],
                "MinCartTotalPrice":500000
            },
            {
                "Id":463,
                "LegacyStoreId":4765,
                "LegacyStoreTypeId":1,
                "PartnerId":3,
                "StoreName":"test",
                "StoreDisplayName":"test",
                "PartnerName":"test",
                "ServiceRadius":3,
                "Longitude":56.995281,
                "Latitude":35.82251,
                "StatusCode":1,
                "CityId":200,
                "CityName":null,
                "Description":"test",
                "IsOk24":false,
                "RouteDistanceInMeter":0,
                "IsRouteDistanceValid":false,
                "IsOutOfOrders":false,
                "AirDistanceInMeter":2593,
                "IsAirDistanceInValid":true,
                "IsDeliveryCoverage":true,
                "IsNonCoverageArea":false,
                "Rate":3.8,
                "Reviews":532,
                "IsHighPriorityStore":false,
                "StoreScore":0,
                "PartnerRank":1,
                "DeliveryCost":"80000",
                "FirstDeliveryTime":"test",
                "Labels":[],
                "MinCartTotalPrice":500000
            }
        ]
    },
    "Id":"f2655da4-c236-4f86-9ca0-8063a4c77da8",
    "CreationAt":"2023-06-18T14:49:11.8545562+03:30"
}

使用Spark Streaming读取数据时,其他字段均能正常解析,但CreationAt字段始终返回null,相关代码如下:

spark= SparkSession \
.builder \
.appName("striming") \
.config("spark.jars.packages","*****************") \
.config('spark.driver.extraClassPath', '/usr/local/spark/resources/jars/sqljdbc42.jar') \
.config('spark.executor.extraClassPath', '/usr/local/spark/resources/jars/sqljdbc42.jar') \
.config("spark.cores.max", "1") \
.config("spark.executor.memory", "1g") \
.config("spark.executor.cores", "1") \
.config("spark.dynamicAllocation.initialExecutors", "1") \
.master("local[1]") \
.getOrCreate()

sc = spark.sparkContext
sqlContext = SQLContext(sc)
kafka_df = spark.readStream.format("kafka").option("kafka.bootstrap.servers","My Kafka Servers").option("subscribe","Event").option("startingOffsets", "earliest").option("failOnDataLoss","false").option('multiline', "True").load()

schema = StructType([StructField("CreationAt", TimestampType(), True),StructField("CustomerId", LongType(), True),
StructField("Id", StringType(), True),StructField("Latitude", DoubleType(), True),
StructField("Longitude", DoubleType(), True),StructField("Response",
    StructType([StructField("Stores",ArrayType(StructType([StructField("Id", IntegerType(), True),
StructField("LegacyStoreId", IntegerType(), True), StructField("LegacyStoreTypeId", IntegerType(), True),
StructField("PartnerId", IntegerType(), True),StructField("StoreName", StringType(), True),
StructField("StoreDisplayName", StringType(), True),StructField("PartnerName", StringType(), True)
])), True), ]), True),])

value_df = kafka_df.select(from_json(col("value").cast("string"), schema).alias("StoreSelection"))

value_df.printSchema()

sites_flat = value_df.selectExpr("StoreSelection").select("StoreSelection.CustomerId", "StoreSelection.Latitude", 
"StoreSelection.Longitude", "StoreSelection.CreationAt", explode_outer("StoreSelection.Response.Stores").alias("Stores"))\
.select( "CustomerId","Latitude","Longitude","CreationAt","Stores.LegacyStoreId", "Stores.StoreName", "Stores.PartnerName")\
.select( "CustomerId","Latitude","Longitude", "CreationAt","LegacyStoreId", "StoreName", "PartnerName")

sites_flat.printSchema()

def foreach_batch_function(df, epoch_id):
    df.write \
    df.show()


invoiceWriterQuery = sites_flat.writeStream.foreachBatch(foreach_batch_function).outputMode("update") \
.option("checkpointLocation", "/usr/local/airflow/dags/otime/log_file").trigger(
processingTime="1 minute").start().awaitTermination()
invoiceWriterQuery.awaitTermination()

解决方案

问题原因

Spark默认的TimestampType无法直接解析带7位小数精度的时区格式时间字符串(示例中CreationAt的毫秒部分为7位),Spark时间解析默认仅支持最多6位小数的毫秒值,超出精度的部分会导致解析失败,最终返回null。

修复方法

有两种可行的修复方式:

方式1:自定义时间格式解析

先将CreationAt以StringType读取,再用to_timestamp函数指定兼容格式转换为Timestamp类型:

  1. 修改schema中CreationAt的类型:
schema = StructType([
    StructField("CreationAt", StringType(), True),
    StructField("CustomerId", LongType(), True),
    StructField("Id", StringType(), True),
    StructField("Latitude", DoubleType(), True),
    StructField("Longitude", DoubleType(), True),
    StructField("Response",
        StructType([
            StructField("Stores",ArrayType(StructType([
                StructField("Id", IntegerType(), True),
                StructField("LegacyStoreId", IntegerType(), True), 
                StructField("LegacyStoreTypeId", IntegerType(), True),
                StructField("PartnerId", IntegerType(), True),
                StructField("StoreName", StringType(), True),
                StructField("StoreDisplayName", StringType(), True),
                StructField("PartnerName", StringType(), True)
            ])), True), 
        ]), True),
])
  1. 在数据转换阶段执行时间格式转换:
from pyspark.sql.functions import to_timestamp, explode_outer

sites_flat = value_df.selectExpr("StoreSelection")\
.select(
    "StoreSelection.CustomerId", 
    "StoreSelection.Latitude", 
    "StoreSelection.Longitude", 
    to_timestamp(col("StoreSelection.CreationAt"), "yyyy-MM-dd'T'HH:mm:ss.SSSSSSSXXX").alias("CreationAt"),
    explode_outer("StoreSelection.Response.Stores").alias("Stores")
)\
.select( 
    "CustomerId","Latitude","Longitude","CreationAt",
    "Stores.LegacyStoreId", "Stores.StoreName", "Stores.PartnerName"
)

方式2:调整Spark配置支持高精度解析

在SparkSession初始化时添加配置,启用Java 8时间API以支持更高精度的时间解析:

spark= SparkSession \
.builder \
.appName("striming") \
.config("spark.jars.packages","*****************") \
.config('spark.driver.extraClassPath', '/usr/local/spark/resources/jars/sqljdbc42.jar') \
.config('spark.executor.extraClassPath', '/usr/local/spark/resources/jars/sqljdbc42.jar') \
.config("spark.sql.datetime.java8API.enabled", "true")  # 启用Java 8时间API
.config("spark.cores.max", "1") \
.config("spark.executor.memory", "1g") \
.config("spark.executor.cores", "1") \
.config("spark.dynamicAllocation.initialExecutors", "1") \
.master("local[1]") \
.getOrCreate()

添加该配置后,原有的TimestampType即可正确解析7位小数的时间字符串,无需修改schema。

额外注意点

代码中的foreach_batch_function存在语法错误,df.write \后缺少完整写入逻辑,建议修正为:

def foreach_batch_function(df, epoch_id):
    df.show()

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 21:37:09