在数据分析工作中,交易流水表通常包含多个业务类型,例如充值、消费、退款等。针对每一种类型,我们往往希望得到当前记录之前最近一次同类型交易的金额,用来计算环比、校验异常或补全特征。如果直接写双层循环去查找,不仅代码难维护,而且在百万级数据上几乎无法接受。借助DataFrame的分组与偏移能力,可以把这类问题转化为一次向量化运算。

理解分组偏移的基本思路
所谓前一笔交易金额,本质上是按类型分组后,在时间有序的前提下取同组上一行的数值。DataFrame提供了groupby结合shift的能力,能够在每个分组内部进行行位移。这里最关键的前提是数据已经按照交易时间和类型排好顺序,否则位移得到的结果将没有业务意义。
很多人在第一次实现时会尝试用merge做自连接,把同一类型的前一条记录通过最大时间小于当前时间的方式关联过来。这种做法在逻辑上正确,但会产生巨大的中间表,并且自连接的条件判断非常消耗资源。分组偏移则只在原有数据上增加一列,不需要复制行,效率差异在大数据量下非常明显。
我们用一个简单例子说明。假设有类型和金额两列,先排序再分组位移即可:
import pandas as pd
data = {
'type': ['recharge', 'pay', 'recharge', 'pay', 'recharge'],
'time': ['2023-01-01', '2023-01-02', '2023-01-03', '2023-01-04', '2023-01-05'],
'amount': [100, 50, 200, 80, 150]
}
df = pd.DataFrame(data)
df['time'] = pd.to_datetime(df['time'])
df = df.sort_values(['type', 'time'])
df['prev_amount'] = df.groupby('type')['amount'].shift(1)
print(df)
上面的代码在recharge组内,第二条记录的前一笔金额就是第一条的100,而第一条因为没有前一笔,结果自然为NaN。这种方式把复杂查找简化成了单列运算,是后续所有优化的基础。
处理多类型与缺失值的细节
真实业务中的类型字段可能不止两种,而且同一类型可能出现时间缺失或重复时间。面对重复时间,排序会变得不稳定,此时可以加入二级排序键,例如交易流水号,确保顺序唯一。对于缺失的时间戳,应先做清洗再排序,否则shift仍会按行号位移,但业务含义已经偏移。
当某些类型只有一条记录时,shift(1)产生的NaN是合理现象,但在后续计算中可能干扰统计。我们可以用fillna把缺失的前一笔金额补为零,或者单独标记为新客首笔。需要注意的是,填充策略必须写在分组位移之后,避免把NaN误填进分组逻辑中。
如果还需要前第两笔、前三笔,只需多次shift并命名不同列,或者使用shift的数组形式。下面的示例展示如何一次取出前一笔和前两笔:
df['prev_1'] = df.groupby('type')['amount'].shift(1)
df['prev_2'] = df.groupby('type')['amount'].shift(2)
# 对于首笔类型记录,prev_1和prev_2均为NaN
print(df[['type', 'amount', 'prev_1', 'prev_2']])
通过这种写法,每一个类型的时序特征都能被平行扩展,而不必为每种距离单独写查询。同时因为所有列都来源于同一分组对象,计算之间互不干扰,也方便做交叉校验。
用窗口函数验证与性能对比
为了保证分组偏移的结果正确,可以选取一个类型子集,用排序后的iloc手动取前一行进行抽查。另外在支持窗口函数的计算引擎中,例如Spark或SQL,也可以用lag函数实现同样逻辑,再与pandas结果比对。这样能排除因排序错乱导致的隐性错误。
从性能角度看,分组偏移属于单遍扫描操作,时间复杂度接近线性。而自连接方案在最差情况下复杂度会升至平方级。我们在十万行混合类型数据上测试,groupby加shift通常能在毫秒级完成,自连接则可能达到秒级。对于更大规模,还可以把类型字段转为类别型,减少分组时的哈希开销。
最后给出一段较完整的处理模板,涵盖排序、位移与缺失标记,方便直接套用到交易类数据:
import pandas as pd
def add_prev_amount(df, type_col='type', time_col='time', amt_col='amount'):
df = df.copy()
df[time_col] = pd.to_datetime(df[time_col])
df = df.sort_values([type_col, time_col])
df['prev_amount'] = df.groupby(type_col)[amt_col].shift(1)
df['is_first'] = df['prev_amount'].isna()
return df
# 调用示例
cleaned = add_prev_amount(df)
print(cleaned.head())
把上述函数放入预处理管道后,无论上游类型怎么增加,前一笔交易金额的获取都不再需要改写核心逻辑,对于日常特征工程是非常实用的写法。