You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

PySpark毫秒级前值填充失效,寻求可行解决方案

在PySpark中按毫秒级时间间隔为每个Node填充缺失时间点的前值

输入DataFrame

nodevaluetimestamp
node17777'2023-10-28 14:22:41.9'
node18888'2023-10-28 14:22:42.5'
node11111'2023-10-28 14:22:42.7'
node22222'2023-10-28 14:22:41.2'
node26666'2023-10-28 14:22:41.5'

期望输出DataFrame

nodevaluetimestamp
node17777'2023-10-28 14:22:41.9'
node17777'2023-10-28 14:22:42.0'
node17777'2023-10-28 14:22:42.1'
node17777'2023-10-28 14:22:42.2'
node17777'2023-10-28 14:22:42.3'
node17777'2023-10-28 14:22:42.4'
node18888'2023-10-28 14:22:42.5'
node18888'2023-10-28 14:22:42.6'
node11111'2023-10-28 14:22:42.7'
node22222'2023-10-28 14:22:41.2'
node22222'2023-10-28 14:22:41.3'
node22222'2023-10-28 14:22:41.4'
node26666'2023-10-28 14:22:41.5'

实现方案

核心思路

  1. 将字符串格式的时间解析为带毫秒精度的Timestamp类型,避免精度丢失
  2. 按node分组,计算每个分组的时间起止范围,生成连续的毫秒级时间序列
  3. 将生成的时间序列与原始数据关联,通过前向填充补全缺失的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.07 14:58:10