如何正确使用Polars的gather_every与rolling_sum计算时序求和?
需求目标
针对频率为1秒的时间序列,为每个时间点查找1分钟、2分钟、3分钟、4分钟、5分钟前的对应秒数值,并对这5个值求和。例如:
- 23:59:59的计算值为23:59:59、23:58:59、23:57:59、23:56:59、23:55:59的数值之和;
- 23:59:58的计算值为23:59:58、23:58:58等五个对应时间点的数值之和,以此类推。
问题与尝试过程
如何用Polars正确实现上述计算?
我尝试用gather_every(60).rolling_sum(5),但触发ShapeError;之后分开调用gather_every和rolling_sum再通过concat合并,结果DataFrame首尾出现null,说明行被错误填充。输出7里,values_gather_every列的正确结果应该是顶部为null,底部有数值,但实际结果相反。
代码及输出
Code 1
import polars as pl import numpy as np from datetime import datetime, timedelta start = datetime(2023, 1, 1) seconds = 86400 end = start + timedelta(seconds=seconds-1) values = np.arange(1, seconds + 1) df = pl.DataFrame({ 'dt': pl.datetime_range(start=start, end=end, interval="1s", eager=True), 'values': values }) print(df)
Out 1
shape: (86_400, 2) ┌─────────────────────┬────────┐ │ dt ┆ values │ │ --- ┆ --- │ │ datetime[μs] ┆ i64 │ ╞═════════════════════╪════════╡ │ 2023-01-01 00:00:00 ┆ 1 │ │ 2023-01-01 00:00:01 ┆ 2 │ │ 2023-01-01 00:00:02 ┆ 3 │ │ 2023-01-01 00:00:03 ┆ 4 │ │ … ┆ … │ │ 2023-01-01 23:59:56 ┆ 86397 │ │ 2023-01-01 23:59:57 ┆ 86398 │ │ 2023-01-01 23:59:58 ┆ 86399 │ │ 2023-01-01 23:59:59 ┆ 86400 │ └─────────────────────┴────────┘
Code 2
df = df.with_columns(pl.col("values").gather_every(60).rolling_sum(5).name.suffix("_gather_every_rolling_sum"))
Out 2
ShapeError: unable to add a column of length 1440 to a dataframe of height 86400
Code 3
gather_every_df = df.select("values").gather_every(60) gather_every_df = gather_every_df.rename({"values": "values_gather_every"}) print(gather_every_df)
Out 3
shape: (1_440, 1) ┌─────────────────────┐ │ values_gather_every │ │ --- │ │ i64 │ ╞═════════════════════╡ │ 1 │ │ 61 │ │ 121 │ │ 181 │ │ 241 │ │ … │ │ 86101 │ │ 86161 │ │ 86221 │ │ 86281 │ │ 86341 │ └─────────────────────┘
Code 4
rolling_sum_df = gather_every_df.select("values_gather_every").rolling_sum(5) print(rolling_sum_df)
Out 4
AttributeError: 'DataFrame' object has no attribute 'rolling_sum'
Code 5
rolling_sum_series = gather_every_df.select("values_gather_every").to_series().rolling_sum(5) print(rolling_sum_series)
Out 5
shape: (1_440,) Series: 'values_gather_every' [i64] [ null null null null 605 905 1205 1505 1805 2105 2405 2705 … 427505 427805 428105 428405 428705 429005 429305 429605 429905 430205 430505 430805 431105 ]
Code 6
rolling_sum_df = pl.DataFrame({ 'values_gather_every_rolling_sum': rolling_sum_series }) print(rolling_sum_df)
Out 6
shape: (1_440, 1) ┌─────────────────────────────────┐ │ values_gather_every_rolling_su… │ │ --- │ │ i64 │ ╞═════════════════════════════════╡ │ null │ │ null │ │ null │ │ null │ │ 605 │ │ … │ │ 429905 │ │ 430205 │ │ 430505 │ │ 430805 │ │ 431105 │ └─────────────────────────────────┘
Code 7
df2 = pl.concat([df, gather_every_df, rolling_sum_df], how="horizontal") print(df2) print(df2.describe())
Out 7
shape: (86_400, 4) ┌─────────────────────┬────────┬─────────────────────┬─────────────────────────────────┐ │ dt ┆ values ┆ values_gather_every ┆ values_gather_every_rolling_su… │ │ --- ┆ --- ┆ --- ┆ --- │ │ datetime[μs] ┆ i64 ┆ i64 ┆ i64 │ ╞═════════════════════╪════════╪═════════════════════╪═════════════════════════════════╡ │ 2023-01-01 00:00:00 ┆ 1 ┆ 1 ┆ null │ │ 2023-01-01 00:00:01 ┆ 2 ┆ 61 ┆ null │ │ 2023-01-01 00:00:02 ┆ 3 ┆ 121 ┆ null │ │ 2023-01-01 00:00:03 ┆ 4 ┆ 181 ┆ null │ │ 2023-01-01 00:00:04 ┆ 5 ┆ 241 ┆ 605 │ │ … ┆ … ┆ … ┆ … │ │ 2023-01-01 23:59:55 ┆ 86396 ┆ null ┆ null │ │ 2023-01-01 23:59:56 ┆ 86397 ┆ null ┆ null │ │ 2023-01-01 23:59:57 ┆ 86398 ┆ null ┆ null │ │ 2023-01-01 23:59:58 ┆ 86399 ┆ null ┆ null │ │ 2023-01-01 23:59:59 ┆ 86400 ┆ null ┆ null │ └─────────────────────┴────────┴─────────────────────┴─────────────────────────────────┘ shape: (9, 5) ┌────────────┬────────────────────────────┬──────────────┬─────────────────────┬─────────────────────────────────┐ │ statistic ┆ dt ┆ values ┆ values_gather_every ┆ values_gather_every_rolling_su… │ │ --- ┆ --- ┆ --- ┆ --- ┆ --- │ │ str ┆ str ┆ f64 ┆ f64 ┆ f64 │ ╞════════════╪════════════════════════════╪══════════════╪═════════════════════╪═════════════════════════════════╡ │ count ┆ 86400 ┆ 86400.0 ┆ 1440.0 ┆ 1436.0 │ │ null_count ┆ 0 ┆ 0.0 ┆ 84960.0 ┆ 84964.0 │ │ mean ┆ 2023-01-01 11:59:59.500000 ┆ 43200.5 ┆ 43171.0 ┆ 215855.0 │ │ std ┆ null ┆ 24941.675966 ┆ 24950.19038 ┆ 124404.541718 │ │ min ┆ 2023-01-01 00:00:00 ┆ 1.0 ┆ 1.0 ┆ 605.0 │ │ 25% ┆ 2023-01-01 06:00:00 ┆ 21601.0 ┆ 21601.0 ┆ 108305.0 │ │ 50% ┆ 2023-01-01 12:00:00 ┆ 43201.0 ┆ 43201.0 ┆ 216005.0 │ │ 75% ┆ 2023-01-01 17:59:59 ┆ 64800.0 ┆ 64741.0 ┆ 323405.0 │ │ max ┆ 2023-01-01 23:59:59 ┆ 86400.0 ┆ 86341.0 ┆ 431105.0 │ └────────────┴────────────────────────────┴──────────────┴─────────────────────┴─────────────────────────────────┘
正确实现方法
核心思路是按时间的秒数分组:所有秒数相同的时间点(如XX:XX:00、XX:XX:01...)会被分到同一组,组内的元素正好是每隔60秒一个(即每分钟的同一秒)。随后在每个组内做窗口大小为5的滚动求和,即可得到当前时间点及前4个同秒数时间点的和,对应需求中当前、1分钟前、2分钟前、3分钟前、4分钟前的5个数值。
代码实现
import polars as pl import numpy as np from datetime import datetime, timedelta start = datetime(2023, 1, 1) seconds = 86400 end = start + timedelta(seconds=seconds-1) values = np.arange(1, seconds + 1) df = pl.DataFrame({ 'dt': pl.datetime_range(start=start, end=end, interval="1s", eager=True), 'values': values }) # 按秒数分组,组内滚动求和 result_df = df.with_columns( pl.col("values") .rolling_sum(window_size=5, by="dt") .over(pl.col("dt").dt.second()) .alias("sum_5_minutes_back") ) # 查看最后10行结果 print(result_df.tail(10))
关键步骤解释
pl.col("dt").dt.second():提取每个时间点的秒数作为分组键,确保同秒数的时间点处于同一组;.rolling_sum(window_size=5, by="dt"):在每个组内,按时间顺序执行窗口大小为5的滚动求和,窗口包含当前行及前4行,恰好对应当前时间点和过去4个同秒数的时间点(即1-4分钟前的对应秒);- 该方法不会出现形状不匹配问题,每个时间点都能得到正确的求和结果,前4个时间点因数据不足会显示null,符合预期。
结果验证
以最后一行2023-01-01 23:59:59为例,其求和结果应为:86400 + (86400-60) + (86400-120) + (86400-180) + (86400-240) = 431400,执行代码后可验证该值正确。
内容的提问来源于stack exchange,提问作者Niq_Lin
相关产品推荐
相关产品推荐

