PySpark毫秒级前值填充失效,寻求可行解决方案
在PySpark中按毫秒级时间间隔为每个Node填充缺失时间点的前值
输入DataFrame
| node | value | timestamp |
|---|---|---|
| node1 | 7777 | '2023-10-28 14:22:41.9' |
| node1 | 8888 | '2023-10-28 14:22:42.5' |
| node1 | 1111 | '2023-10-28 14:22:42.7' |
| node2 | 2222 | '2023-10-28 14:22:41.2' |
| node2 | 6666 | '2023-10-28 14:22:41.5' |
期望输出DataFrame
| node | value | timestamp |
|---|---|---|
| node1 | 7777 | '2023-10-28 14:22:41.9' |
| node1 | 7777 | '2023-10-28 14:22:42.0' |
| node1 | 7777 | '2023-10-28 14:22:42.1' |
| node1 | 7777 | '2023-10-28 14:22:42.2' |
| node1 | 7777 | '2023-10-28 14:22:42.3' |
| node1 | 7777 | '2023-10-28 14:22:42.4' |
| node1 | 8888 | '2023-10-28 14:22:42.5' |
| node1 | 8888 | '2023-10-28 14:22:42.6' |
| node1 | 1111 | '2023-10-28 14:22:42.7' |
| node2 | 2222 | '2023-10-28 14:22:41.2' |
| node2 | 2222 | '2023-10-28 14:22:41.3' |
| node2 | 2222 | '2023-10-28 14:22:41.4' |
| node2 | 6666 | '2023-10-28 14:22:41.5' |
实现方案
核心思路
- 将字符串格式的时间解析为带毫秒精度的Timestamp类型,避免精度丢失
- 按
node分组,计算每个分组的时间起止范围,生成连续的毫秒级时间序列 - 将生成的时间序列与原始数据关联,通过前向填充补全缺失的
value值
代码实现
from pyspark.sql import SparkSession from pyspark.sql.functions import col, last, sequence, explode, to_timestamp, min, max from pyspark.sql.types import TimestampType # 初始化SparkSession spark = SparkSession.builder.appName("MillisecondFill").getOrCreate() # 构造测试数据 data = [ ("node1", 7777, "2023-10-28 14:22:41.9"), ("node1", 8888, "2023-10-28 14:22:42.5"), ("node1", 1111, "2023-10-28 14:22:42.7"), ("node2", 2222, "2023-10-28 14:22:41.2"), ("node2", 6666, "2023-10-28 14:22:41.5") ] df = spark.createDataFrame(data, ["node", "value", "timestamp_str"]) # 转换时间字符串为带毫秒精度的Timestamp类型 df = df.withColumn("timestamp", to_timestamp(col("timestamp_str"), "yyyy-MM-dd HH:mm:ss.S")) # 计算每个node的时间起止范围 time_ranges = df.groupBy("node").agg( min("timestamp").alias("min_ts"), max("timestamp").alias("max_ts") ) # 生成每个node的连续毫秒级时间序列 time_sequences = time_ranges.withColumn( "timestamp", explode(sequence(col("min_ts"), col("max_ts"), interval 1 millisecond)) ).select("node", "timestamp") # 关联原始数据并执行前向填充 result = time_sequences.join(df, on=["node", "timestamp"], how="left") \ .orderBy("node", "timestamp") \ .groupBy("node") \ .agg( last("timestamp", ignorenulls=True).alias("timestamp"), last("value", ignorenulls=True).alias("value") ) \ .orderBy("node", "timestamp") # 输出结果 result.show(truncate=False)
关键细节说明
- 时间精度保证:使用
to_timestamp指定格式yyyy-MM-dd HH:mm:ss.S,确保毫秒部分被正确解析和存储 - 毫秒序列生成:
sequence函数搭配interval 1 millisecond参数,直接生成连续的毫秒级时间点,这是处理毫秒间隔的核心 - 前向填充逻辑:通过
last(..., ignorenulls=True)配合分组排序,自动用最近的非空值填充缺失的value,保证每个毫秒时间点都有对应值
内容的提问来源于stack exchange,提问作者Waleed saeed
相关产品推荐
相关产品推荐

