Spark结构化流:将Kafka传来的毫秒级Epoch时间转换为Timestamp
嗨,我来帮你把这段代码补全,完美实现毫秒级时间戳转Timestamp类型的需求~
你已经完成了解析Kafka消息的核心步骤,现在只需要在最后一步处理timestamp字段就行。因为你的timestamp是毫秒级的Epoch长整型,而Spark的to_timestamp默认接受秒级数值,所以得先把毫秒数转成秒,再转换类型。
完整的代码如下:
val schema = StructType( List( StructField("timestamp", LongType, true), StructField("id", StringType, true), StructField("value", DoubleType, true) ) ) val dfNew = df.selectExpr("CAST(value AS STRING)") .as[String] .select(from_json($"value", schema) as "record") .select( $"record.id", $"record.value", // 将毫秒级时间戳转换为Timestamp类型 to_timestamp(col("record.timestamp") / 1000) as "timestamp" )
如果你的Spark版本较低,to_timestamp直接处理数值可能有兼容问题,也可以用from_unixtime先转成标准时间字符串,再转成Timestamp:
.select( $"record.id", $"record.value", to_timestamp(from_unixtime(col("record.timestamp") / 1000)) as "timestamp" )
这样处理后,timestamp字段就变成了Spark的Timestamp类型,后续不管是做时间窗口聚合还是格式化输出都方便多啦!
内容的提问来源于stack exchange,提问作者Abhijith V R
相关产品推荐
相关产品推荐

