PySpark中无法将毫秒级时间戳转换为dd-MM-yyyy HH:mm:ss格式
PySpark时间戳转换问题解决
我读取了从S3下载的、由Kinesis Firehose批处理生成的JSON文件,需要将其中的ts字段(毫秒级时间戳)转换为dd-MM-yyyy HH:mm:ss格式,当前代码转换后得到的是错误的未来日期,求修正。
当前代码
import pyspark from pyspark.sql import SparkSession from pyspark.sql import functions as f from pyspark.sql import types as t spark = SparkSession.builder.appName('demo').getOrCreate() d = spark.read.json('firehose-ds-iot-data-to-s3-1-2023-01-10-14-31-00-f23001c7-7759-38ca-9f14-4682cc39ae89') d.show() d.withColumn('ts',f.date_format(d.ts.cast(dataType=t.TimestampType()),"yyyy-MM-dd HH:mm:ss")) d.select('ts').show(truncate=False)
原始数据集
+-------------+------------+-----------+-----------+---------+-----------+-------------+ |Active-Import|Active-Power|Pyranometer|Temperature|Voltage-1|device_name| ts| +-------------+------------+-----------+-----------+---------+-----------+-------------+ | 2.57| 0| 0| 25.3| 239.65| inHand-RTU|1673361060486| | 2.57| 0| 0| 25.3| 239.44| inHand-RTU|1673361075375| | 2.57| 0| 0| 25.3| 239.44| inHand-RTU|1673361090384| | 2.57| 0| 0| 25.3| 239.44| inHand-RTU|1673361105397| | 2.57| 0| 0| 25.3| 239.44| inHand-RTU|1673361120532| | 2.57| 0| 0| 25.3| 239.43| inHand-RTU|1673361135503| | 2.57| 0| 0| 25.3| 239.43| inHand-RTU|1673361150520| | 2.57| 0| 0| 25.3| 239.27| inHand-RTU|1673361165428| | 2.57| 0| 0| 25.3| 236.14| inHand-RTU|1673361180435| | 2.57| 0| 0| 25.3| 236.14| inHand-RTU|1673361195440| | 2.57| 0| 0| 25.3| 236.03| inHand-RTU|1673361210450| | 2.57| 0| 0| 25.3| 236.03| inHand-RTU|1673361225498| | 2.57| 0| 0| 25.3| 236.08| inHand-RTU|1673361240595| | 2.57| 0| 0| 25.3| 236.08| inHand-RTU|1673361255512| | 2.57| 0| 0| 25.3| 236.09| inHand-RTU|1673361270490| | 2.57| 0| 0| 25.3| 235.96| inHand-RTU|1673361285544| | 2.57| 0| 0| 25.3| 235.96| inHand-RTU|1673361300800| | 2.57| 0| 0| 25.3| 235.94| inHand-RTU|1673361315630| | 2.57| 0| 0| 25.3| 235.94| inHand-RTU|1673361330528| | 2.57| 0| 0| 25.3| 235.75| inHand-RTU|1673361345566| +-------------+------------+-----------+-----------+---------+-----------+-------------+ only showing top 20 rows
当前错误输出
+---------------------+ |ts | +---------------------+ |+54996-09-13 00:00:00| |+54996-09-13 00:00:00| |+54996-09-13 00:00:00| |+54996-09-13 00:00:00| |+54996-09-13 00:00:00| |+54996-09-13 00:00:00| |+54996-09-14 00:00:00| |+54996-09-14 00:00:00| |+54996-09-14 00:00:00| |+54996-09-14 00:00:00| |+54996-09-14 00:00:00| |+54996-09-15 00:00:00| |+54996-09-15 00:00:00| |+54996-09-15 00:00:00| |+54996-09-15 00:00:00| |+54996-09-15 00:00:00| |+54996-09-15 00:00:00| |+54996-09-16 00:00:00| |+54996-09-16 00:00:00| |+54996-09-16 00:00:00| +---------------------+ only showing top 20 rows
问题原因及修正方案
问题1:时间戳单位不匹配
你的ts字段是毫秒级时间戳(比如1673361060486),但Spark将数字直接转换为TimestampType时,默认把数字当作秒级时间戳,导致时间被放大1000倍,出现未来日期。
问题2:DataFrame不可变性
Spark DataFrame是不可变的,withColumn方法会返回新的DataFrame,你没有将结果赋值给变量,所以原DataFramed的ts字段没有被修改。
问题3:格式不符合需求
你代码里用的是yyyy-MM-dd HH:mm:ss,但需求是dd-MM-yyyy HH:mm:ss,需要调整格式字符串。
修正后的代码
import pyspark from pyspark.sql import SparkSession from pyspark.sql import functions as f from pyspark.sql import types as t spark = SparkSession.builder.appName('demo').getOrCreate() d = spark.read.json('firehose-ds-iot-data-to-s3-1-2023-01-10-14-31-00-f23001c7-7759-38ca-9f14-4682cc39ae89') # 修正时间戳转换:先转成double除以1000得到秒级,再转Timestamp,最后格式化 d = d.withColumn( 'ts', f.date_format( f.to_timestamp(d.ts.cast(t.DoubleType()) / 1000), "dd-MM-yyyy HH:mm:ss" ) ) # 查看转换后的结果 d.select('ts').show(truncate=False)
说明
d.ts.cast(t.DoubleType()) / 1000:将毫秒级时间戳转为秒级浮点数f.to_timestamp(...):将秒级数值转为Spark Timestamp类型f.date_format(..., "dd-MM-yyyy HH:mm:ss"):将Timestamp格式化为目标字符串格式- 将
withColumn的结果重新赋值给d,确保修改生效
执行后,ts字段会显示正确的日期,比如10-01-2023 14:31:00,和文件名中的时间匹配。
内容的提问来源于stack exchange,提问作者anonymous_33008899
相关产品推荐
相关产品推荐

