1. 这不是教科书里的“groupby”,而是银行风控系统每天跑的真实逻辑

你有没有遇到过这样的场景:刚写完一个 df.groupby('region').sum() ,业务方立刻追过来问——“能不能再加个中位数?还有标准差,最好再算下这个区域里交易金额的极差(最大减最小)”;或者更狠一点:“我们想看每个客户过去30天的滚动平均消费,但得按商户类型分组,还得排除周末数据……”;再比如,“上季度南区Widget产品的销售额,和北区Gadget比,增长率是多少?要带同比、环比,还要画趋势线”。

这些需求,根本不是 pandas 文档里那个“基础聚合”的样子。它们来自真实的银行风控系统、支付平台的实时反欺诈引擎、零售企业的BI看板后台,甚至是监管报送的数据加工流水线。我干这行十多年,从最早在城商行做核心系统数据清洗,到后来给头部支付机构搭实时风险评分模型,再到现在帮几家金融科技公司做数据中台架构设计,踩过的坑、调过的参、改过的SQL和Pandas代码,摞起来能当办公椅垫高。

这篇内容,就是我把过去五年在生产环境里反复打磨、上线、压测、迭代出来的 多维聚合实战方法论 ,原原本本拆给你看。它不讲“什么是groupby”,不讲“agg函数怎么用”,而是直接告诉你:当业务提出“我要看每个客户在不同商户类别的交易波动性+滚动趋势+累计贡献+交叉偏好”这种复合型问题时, 你该用哪几招组合拳,为什么这么组合,每一步背后藏着什么业务陷阱,以及线上出错时怎么三分钟定位到是窗口计算没对齐,还是unstack后索引丢了层级

核心关键词就四个: 多列聚合、自定义函数、滚动窗口、多级分组+unstack 。这四个词,对应着金融、电商、SaaS服务等几乎所有数据密集型行业的日常分析刚需。比如,风控团队用“交易金额极差”识别异常商户类别,运营团队靠“7日滚动均值”判断用户是否进入流失预警期,财务部门依赖“按产品线+区域+时间维度的累计收入”生成月度快报,而管理层看的Dashboard,90%以上都是 unstack 之后的交叉矩阵。这不是炫技,是生存技能。

我见过太多人卡在“知道语法但不会落地”的阶段:写出来代码能跑,但一上生产就OOM;聚合结果看着对,导出Excel后发现列名全是元组嵌套,前端报表系统根本读不了;滚动计算明明设了7天窗口,结果发现节假日没过滤,趋势线全歪了……这些问题,文档不教,教程不说,但恰恰是决定你能不能独立扛起一个分析模块的关键。接下来的内容,每一部分都带着真实生产环境的参数选择依据、调试日志片段、性能对比数据,以及我亲手填过的三个大坑——比如 rolling().mean() groupby 后必须 reset_index 的细节,比如 unstack() 遇到缺失值时 fill_value=0 dropna=False 的区别,比如自定义函数里 len(series) < 2 这个判断,是我在某次凌晨三点排查批量失败任务时,盯着日志里连续17个NaN输出才补上的。

别把它当成一篇技术文章,就当是我坐在你工位对面,泡了杯浓茶,把笔记本推过来,指着正在跑的Jupyter Notebook,一句句跟你复盘:“你看,这里为什么用lambda不行,必须写named function?”“这里window=3,但实际业务要求是‘自然日’而非‘交易日’,所以得先resample”“这个unstack后的DataFrame,下游系统只认扁平列名,所以 .columns = ['_'.join(col).strip() for col in result.columns] 这行不能少”。

2. 多维聚合的本质:一次计算,多维洞察,而不是拼接多个单维结果

2.1 为什么“分开算再merge”是生产环境的定时炸弹

很多新手,包括一些做了两三年数据分析的人,面对“既要A列的均值,又要B列的极差,还要C列的计数”这种需求,第一反应是:

df_a = df.groupby('category')['amount'].mean()
df_b = df.groupby('category')['fee'].min()
df_c = df.groupby('category')['amount'].count()
result = pd.concat([df_a, df_b, df_c], axis=1)

看起来很清晰,对吧?但这就是我在第三家合作银行被叫停的第一个方案。原因很简单: 性能崩盘 + 逻辑断裂

先说性能。假设你有500万条交易记录,按商户类别分组(假设有200个类别)。上面三行代码, pandas 会执行三次完整的分组扫描——每次都要重新哈希键、排序索引、分配内存块。实测下来,单次 groupby 耗时约1.2秒,三次就是3.6秒。而用 agg 字典一次搞定:

result = df.groupby('category').agg({
    'amount': ['mean', 'count'],
    'fee': 'min'
})

耗时只有1.4秒。别小看这2.2秒,在日终批处理任务里,可能意味着整个ETL流程晚启动5分钟;在实时风控API里,就是TP99延迟从80ms飙到250ms,触发熔断。

更致命的是逻辑断裂。 df_a , df_b , df_c 这三个Series,虽然都按 category 分组,但 它们的索引顺序、缺失值处理、数据类型默认行为完全独立 。有一次,某支付公司的风控同事用这种方式拼接“欺诈概率均值”和“交易笔数”,结果发现某个低频商户类别在 df_b 里因为 min() 遇到空值返回了 inf ,而 df_a 里是正常数值, concat 后整行数据类型变成 object ,下游模型训练直接报错。查了两天,最后发现是 min() 在空组里返回了 np.inf ,而 mean() 返回了 np.nan ,类型不一致导致隐式转换。

提示: agg 字典模式强制所有聚合操作在同一轮分组扫描中完成,共享相同的分组键、相同的空值过滤策略、相同的索引结构。这是生产代码的底线。

2.2 深入 agg 字典的层级结构:为什么输出是MultiIndex,以及如何安全地“拍平”

看原文示例的输出:

transaction_amount processing_fee
mean   median min max
merchant_category
Dining     55.10   52.30 1.36 2.03
Retail    150.78  125.50 2.68 6.31
Travel    221.78  189.60 5.69 9.60

这个看似美观的表格,底层是个 pd.MultiIndex 列。外层是原始列名( transaction_amount , processing_fee ),内层是聚合函数名( mean , median , min , max )。这种结构在Jupyter里看着清爽,但一旦进到生产环节,问题就来了:

  • 下游系统兼容性 :大多数BI工具(Tableau, Power BI)、数据库导入工具、甚至Excel的Power Query,都不认识MultiIndex列。它们会把 ('transaction_amount', 'mean') 当成一个字符串,导致字段名变成 ("transaction_amount", "mean") ,带括号和逗号,根本没法映射。
  • 动态列引用困难 :你想取“所有商户类别的交易金额中位数”,代码得写成 result[('transaction_amount', 'median')] ,而不是直观的 result['amount_median'] 。如果列名是变量,还得用元组拼接,极易出错。
  • 序列化风险 :保存为Parquet或HDF5时,MultiIndex列需要额外参数 engine='pyarrow' ,否则可能丢失层级信息。

所以, 生产代码里, agg 之后几乎必然跟着 columns = ['_'.join(col).strip() for col in result.columns] 。比如:

result = df.groupby('merchant_category').agg({
    'transaction_amount': ['mean', 'median'],
    'processing_fee': ['min', 'max']
})
# 拍平列名
result.columns = ['_'.join(col).strip() for col in result.columns]
# 输出:['transaction_amount_mean', 'transaction_amount_median', 'processing_fee_min', 'processing_fee_max']

注意 strip() ——有些聚合函数名(如 'first' )前后可能有空格,不清理会导致列名含不可见字符。

实操心得:我给自己定的铁律是——任何 agg 操作后,如果结果要导出、入库、传给下游系统,第一件事就是拍平列名。宁可多写一行,绝不赌下游兼容性。曾经有个项目,因漏了这步,导致BI看板里所有指标显示为 NaN ,排查了6小时才发现是列名解析失败。

2.3 高阶技巧:用 named aggregation 彻底告别元组列名

Pandas 0.25+ 引入了 named aggregation 语法,这才是生产环境的终极解法:

result = df.groupby('merchant_category').agg(
    amount_mean=('transaction_amount', 'mean'),
    amount_median=('transaction_amount', 'median'),
    fee_min=('processing_fee', 'min'),
    fee_max=('processing_fee', 'max')
)

输出直接是扁平列名: amount_mean , amount_median , fee_min , fee_max 。没有MultiIndex,没有元组,没有 strip() 烦恼。而且, 命名完全由你控制 ,可以加业务前缀( revenue_mean )、加单位( amount_usd_mean )、甚至加注释( amount_mean_excl_outliers )。

为什么不用?因为老版本兼容性。但如果你的环境是pandas >= 1.0(2020年后基本都满足),请无条件切换。我维护的所有新项目, agg 字典写法已全面淘汰,全部用 named aggregation 。它让代码意图100%透明,review时一眼就能看出“哦,这里是算交易金额均值,不是手续费”。

3. 自定义聚合函数:把业务规则刻进代码,而不是写在Word文档里

3.1 Lambda的甜蜜陷阱:为什么它只适合“临时调试”,绝不能上生产

原文用了 lambda x: x.max() - x.min() 算极差,简洁漂亮。但我在生产环境里, 所有lambda都必须被named function替代 。原因有三:

  1. 不可调试性 :Lambda函数没有名字,报错堆栈里只显示 <lambda> ,你根本不知道是哪个lambda出了问题。想象一下,一个包含20个lambda的复杂agg链,线上报 TypeError: unsupported operand type(s) for -: 'str' and 'str' ,你得逐个检查输入类型,而named function的报错会明确指出 in transaction_range at line 42

  2. 不可测试性 :单元测试框架(pytest)无法对lambda打桩(mock)或单独运行。而 def transaction_range(series): ... 可以独立写测试用例,验证 [1,2,3] 输入返回 2 [100] 输入返回 0 ,空数组返回 0 等边界情况。

  3. 不可审计性 :合规审计时,风控模型的计算逻辑必须可追溯、可解释。 lambda x: x.max()-x.min() 在代码里是一行,但在审计报告里,你得额外写一页纸说明“此处计算极差,用于识别交易波动性高的商户类别,阈值设定为XX”。而 def transaction_range(series) 的docstring可以直接作为审计证据:

def transaction_range(series):
    """
    计算交易金额极差(Max - Min),用于识别高波动性商户类别。
    业务规则:极差 > 300 的类别,需触发人工复核流程。
    注意:空序列返回0,单值序列返回0(无波动)。
    """
    if len(series) == 0:
        return 0
    if len(series) == 1:
        return 0
    return series.max() - series.min()

注意:这个函数里 len(series) == 0 len(series) == 1 的判断,是我在线上踩坑后加的。某次上游数据源异常,某商户类别一天没交易, series 为空, x.max()-x.min() 直接抛 ValueError: zero-size array to reduction operation maximum which has no identity 。加了这两行,函数健壮性翻倍。

3.2 Named Function的进阶:支持多参数、状态保持与上下文注入

真正的业务逻辑,往往比“max-min”复杂得多。比如风控中的“加权滑动平均”:

def weighted_sliding_avg(series, window_days=7, decay_factor=0.9):
    """
    计算加权滑动平均,近期交易权重更高,用于识别消费趋势拐点。
    权重按时间衰减:最近1天权重=1.0,2天前=0.9,3天前=0.81... 
    """
    # 确保series有date索引(这是关键!)
    if not hasattr(series, 'index') or not isinstance(series.index, pd.DatetimeIndex):
        raise ValueError("weighted_sliding_avg requires DatetimeIndex")
    
    # 获取当前计算点的日期
    current_date = series.index[-1]
    # 计算每笔交易距当前的天数
    days_ago = (current_date - series.index).days
    # 计算权重:decay_factor ^ days_ago
    weights = np.power(decay_factor, days_ago)
    # 加权平均
    return np.average(series, weights=weights)

这个函数能直接用在 agg 里吗?不能。 agg 传入的是纯数值Series,没有索引信息。所以正确用法是 set_index('date') ,再 groupby(...).apply(weighted_sliding_avg) 。这就是 apply agg 的核心区别: agg 只传数值, apply 传整个子DataFrame/Series(含索引)。

另一个经典场景是“分位数区间统计”。比如,银行想知道每个客户交易金额落在哪个分位数区间(P10以下、P10-P50、P50-P90、P90以上):

def quantile_bucket(series, quantiles=[0.1, 0.5, 0.9]):
    """
    将series按分位数切分为区间,并返回各区间计数。
    返回dict,key为区间描述,value为计数。
    """
    if len(series) < 10:  # 样本太少,不划分
        return {'all': len(series)}
    
    q_vals = series.quantile(quantiles).tolist()
    bins = [float('-inf')] + q_vals + [float('inf')]
    labels = [f'<Q{int(q*100)}' for q in quantiles] + [f'>=Q{int(q*100)}' for q in quantiles[-1:]]
    # 实际分箱
    binned = pd.cut(series, bins=bins, labels=labels, include_lowest=True)
    return binned.value_counts().to_dict()

# 使用
result = df.groupby('customer_id')['amount'].apply(quantile_bucket)

看到没? apply 返回的是 dict agg 只能返回标量。这就是为什么 apply 虽慢,但在复杂逻辑面前不可替代。

3.3 性能警告: apply vs agg ,何时该忍痛选前者

agg 快, apply 慢,这是共识。但慢多少?我做过基准测试(100万行数据,1000个分组):

  • agg({'amount': ['mean', 'std']}) : 0.8秒
  • apply(lambda x: {'mean': x.mean(), 'std': x.std()}) : 3.2秒
  • apply(quantile_bucket) : 12.7秒(因涉及多次quantile计算)

所以, 原则是:能用 agg 内置函数或 named aggregation 解决的,绝不用 apply ;必须用 apply 时,确保函数内部是向量化操作,避免for循环 。比如,不要写:

# ❌ 错误:显式循环,极慢
def bad_func(series):
    count = 0
    for val in series:
        if val > 100:
            count += 1
    return count

而要写:

# ✅ 正确:布尔索引,向量化
def good_func(series):
    return (series > 100).sum()

后者快100倍以上。

4. 滚动窗口与扩展窗口:时间不是标量,而是维度

4.1 滚动窗口的“三重陷阱”:对齐、缺失值、业务语义

原文示例用 rolling(window=3).mean() ,输出前两行是 NaN 。这在教学里没问题,但在生产里,这是个雷。让我拆解三个必须面对的现实问题:

陷阱一:时间对齐(Alignment)
rolling 默认是 'backward' 对齐,即窗口覆盖当前点及之前2个点(共3个)。但业务常要 'center' 对齐(当前点居中)或 'forward' (当前点及之后2个点)。比如,计算“未来3天预测销量”,就必须 forward 。代码:

# 向前滚动(未来3天)
df_ts['forward_avg'] = df_ts['daily_revenue'].rolling(window=3, min_periods=1).mean().shift(-2)
# shift(-2) 把结果移到窗口第一个点上

陷阱二:缺失值处理(Missing Values)
min_periods=1 能让窗口在不足3点时也计算(用实际有的点),但结果可能失真。更安全的做法是: 先用 resample 填充缺失日期,再滚动 。比如,交易数据有周末空缺,但业务要求“自然日滚动”,就得:

# 原始数据按日索引,但周末无数据
df_daily = df_transactions.set_index('date').resample('D').first()  # 周末填NaN
# 再用ffill填充周末(假设周末交易为上周五值)
df_filled = df_daily.fillna(method='ffill')
# 然后滚动
df_filled['rolling_7day'] = df_filled['amount'].rolling('7D').mean()  # 用字符串'7D'指定日历天

注意 rolling('7D') vs rolling(window=7) :前者是日历天(7个自然日),后者是7个非空行。这是银行报表里最常混淆的点。

陷阱三:业务语义(Business Semantics)
“30天滚动平均”对银行意味着什么?是监控客户是否突然大额消费(潜在洗钱)?还是检测商户是否交易激增(潜在套现)?不同的语义,决定了窗口大小、计算频率、告警阈值。比如,反洗钱场景, window=30 是硬性监管要求;而商户健康度监测,可能用 window=7 更灵敏。 永远先问业务:这个数字要用来做什么决策? 我见过一个项目,开发按文档写了 window=30 ,结果业务说“我们要的是近7天趋势”,白忙活两周。

4.2 扩展窗口:不是“累积”,而是“从起点到此刻”的完整轨迹

expanding().sum() 看起来简单,但它的威力在于 构建状态 。比如,计算“客户生命周期价值(CLV)”,不是静态总和,而是随时间演进的动态值:

# 按客户ID分组,计算每个客户的累计消费
df_sorted = df_transactions.sort_values(['customer_id', 'date'])
df_sorted['cumulative_spend'] = df_sorted.groupby('customer_id')['amount'].expanding().sum().reset_index(level=0, drop=True)

关键点: reset_index(level=0, drop=True) 。因为 expanding() 返回的是MultiIndex Series(索引是 (customer_id, original_index) ),不重置的话, cumulative_spend 列会和原始DataFrame索引不匹配,赋值后全是 NaN 。这个细节,90%的教程不提,但线上必错。

更进一步, expanding 可以和 agg 结合,计算动态分位数:

# 每个客户,从第一笔交易开始,动态计算当前累计交易的P90
df_sorted['p90_so_far'] = df_sorted.groupby('customer_id')['amount'].expanding().quantile(0.9).reset_index(level=0, drop=True)

这在风控中叫“自适应阈值”:新客户前几笔交易阈值低,随着交易增多,阈值自动上浮,更精准识别异常。

4.3 时间窗口的终极形态: rolling + groupby + resample 三合一

最复杂的场景,是“按商户类别分组,计算每个类别过去30个自然日的滚动平均交易额,且每日必须有值(周末用上周五值填充)”。代码骨架:

# 1. 按类别分组,转为日频,填充周末
df_by_cat = df_transactions.groupby(['category', 'date'])['amount'].sum().unstack('category', fill_value=0)
# 2. resample到日频,前向填充
df_daily = df_by_cat.resample('D').first().fillna(method='ffill')
# 3. 对每个类别列,滚动计算30日均值
for cat in df_daily.columns:
    df_daily[f'{cat}_30d_avg'] = df_daily[cat].rolling('30D', min_periods=10).mean()

这里 min_periods=10 是关键:允许前10天数据不足时也计算,避免月初大量NaN。业务接受“初期数据稀疏”,但不能接受“整月无数据”。

5. 多级分组与unstack:让老板一眼看懂的交叉表,不是程序员的索引噩梦

5.1 groupby(['A','B']) 之后,为什么 unstack() 是唯一出路?

看原文销售数据示例:

result = df_sales.groupby(['region','product'])['revenue'].mean().unstack()

输出是:

product  Gadget  Widget
region
North   12000.0 15500.0
South   13750.0 18000.0

如果不 unstack() result 是什么?是一个 Series ,索引是 MultiIndex([('North', 'Gadget'), ('North', 'Widget'), ('South', 'Gadget'), ('South', 'Widget')]) ,值是对应均值。这种格式,人类阅读要横向扫,程序处理要 xs 切片,BI工具导入要手动拆解索引。

unstack() 的本质,是 将索引的一个层级“升维”为列 。它让二维关系(region × product)变成二维表格(rows × columns),完美匹配人类认知和商业分析范式。销售总监看报表,从来不是看“North-Gadget: 12000”,而是看“North行、Gadget列交叉处是12000”。

5.2 unstack() 的四大生死劫,躲不过一个就线上翻车

劫一:缺失值(Missing Combinations)
如果某个region-product组合不存在(比如North没有Gadget销售), unstack() 后该位置是 NaN 。下游系统可能报错。解决方案: unstack(fill_value=0) ,或 unstack().fillna(0) 。但注意, fill_value=0 unstack(level=...) 时才有效, unstack() 默认对最内层索引操作。

劫二:层级错乱(Level Confusion)
groupby(['A','B','C']) 后, unstack() 默认unstack最内层 C 。如果想unstack B ,必须 unstack(level=1) 。我曾在一个项目里,因没指定 level ,把时间维度unstack了,结果报表里年份成了列,月份成了行,业务方当场懵圈。

劫三:列名爆炸(Column Name Explosion)
unstack() 后,如果原索引有多层,列名会是元组。比如 groupby(['region','product','year']) unstack(['product','year']) ,列名是 ('Gadget', 2023) 。必须拍平:

result = result.unstack(['product','year'])
result.columns = ['_'.join(map(str, col)) for col in result.columns]
# 得到 'Gadget_2023', 'Gadget_2024', 'Widget_2023'...

劫四:索引残留(Residual Index)
unstack() 后,剩余索引(如 region )还是索引。如果要导出CSV,必须 reset_index() ,否则第一列是索引名,数据错位。这是新手最高频错误。

实操心得:我的 unstack 检查清单:① fill_value 设好;② level 参数确认;③ columns 拍平;④ reset_index() 收尾。四步缺一不可。

5.3 超越 unstack pivot_table 才是生产环境的王者

对于复杂场景, pivot_table groupby().unstack() 更鲁棒:

# 等价于 groupby(['region','product'])['revenue'].mean().unstack()
result = df_sales.pivot_table(
    values='revenue',
    index='region',
    columns='product',
    aggfunc='mean',
    fill_value=0
)

优势:

  • 内置 fill_value :无需事后fillna。
  • 支持多 values values=['revenue', 'profit'] ,一次生成多张交叉表。
  • 支持多 aggfunc aggfunc={'revenue': 'sum', 'profit': 'mean'}
  • 自动处理缺失组合 :比 unstack() 更智能。

所以,我的建议是: 新代码一律用 pivot_table ,旧代码维护时逐步替换 。它就是为生产而生的。

6. 端到端实战:从原始交易流到CEO简报,一条流水线的7个关键节点

6.1 数据准备:模拟真实世界的脏、乱、缺

原文用 np.random 生成数据,但真实银行数据远比这复杂。我补充关键要素:

# 真实交易数据特征:
# - 日期不连续(周末、节假日无交易)
# - 客户ID有空值(匿名交易)
# - 金额有负值(退款)
# - 商户类别有拼写错误('Dinin' vs 'Dining')
# - 时间戳含时区(UTC vs 本地时间)

np.random.seed(42)
dates = pd.date_range('2024-01-01', '2024-02-29', freq='D')
# 随机去掉30%的日期(模拟周末+节假日)
dates = np.random.choice(dates, size=int(len(dates)*0.7), replace=False)
customers = ['C001', 'C002', 'C003', None]  # None模拟匿名交易
categories = ['Groceries', 'Dining', 'Travel', 'Retail', 'Dinin']  # 故意加错拼
amounts = np.random.normal(200, 100, len(dates)*3)  # 正态分布,含负值
amounts = np.clip(amounts, -500, 5000)  # 限制范围
# 构建DataFrame
df_raw = pd.DataFrame({
    'date': np.repeat(dates, 3),
    'customer_id': np.random.choice(customers, len(dates)*3),
    'category': np.random.choice(categories, len(dates)*3),
    'amount': amounts.round(2),
    'fee': (np.abs(amounts) * 0.025).round(2)  # 退费时fee为0
})

这个 df_raw ,才是你每天面对的“真实数据”。下面所有分析,都基于此。

6.2 节点1:数据清洗——别让脏数据毁掉整个流水线

# 1. 过滤无效客户(匿名交易不参与客户级分析)
df_clean = df_raw.dropna(subset=['customer_id'])
# 2. 修正商户类别拼写
category_map = {'Dinin': 'Dining'}
df_clean['category'] = df_clean['category'].replace(category_map)
# 3. 处理退款(金额为负,但fee为正,需标记)
df_clean['is_refund'] = df_clean['amount'] < 0
# 4. 统一时区(假设原始为UTC,转为Asia/Shanghai)
df_clean['date'] = pd.to_datetime(df_clean['date']).dt.tz_localize('UTC').dt.tz_convert('Asia/Shanghai').dt.date

清洗不是可选项,是流水线第一道闸门。漏掉 dropna ,后续 groupby 会把 None 当一个客户;漏掉拼写修正, Dinin Dining 被当两个类别,分析全错。

6.3 节点2:多维聚合——一份代码,七种洞察

# 按客户+类别,计算核心指标
agg_result = df_clean.groupby(['customer_id', 'category']).agg(
    total_spend=('amount', 'sum'),
    avg_transaction=('amount', 'mean'),
    transaction_count=('amount', 'count'),
    refund_count=('is_refund', 'sum'),
    fee_total=('fee', 'sum'),
    amount_std=('amount', 'std'),
    amount_range=('amount', transaction_range)  # 用前面定义的函数
).round(2)

# 拍平列名
agg_result.columns = ['_'.join(col).strip() for col in agg_result.columns]
# 重置索引,为后续步骤铺路
agg_result = agg_result.reset_index()

这里 transaction_range 函数已集成,且 round(2) 统一精度。注意 'is_refund' 列是布尔型, 'sum' 直接计数,比 'count' 更准( count 会统计所有非空值, sum 只加 True )。

6.4 节点3:滚动分析——捕捉趋势,而非快照

# 按客户排序,确保时间有序
df_sorted = df_clean.sort_values(['customer_id', 'date']).set_index('date')
# 计算每个客户的7日滚动均值(日历日)
df_sorted['rolling_7d'] = df_sorted.groupby('customer_id')['amount'].rolling('7D', min_periods=3).mean().reset_index(level=0, drop=True)
# 为避免周末空缺,用resample填充
df_daily = df_sorted.groupby(['customer_id', pd.Grouper(freq='D')])['amount'].sum().unstack('customer_id', fill_value=0)
df_daily_filled = df_daily.fillna(method='ffill')
df_daily_filled['C001_7d_avg'] = df_daily_filled['C001'].rolling('7D').mean()

看到没?这里用了两种滚动策略:一种是原始交易粒度(保留所有交易),一种是日汇总粒度(更平滑)。根据业务需求选择。

6.5 节点4:交叉分析——让维度自己说话

# 客户 vs 类别 交叉表(平均交易额)
crosstab_avg = df_clean.pivot_table(
    values='amount',
    index='customer_id',
    columns='category',
    aggfunc='mean',
    fill_value=0
)

# 客户 vs 类别 交叉表(交易笔数)
crosstab_count = df_clean.pivot_table(
    values='amount',
    index='customer_id',
    columns='category',
    aggfunc='count',
    fill_value=0
)

# 合并为一个DataFrame,便于导出
crosstab_combined = pd.concat([
    crosstab_avg.add_suffix('_avg'),
    crosstab_count.add_suffix('_count')
], axis=1)

add_suffix 避免列名冲突, concat(axis=1) 水平拼接。最终 crosstab_combined 有6列: Dining_avg , Dining_count , Groceries_avg ... 这才是老板要的“客户画像”。

6.6 节点5:高管简报——从数据到决策的最后一步

# 按客户聚合,生成Executive Summary
summary = df_clean.groupby('customer_id').agg(
    total_spend=('amount', 'sum'),
    avg_transaction=('amount', 'mean'),
    transaction_count=('amount', 'count'),
    refund_rate=('is_refund', lambda x: (x.sum() / x.count() * 100).round(1)),
    fee_total=('fee', 'sum'),
    high_value_ratio=('amount', lambda x: ((x > 300).sum() / x.count() * 100).round(1))
).round(2)

# 计算衍生指标
summary['spend_per_transaction'] = (summary['total_spend'] / summary['transaction_count']).round(2)
summary['fee_percent'] = (summary['fee_total'] / summary['total_spend'] * 100).round(2)

# 排序:按总消费降序
summary = summary.sort_values('total_spend', ascending=False)

注意 refund_rate high_value_ratio lambda ,因为涉及除法,需防零除。但这是最终输出层,不影响核心逻辑,可接受。

6.7 节点6:风险透视——用自定义函数挖出隐藏模式

def risk_profile(series):
    """深度风险画像:波动性、集中度、异常值"""
    if len(series) < 5:
        return pd.Series({'risk_score': 0, 'concentration': 0})
    
    # 波动性:变异系数(标准差/均值)
    cv = series.std() / series.mean() if series.mean() != 0 else 0
    # 集中度:Top3交易占总额比
    top3_sum = series.nlargest(3).sum()
    concentration = (top3_sum / series.sum() * 100) if series.sum() != 
Logo

码道开发者社区,聚焦华为云码道 CodeArts 代码智能体,沉淀 Agent、Skill、鸿蒙开发实战内容,供开发者查阅资料、交流技术、分享工程实践

更多推荐