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

如何在Polars中进行上采样且保留原有数据?

问题:Polars上采样时保留原有非整点记录并填充间隙

原始DataFrame:

┌───────────┬─────────────────────┬───────────┐
│ sensor_id ┆ ts                  ┆ value     │
│ ---       ┆ ---                 ┆ ---       │
│ i32       ┆ datetime[ms]        ┆ f64       │
╞═══════════╪═════════════════════╪═══════════╡
│ 1551      ┆ 2025-12-25 00:00:06 ┆ -1.464e6  │
│ 1551      ┆ 2025-12-25 00:05:17 ┆ -1.4639e6 │
│ 1551      ┆ 2025-12-25 00:15:43 ┆ -1.4638e6 │
│ 1551      ┆ 2025-12-25 00:45:51 ┆ -1.4637e6 │
│ 1551      ┆ 2025-12-25 01:45:06 ┆ -1.4636e6 │
│ 1551      ┆ 2025-12-25 03:45:17 ┆ -1.4635e6 │
│ 1551      ┆ 2025-12-25 03:50:43 ┆ -1.4634e6 │
│ 1551      ┆ 2025-12-25 04:00:51 ┆ -1.4633e6 │
│ 1551      ┆ 2025-12-25 04:30:06 ┆ -1.4632e6 │
│ 1551      ┆ 2025-12-25 05:30:17 ┆ -1.4631e6 │
└───────────┴─────────────────────┴───────────┘

需求:对数据间隙进行5分钟间隔的上采样,同时保留原有秒数不同的记录,不能丢失原始的非整5分钟行。

当前使用代码:

df = df.upsample(["sensor_id"], "ts", polars::time::Duration::parse("5m")).unwrap();
dbg!(df.clone());

得到的问题结果:丢失了原始的非整5分钟记录,仅保留间隙处的空值行:

┌─────────────────────┬───────────┬──────────┐
│ ts                  ┆ sensor_id ┆ value    │
│ ---                 ┆ ---       ┆ ---      │
│ datetime[ms]        ┆ i32       ┆ f64      │
╞═════════════════════╪═══════════╪══════════╡
│ 2025-12-25 00:00:06 ┆ 1551      ┆ -1.464e6 │
│ 2025-12-25 00:05:06 ┆ null      ┆ null     │
│ 2025-12-25 00:10:06 ┆ null      ┆ null     │
│ 2025-12-25 00:15:06 ┆ null      ┆ null     │
│ 2025-12-25 00:20:06 ┆ null      ┆ null     │
│ …                   ┆ …         ┆ …        │
│ 2025-12-25 05:10:06 ┆ null      ┆ null     │
│ 2025-12-25 05:15:06 ┆ null      ┆ null     │
│ 2025-12-25 05:20:06 ┆ null      ┆ null     │
│ 2025-12-25 05:25:06 ┆ null      ┆ null     │
│ 2025-12-25 05:30:06 ┆ null      ┆ null     │
└─────────────────────┴───────────┴──────────┘

解决方案

Polars原生的upsample方法会生成规整的时间间隔序列,但不会保留原始的非间隔时间点。要实现保留原始数据+填充间隙时间点的效果,需要先生成完整的5分钟间隔序列,再与原始数据做外连接:

代码实现

use polars::prelude::*;
use chrono::NaiveDateTime;

// 先生成每个sensor_id对应的5分钟规整时间序列
let upsample_df = df
    .clone()
    .group_by(["sensor_id"])
    .agg([
        // 计算时间范围并对齐到5分钟间隔的边界
        col("ts")
            .min()
            .dt()
            .floor("5m")
            .alias("start"),
        col("ts")
            .max()
            .dt()
            .ceil("5m")
            .alias("end"),
    ])
    .with_columns([
        // 生成从start到end的所有5分钟间隔时间点
        col("start")
            .dt()
            .range(
                col("end"),
                Duration::parse("5m").unwrap(),
            )
            .alias("ts"),
    ])
    .explode(["ts"])  // 将数组展开为多行
    .select(["sensor_id", "ts"]);

// 外连接原始数据,保留所有原始行和规整时间点
let result_df = upsample_df
    .join(df, on=["sensor_id", "ts"], how="outer")
    .sort(["sensor_id", "ts"], Default::default())
    .unwrap();

dbg!(result_df);

逻辑说明

  1. 生成规整时间序列:按sensor_id分组,计算每个传感器数据的时间范围,将起始时间向下取整到最近的5分钟、结束时间向上取整到最近的5分钟,然后生成该范围内所有5分钟间隔的时间点。
  2. 外连接合并数据:用outer连接方式将规整时间序列和原始数据合并,这样原始的非整点记录和间隙的整点记录都会被保留,value字段为空的就是间隙需要填充的行。
  3. 排序:按sensor_id和ts排序,得到有序的完整数据。

测试数据生成函数(中文注释版):

pub(crate) fn create_test_dataframe_many_gaps() -> DataFrame {
    let start = NaiveDateTime::parse_from_str("2025-12-25 00:00:00", "%Y-%m-%d %H:%M:%S").unwrap();
    let end = NaiveDateTime::parse_from_str("2025-12-25 06:00:00", "%Y-%m-%d %H:%M:%S").unwrap();

    let mut df = DataFrame::empty();

    // 固定秒数偏移,生成非整点时间
    let second_offsets = [6, 17, 43, 51];
    // 固定分钟跳跃,制造数据间隙
    let minute_jumps = [5, 10, 30, 60, 120];

    for sensor_id in [1551] {
        let mut sensor_ids = Vec::new();
        let mut ts_values = Vec::new();
        let mut values = Vec::new();

        let mut current = start;
        let mut idx = 0usize;

        while current <= end {
            // 给当前时间添加秒数偏移
            let ts = current
                .with_second(second_offsets[idx % second_offsets.len()])
                .unwrap();

            sensor_ids.push(sensor_id);
            ts_values.push(ts);
            // 生成单调递增的测试值
            values.push(-1_464_000.0 + (idx as f64 * 100.0));

            // 跳跃前进制造间隙
            let jump = minute_jumps[idx % minute_jumps.len()];
            current += chrono::Duration::minutes(jump);
            idx += 1;
        }

        let sensor_df = DataFrame::new(vec![
            Column::from(Series::new(PlSmallStr::from_str("sensor_id"), sensor_ids)),
            Column::from(Series::new(PlSmallStr::from_str("ts"), ts_values)),
            Column::from(Series::new(PlSmallStr::from_str("value"), values)),
        ]).unwrap();

        df = df.vstack(&sensor_df).unwrap();
    }

    df
}

内容的提问来源于stack exchange,提问作者Glenn Pierce

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.11 18:44:49