如何用PySpark处理Unix timestamp生成10分钟区间的ts_start与ts_end
在PySpark SQL中生成10分钟时间区间的起始与结束时间
核心思路
先将Unix时间戳转换为Spark支持的Timestamp类型,再通过数值计算定位到所属10分钟区间的起始点,最后给起始点加上10分钟得到结束点。
具体实现步骤及代码
假设你的数据框df包含Unix时间戳列unix_ts(秒级,若为毫秒级需先做转换):
转换Unix时间戳为Timestamp类型
用from_unixtime函数将秒级Unix时间戳转成Spark可处理的Timestamp格式:from pyspark.sql import functions as F df = df.withColumn("ts", F.from_unixtime(F.col("unix_ts")))计算10分钟区间起始时间
ts_start
把时间戳转成秒数后,按10分钟(600秒)为单位取整,再转回Timestamp:df = df.withColumn( "ts_start", F.to_timestamp((F.floor(F.unix_timestamp(F.col("ts")) / 600) * 600)) )解释:
unix_timestamp(ts)获取当前时间的秒数,除以600取整后再乘600,得到当前时间所属10分钟区间的起始秒数,最后转成Timestamp格式。计算10分钟区间结束时间
ts_end
给ts_start加上600秒(10分钟)即可:df = df.withColumn( "ts_end", F.timestamp_add(F.col("ts_start"), 600) )
用Spark SQL语句实现
如果习惯用SQL语法,可将数据注册为临时视图后执行:
df.createOrReplaceTempView("time_data") result_df = spark.sql(""" SELECT unix_ts, from_unixtime(unix_ts) AS ts, to_timestamp((FLOOR(unix_timestamp(from_unixtime(unix_ts)) / 600) * 600)) AS ts_start, timestamp_add(to_timestamp((FLOOR(unix_timestamp(from_unixtime(unix_ts)) / 600) * 600)), 600) AS ts_end FROM time_data """)
毫秒级Unix时间戳适配
如果你的Unix时间戳是毫秒级,只需在转换时先除以1000:
df = df.withColumn("ts", F.from_unixtime(F.col("unix_ts") / 1000))
内容的提问来源于stack exchange,提问作者Любовь Пономарева
相关产品推荐
相关产品推荐

