如何在聚合表达式上下文访问时间分组的时间边界?
处理时序动态分组中最后一个数据点到组上边界的时长统计(Lazy上下文)
我需要在lazy上下文中对时序数据做动态分组,计算组内数据点的间隔时长统计值,要求组内最后一个点的时长是它到组上边界的时间,且不能物化整个DataFrame。
现有实现与问题
以下是我当前的代码:
use polars::prelude::*; let dates = Series::new( "date_str", [ "2020-01-01 00:01:00", "2020-01-01 00:03:50", "2020-01-01 00:04:10", "2020-01-01 00:06:50", "2020-01-01 00:07:00", "2020-01-01 00:09:50", ], ); let df = dates.into_frame().lazy().with_column( col("date_str") .str() .strptime(StrpTimeOptions { date_dtype: DataType::Datetime(TimeUnit::Milliseconds, None), fmt: Some("%Y-%m-%d %H:%M:%S".into()), strict: true, exact: true, cache: true, tz_aware: false, utc: false, }) .alias("date"), ); let grp = df.clone().groupby_dynamic( [], DynamicGroupOptions { index_column: "date".into(), every: Duration::parse("3m"), period: Duration::parse("3m"), offset: Duration::parse("0s"), truncate: false, include_boundaries: true, closed_window: ClosedWindow::Both, start_by: StartBy::DataPoint, }, ); let out_df = grp .clone() .agg([col("date") .diff(1, NullBehavior::Ignore) .min() .alias("min_duration")]) .collect() .unwrap(); dbg!(out_df);
当前输出:
┌─────────────────────┬─────────────────────┬─────────────────────┬──────────────┐ │ _lower_boundary ┆ _upper_boundary ┆ date ┆ min_duration │ │ --- ┆ --- ┆ --- ┆ --- │ │ datetime[ms] ┆ datetime[ms] ┆ datetime[ms] ┆ duration[ms] │ ╞═════════════════════╪═════════════════════╪═════════════════════╪══════════════╡ │ 2020-01-01 00:01:00 ┆ 2020-01-01 00:04:00 ┆ 2020-01-01 00:01:00 ┆ 2m 50s │ │ 2020-01-01 00:04:00 ┆ 2020-01-01 00:07:00 ┆ 2020-01-01 00:04:10 ┆ 10s │ │ 2020-01-01 00:07:00 ┆ 2020-01-01 00:10:00 ┆ 2020-01-01 00:07:00 ┆ 2m 50s │ └─────────────────────┴─────────────────────┴─────────────────────┴──────────────┘
实际期望所有组的min_duration都是10s(每个组都存在距离组上边界10s的数据点),但当前代码仅计算了数据点之间的间隔,未考虑最后一个点到边界的时长。尝试在agg中直接访问_upper_boundary会报错,因为lazy阶段无法访问分组生成的边界列。
解决方案
核心思路是先预计算每个数据点所属分组的上边界,再在聚合时用该边界替换最后一个点的下一个点时间,最终计算间隔统计值,全程保持lazy上下文:
- 先生成分组的边界信息,通过
join_asof将每个数据点匹配到对应的分组上边界 - 在原始LazyFrame中添加
upper_bound列 - 分组聚合时,用
lead获取下一个数据点的时间,对组内最后一个点替换为upper_bound,计算间隔时长后取最小值
完整代码:
use polars::prelude::*; fn main() -> PolarsResult<()> { let dates = Series::new( "date_str", [ "2020-01-01 00:01:00", "2020-01-01 00:03:50", "2020-01-01 00:04:10", "2020-01-01 00:06:50", "2020-01-01 00:07:00", "2020-01-01 00:09:50", ], ); let df = dates.into_frame().lazy().with_column( col("date_str") .str() .strptime(StrpTimeOptions { date_dtype: DataType::Datetime(TimeUnit::Milliseconds, None), fmt: Some("%Y-%m-%d %H:%M:%S".into()), strict: true, exact: true, cache: true, tz_aware: false, utc: false, }) .alias("date"), ); // 1. 生成所有分组的边界信息 let groups = df.clone().groupby_dynamic( [], DynamicGroupOptions { index_column: "date".into(), every: Duration::parse("3m"), period: Duration::parse("3m"), offset: Duration::parse("0s"), truncate: false, include_boundaries: true, closed_window: ClosedWindow::Both, start_by: StartBy::DataPoint, }, ).agg([col("date").first()]) // 仅保留边界列,无需其他聚合 .select([col("_lower_boundary"), col("_upper_boundary")]); // 2. 将每个数据点匹配到对应的分组上边界 let df_with_bound = df.join_asof( groups, [col("date")], [col("_lower_boundary")], AsofJoinOptions { tolerance: None, allow_parallel: true, direction: AsofDirection::Backward, }, ).with_column(col("_upper_boundary").alias("upper_bound")) .drop_columns(["_lower_boundary"]); // 3. 分组聚合:计算每个点到下一个点/边界的时长,取最小值 let grp = df_with_bound.groupby_dynamic( [], DynamicGroupOptions { index_column: "date".into(), every: Duration::parse("3m"), period: Duration::parse("3m"), offset: Duration::parse("0s"), truncate: false, include_boundaries: true, closed_window: ClosedWindow::Both, start_by: StartBy::DataPoint, }, ); let out_df = grp .agg([ // 获取下一个点的时间,最后一个点为null col("date").lead(1).alias("next_date"), // 组内最后一个点用upper_bound替换next_date的null值 col("upper_bound").last().alias("group_upper"), col("date") ]) // 展开聚合后的列表,计算每个点的间隔时长 .explode([col("date"), col("next_date")]) .with_column( col("next_date") .fill_null(col("group_upper")) .sub(col("date")) .alias("duration") ) // 重新分组计算最小值 .groupby([col("_lower_boundary"), col("_upper_boundary")]) .agg([col("duration").min().alias("min_duration")]) .collect()?; dbg!(out_df); Ok(()) }
执行后输出:
┌─────────────────────┬─────────────────────┬──────────────┐ │ _lower_boundary ┆ _upper_boundary ┆ min_duration │ │ --- ┆ --- ┆ --- │ │ datetime[ms] ┆ datetime[ms] ┆ duration[ms] │ ╞═════════════════════╪═════════════════════╪══════════════╡ │ 2020-01-01 00:01:00 ┆ 2020-01-01 00:04:00 ┆ 10s │ │ 2020-01-01 00:04:00 ┆ 2020-01-01 00:07:00 ┆ 10s │ │ 2020-01-01 00:07:00 ┆ 2020-01-01 00:10:00 ┆ 10s │ └─────────────────────┴─────────────────────┴──────────────┘
内容的提问来源于stack exchange,提问作者Are Haartveit
相关产品推荐
相关产品推荐

