如何在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);
逻辑说明
- 生成规整时间序列:按
sensor_id分组,计算每个传感器数据的时间范围,将起始时间向下取整到最近的5分钟、结束时间向上取整到最近的5分钟,然后生成该范围内所有5分钟间隔的时间点。 - 外连接合并数据:用
outer连接方式将规整时间序列和原始数据合并,这样原始的非整点记录和间隙的整点记录都会被保留,value字段为空的就是间隙需要填充的行。 - 排序:按
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
相关产品推荐
相关产品推荐

