基于Pandas/DuckDB按分区迭代计算滚动总和与交易预测
解决方案:高效填充TRANSACTIONS_FORECAST
一、DuckDB递归CTE性能优化方案
递归CTE的性能瓶颈通常源于无限制迭代或重复计算,通过以下优化可大幅提升效率:
核心优化思路
- 预计算分区内
TRANSACTIONS_WEEK总和,避免递归过程中重复求和 - 为分区内每行添加顺序编号,控制迭代仅处理下一行,缩小递归范围
- 仅对未填充的行进行迭代计算,自动终止无空值的分区
示例SQL代码
-- 1. 预处理:分区内排序并添加行号,预计算周总和 WITH preprocessed AS ( SELECT CODE, DAY, TIME, TRANSACTIONS, TRANSACTIONS_FORECAST, TRANSACTIONS_WEEK, -- 替换<排序字段>为实际行顺序字段(如时间戳) ROW_NUMBER() OVER (PARTITION BY CODE, DAY, TIME ORDER BY <排序字段>) AS row_idx, SUM(TRANSACTIONS_WEEK) OVER (PARTITION BY CODE, DAY, TIME) AS total_week FROM your_table ), -- 2. 递归填充:仅处理未填充的行 recursive_fill AS ( -- 初始种子:已有非空值的行 SELECT * FROM preprocessed WHERE TRANSACTIONS IS NOT NULL OR TRANSACTIONS_FORECAST IS NOT NULL UNION ALL -- 迭代计算:取前2个非空值的和,生成预测值 SELECT p.CODE, p.DAY, p.TIME, p.TRANSACTIONS, -- 计算step1/step2*TRANSACTIONS_WEEK (SUM(COALESCE(r.TRANSACTIONS, r.TRANSACTIONS_FORECAST)) OVER (PARTITION BY p.CODE, p.DAY, p.TIME ORDER BY r.row_idx ROWS BETWEEN 2 PRECEDING AND 1 PRECEDING)) / p.total_week * p.TRANSACTIONS_WEEK AS TRANSACTIONS_FORECAST, p.TRANSACTIONS_WEEK, p.row_idx, p.total_week FROM preprocessed p JOIN recursive_fill r ON p.CODE = r.CODE AND p.DAY = r.DAY AND p.TIME = r.TIME AND p.row_idx = r.row_idx + 1 WHERE p.TRANSACTIONS IS NULL AND p.TRANSACTIONS_FORECAST IS NULL ) -- 3. 去重取最终结果:每个行仅保留最后一次迭代的填充值 SELECT DISTINCT ON (CODE, DAY, TIME, row_idx) CODE, DAY, TIME, TRANSACTIONS, TRANSACTIONS_FORECAST, TRANSACTIONS_WEEK FROM recursive_fill ORDER BY CODE, DAY, TIME, row_idx, TRANSACTIONS_FORECAST DESC;
二、Pandas高效迭代实现方案
通过groupby按分区处理,结合滚动窗口和局部迭代,避免全表低效循环:
核心优化思路
- 按
CODE/DAY/TIME分区,每个分区独立处理 - 用滚动窗口计算前2个非空值的和,替代逐行遍历
- 每次迭代仅更新空值行,直到分区内无空值时停止
示例Python代码
import pandas as pd def fill_forecast_group(group): # 按业务逻辑排序(替换<排序字段>为实际顺序字段) group = group.sort_values(by='<排序字段>').reset_index(drop=True) # 预计算分区内TRANSACTIONS_WEEK总和 total_week = group['TRANSACTIONS_WEEK'].sum() # 创建合并非空值的临时列 group['combined'] = group['TRANSACTIONS'].fillna(group['TRANSACTIONS_FORECAST']) # 迭代填充直到无空值 while group['combined'].isnull().any(): # 计算前2个非空值的滚动和,偏移到当前空值行 rolling_sum = group['combined'].rolling(window=3, min_periods=2).sum().shift(-1) # 筛选空值行并计算预测值 null_mask = group['combined'].isnull() group.loc[null_mask, 'TRANSACTIONS_FORECAST'] = (rolling_sum[null_mask] / total_week) * group.loc[null_mask, 'TRANSACTIONS_WEEK'] # 更新合并列 group['combined'] = group['TRANSACTIONS'].fillna(group['TRANSACTIONS_FORECAST']) # 清理临时列 return group.drop(columns=['combined']) # 按分区分组处理 result_df = df.groupby(['CODE', 'DAY', 'TIME'], group_keys=False).apply(fill_forecast_group)
内容的提问来源于stack exchange,提问作者AK91
相关产品推荐
相关产品推荐

