如何从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类型:
- 修改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), ])
- 在数据转换阶段执行时间格式转换:
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
相关产品推荐
相关产品推荐

