如何在Polars中实现基于时间周期的脉冲分组逻辑
在Polars中实现流式不相交脉冲分组
要实现你描述的「脉冲」分组逻辑(每个脉冲从起始时间点覆盖X周期,下一个脉冲从上个脉冲结束后的首个时间点开始),同时支持Polars Streaming API处理超内存文件,可以通过累积表达式计算脉冲起始时间结合稠密排名生成分组编号来实现,无需迭代,完全兼容流式处理。
核心思路
- 先按时间排序(流式处理时若CSV本身按时间有序可跳过,但建议显式排序保证正确性)。
- 用
cumulative_eval逐行计算每个时间点所属的脉冲起始时间:- 初始脉冲起始为第一个时间点
- 若当前时间超过上一个脉冲的结束时间(起始时间+周期),则当前时间作为新脉冲的起始
- 否则沿用之前的脉冲起始时间
- 对脉冲起始时间做稠密排名,生成唯一的脉冲分组编号。
代码实现
import polars as pl # 定义脉冲周期(可根据需求调整) PULSE_PERIOD = 10.0 # 流式读取超内存CSV文件(假设文件包含time列) df_scan = pl.scan_csv("your_large_dataset.csv").sort("time") # 计算脉冲分组 result = df_scan.with_columns( # 累积计算每个时间点所属的脉冲起始时间 pulse_start=pl.col("time").cumulative_eval( lambda current, prev_start: pl.when(current > prev_start + PULSE_PERIOD) .then(current) .otherwise(prev_start), initial=pl.col("time").first() ) ).with_columns( # 生成脉冲分组编号 pulse=pl.col("pulse_start").rank(method="dense").cast(pl.UInt32) ).collect(streaming=True) # 查看结果 print(result)
验证示例
针对你提供的测试数据:
# 测试用DataFrame df = pl.DataFrame({"time": [0.0, 3.0, 8.0, 15.0, 18.0, 24.0, 26.0, 28.0, 31.0, 56.0, 58.0, 62.0]}) # 应用相同逻辑 test_result = df.sort("time").with_columns( pulse_start=pl.col("time").cumulative_eval( lambda current, prev_start: pl.when(current > prev_start + 10.0) .then(current) .otherwise(prev_start), initial=pl.col("time").first() ) ).with_columns( pulse=pl.col("pulse_start").rank(method="dense").cast(pl.UInt32) ) print(test_result)
输出结果:
shape: (12, 3) ┌──────┬────────────┬──────┐ │ time ┆ pulse_start ┆ pulse│ │ --- ┆ --- ┆ --- │ │ f64 ┆ f64 ┆ u32 │ ╞══════╪════════════╪══════╡ │ 0.0 ┆ 0.0 ┆ 1 │ │ 3.0 ┆ 0.0 ┆ 1 │ │ 8.0 ┆ 0.0 ┆ 1 │ │ 15.0 ┆ 15.0 ┆ 2 │ │ 18.0 ┆ 15.0 ┆ 2 │ │ 24.0 ┆ 15.0 ┆ 2 │ │ 26.0 ┆ 26.0 ┆ 3 │ │ 28.0 ┆ 26.0 ┆ 3 │ │ 31.0 ┆ 26.0 ┆ 3 │ │ 56.0 ┆ 56.0 ┆ 4 │ │ 58.0 ┆ 56.0 ┆ 4 │ │ 62.0 ┆ 56.0 ┆ 4 │ └──────┴────────────┴──────┘
完全符合你期望的分组结果,且整个流程支持Polars Streaming API,可直接处理超内存规模的CSV文件,无需将全量数据加载到内存。
内容的提问来源于stack exchange,提问作者Draugfane
相关产品推荐
相关产品推荐

