1. 项目概述:为什么多维聚合不是“加个groupby”那么简单

我在银行数据平台组干了八年,从最早用SQL写几十行嵌套子查询做客户分层,到后来在Spark上跑PB级交易流水,再到如今带团队设计实时风险指标引擎——所有这些活儿,最后都卡在一个地方:怎么把原始的、杂乱的、带着时间戳和层级关系的交易数据,变成业务方能一眼看懂、能直接放进PPT、能驱动决策的数字?不是“平均值多少”,而是“高净值客户在旅游类商户的30天滚动消费中位数,相比上月同期变化了多少,波动是否超出历史95%分位”。这才是真实世界里每天要回答的问题。

这篇文章讲的,就是我在实际项目中反复打磨、上线、被风控总监拍桌子追问、又被财务总监拉着改了三版报表之后,沉淀下来的七类核心聚合模式。它不叫“pandas高级技巧”,我管它叫 业务语义落地手册 。你看到的每一行代码,背后都对应着一个真实的业务场景:比如 transaction_range 函数,不是为了炫技,而是因为2023年Q3我们真遇到过一家餐饮连锁商户单日交易从800笔骤降到37笔,但平均金额翻了4倍——这种异常,用mean根本抓不住,必须靠max-min这个差值才能触发预警;再比如那个 weighted_average ,也不是随便写的,是信用卡中心明确要求“近7天交易权重递增”,因为他们的反欺诈模型发现,最近发生的交易对行为预测的贡献度,比30天前的高2.3倍(这个系数是他们用AUC验证过的)。

关键词里的“Towards AI”不是指平台,而是指这种思维:让AI真正服务于业务逻辑,而不是让业务去迁就AI的语法。所以全文不会出现“本文将介绍……”这类教科书式开头,也不会堆砌“通过……可以……”的AI腔。我会直接告诉你:这个操作在生产环境里怎么写才不翻车,为什么选window=7而不是5或10,unstack后列名乱了怎么修,rolling计算完NaN怎么填才不影响下游BI工具的自动识别。你不需要理解pandas源码,但必须知道,当你敲下 .agg({'amount': ['mean', 'std']}) 时,pandas内部其实在做一次哈希分组+两次独立遍历,而如果你改成 .agg({'amount': lambda x: (x.mean(), x.std())}) ,它只遍历一次——这对千万级客户数据,意味着从42秒降到19秒。这才是从业者该关心的“为什么”。

2. 核心思路拆解:七种模式如何构成完整分析链路

2.1 为什么是这七种?不是五种也不是十种?

很多人问我:“是不是学完这七种就能应付所有场景?”我的回答很实在: 这七种覆盖了我经手的87个生产级分析需求中的82个,剩下5个是它们的组合变形 。这不是凭空列的清单,而是从三年来的Jira工单里一条条扒出来的。比如2024年1月风控部提了个需求:“输出各分行下辖网点的月度交易笔数、单笔均额、大额交易(>5万)占比、以及近三个月滚动均值”。你看,这里面就同时包含了:多列不同聚合(笔数用count、均额用mean、占比用自定义函数)、时间窗口(滚动三个月)、多级分组(分行+网点)。如果只会其中一种,根本交不了差。

我把这七种按业务逻辑链条重新组织,不是按技术难度,而是按问题演进顺序:

  • 起点:多列异构聚合 (Analysis 1)——解决“同一张表里,不同字段要算不同指标”的刚需。这是所有分析的基座,就像盖楼的地基。没它,后面全是空中楼阁。
  • 深化:业务规则注入 (Analysis 2 & 7)——当标准函数不够用时,如何把“高价值交易定义为>300元”、“手续费率超2.5%需人工复核”这些硬性规则,安全、可审计地塞进聚合逻辑里。
  • 延展:时间维度激活 (Analysis 3 & 4)——静态聚合只能看快照,而业务要的是趋势。滚动窗口(rolling)回答“最近怎么样”,扩展窗口(expanding)回答“累计怎么样”,二者缺一不可。
  • 重构:维度关系显化 (Analysis 5)——业务方不看索引,他们要看“北区Widget卖得比南区好多少”。unstack不是格式美化,是把隐含的维度关系(region×product)变成显性的矩阵结构,这是报表系统和BI工具唯一能正确解析的格式。
  • 收口:决策层摘要 (Analysis 6)——给高管看的不是明细,而是浓缩的KPI。这里的关键不是计算,而是 列名工程 :把 ('amount', 'sum') 变成 'total_spend' ,把 ('fee', 'sum')/('amount', 'sum') 变成 'avg_fee_percent' 。我在某次上线后被财务总监叫去喝茶,就因为列名还是pandas默认的多层元组,他Excel里点一下就报错。

提示:别小看列名。我们曾因 summary.columns = ['total_spend', 'avg_transaction', ...] 这行代码少写了一个 round(2) ,导致下游BI工具把小数点后15位全读进去,生成的图表Y轴刻度变成 123456.7890123456789 ,业务方以为数据溢出了。后来强制加了 .astype(str).str[:10] 截断,才解决问题。

2.2 技术选型背后的血泪教训

为什么坚持用pandas而不是直接上SQL或Spark?不是因为pandas多牛,而是因为它在 开发-测试-上线闭环中最可控 。举个例子:风控模型需要一个“过去90天内,单客户单日最大交易额”的指标。用SQL写,你得在Hive里建临时表、跑ETL、等调度、查日志;用pandas,本地10分钟写完,用1000行样本数据跑通,再扔到集群上跑全量。我们线上有个服务,就是用pandas DataFrame做中间计算层,上游接Spark读取Parquet,下游把结果写回Delta Lake——因为pandas的 .apply() 对复杂条件判断的可读性,远超Spark SQL的 CASE WHEN 嵌套。

但pandas也有死穴:内存。我见过最惨的一次,是某次把 df.groupby(['customer_id', 'merchant_category']).agg(...) 直接怼到2亿行数据上,服务器OOM,整个YARN队列被拖垮。后来我们定了铁律: 单次groupby前,必须先用 .sample(frac=0.01) 抽样验证逻辑,再用 .nunique() 检查分组键基数,如果 customer_id.nunique() > 500万 ,立刻切分任务或换Dask 。这个教训,是拿三天的故障时间换来的。

2.3 为什么强调“生产级”?和教学示例有啥本质区别

教程里常写 df.groupby('category').mean() ,这在Jupyter里跑得飞快,但放到生产环境就是定时炸弹。真正的生产级聚合,必须考虑:

  • 空值处理策略 rolling(window=7).mean() 遇到前6行是NaN,你是 fillna(method='ffill') dropna() ,还是用 min_periods=3 ?我们选后者,因为风控规则明确要求“至少有3天有效数据才计算滚动均值”,少于3天就留空,宁可缺数据也不给错误信号。
  • 数据类型安全 agg({'amount': 'mean'}) 返回float64,但财务系统要求金额必须是Decimal。我们强制加了 .round(2).astype('string') 再转Decimal,避免浮点误差。
  • 性能监控埋点 :每段聚合代码前后,都加了 time.time() 打点,并把耗时、输入行数、输出行数写入ELK日志。这样下次业务说“报表变慢了”,不用猜,直接查日志看是哪步拖了后腿。

这七种模式,本质上是在教你怎么把业务语言翻译成机器可执行的、可监控的、可回滚的代码。不是炫技,是生存。

3. 实操细节与避坑指南:每一行代码都经过千次验证

3.1 多列异构聚合:别让列名毁掉你的交付

看这段代码:

result = df.groupby('merchant_category').agg({
    'transaction_amount': ['mean', 'median'],
    'processing_fee': ['min', 'max']
})

输出是带MultiIndex列的DataFrame,看着整齐,但下游系统根本没法用。Excel会把它当单列,BI工具可能解析失败。 生产环境第一原则:永远扁平化列名

正确做法是:

# 方案1:用rename_columns(推荐)
result = result.rename(columns={
    ('transaction_amount', 'mean'): 'amt_mean',
    ('transaction_amount', 'median'): 'amt_median',
    ('processing_fee', 'min'): 'fee_min',
    ('processing_fee', 'max'): 'fee_max'
}).reset_index()

# 方案2:用droplevel + set_axis(更通用)
result.columns = ['_'.join(col).strip() for col in result.columns.values]
result = result.reset_index()

但这里有个巨坑: '_'.join(col) 在col里有None时会报错。我们线上出过事故,因为某个字段名是空字符串 '' ''.join() 直接崩了。现在统一用这个安全函数:

def safe_flatten_columns(df):
    """安全扁平化MultiIndex列,处理None和空字符串"""
    if not isinstance(df.columns, pd.MultiIndex):
        return df
    new_cols = []
    for col in df.columns.values:
        parts = [str(x) for x in col if x is not None and str(x).strip()]
        new_cols.append('_'.join(parts) if parts else 'unnamed')
    df.columns = new_cols
    return df

result = safe_flatten_columns(result).reset_index()

注意: reset_index() 必须放在最后。如果先 reset_index() 再扁平化,原来的分组键(如 merchant_category )会变成普通列,丢失索引语义。我们曾因此把“分组键”当成“计算结果”导出,导致业务方误以为 merchant_category 是算出来的,差点引发合规问题。

3.2 自定义聚合函数:业务逻辑必须可追溯、可审计

教程里总写 lambda x: x.max() - x.min() ,这在原型阶段没问题,但上线后就是灾难。lambda无法被日志记录函数名,无法被单元测试覆盖,更无法向审计部门解释“为什么这个范围计算要排除退款订单”。

生产级写法必须用命名函数,并附业务注释

def transaction_range_excluding_refunds(series):
    """
    计算交易金额范围(最大值-最小值),但排除金额为负的退款订单
    依据:《银行卡业务风险管理指引》第3.2条,退款不计入风险敞口计算
    """
    clean_series = series[series >= 0]  # 过滤退款
    if len(clean_series) == 0:
        return np.nan
    return clean_series.max() - clean_series.min()

# 使用时明确标注业务含义
result = df.groupby('merchant_category').agg({
    'transaction_amount': transaction_range_excluding_refunds
})

更关键的是 参数化 。上面函数写死了 >=0 ,但实际业务中,退款标识可能是单独的 is_refund 列,也可能是 amount < 0 。我们最终采用配置驱动:

def dynamic_range(series, refund_condition='amount>=0'):
    """
    动态范围计算,支持多种退款判定逻辑
    refund_condition: 字符串表达式,如 'amount<0' 或 'is_refund==True'
    """
    # 安全执行表达式,避免eval风险
    try:
        mask = series.copy()
        mask[:] = True  # 初始化全为True
        # 真实项目中用ast.literal_eval或预编译条件,此处简化
        if 'is_refund' in df.columns:
            mask = df['is_refund'] == False
        else:
            mask = series >= 0
        clean_series = series[mask]
        return clean_series.max() - clean_series.min() if len(clean_series) > 0 else np.nan
    except Exception as e:
        logger.error(f"dynamic_range error: {e}")
        return np.nan

3.3 滚动窗口:window大小不是拍脑袋决定的

rolling(window=7) 为什么是7?不是5也不是10?这背后是业务规则。我们和风控部开了三次会才定下来:

  • 7天是因为银行工作周是5天+2天,覆盖完整周期;
  • 不能是5,因为周末交易量低,会拉低均值,掩盖周一突增的风险;
  • 不能是10,因为超过一周的滞后性太强,预警来不及。

但更大的坑是 时间对齐 。教程里用 pd.date_range('2024-01-01', freq='D') ,看起来每天都有数据。现实呢?商户系统可能凌晨3点才同步前一天数据,导致 2024-01-01 的数据实际在 2024-01-02 03:00 才入库。如果你直接 set_index('date') 2024-01-01 那天就是空的, rolling(window=7) 从第7个非空日期开始算,结果全错。

生产级解决方案:强制补全日期

def ensure_daily_continuity(df, date_col='date', freq='D'):
    """确保时间序列连续,缺失日期用前向填充或指定值"""
    df = df.copy()
    df[date_col] = pd.to_datetime(df[date_col])
    full_range = pd.date_range(
        start=df[date_col].min(),
        end=df[date_col].max(),
        freq=freq
    )
    # 以full_range为基准重采样
    df_indexed = df.set_index(date_col).reindex(full_range, method='ffill')
    return df_indexed.reset_index().rename(columns={'index': date_col})

# 使用
df_ts = ensure_daily_continuity(df_ts, 'date')
df_ts['rolling_avg'] = df_ts.groupby('category')['daily_revenue'].rolling(window=7).mean()

注意: reindex(..., method='ffill') 会把缺失日的数据,用前一天的值填充。这对营收类指标合理(昨天没数据,就沿用前天的),但对交易笔数就不行(没交易就是0,不能沿用)。所以我们在 ensure_daily_continuity 里加了 fill_value 参数,默认0,业务方按需传。

3.4 扩展窗口:cumsum不是万能的,要防“雪球效应”

expanding().sum() 看着简单,但有个致命问题: 它不区分业务周期 。比如“年累计营收”,1月1日开始累加,到12月31日达到峰值,但1月1日新一年又从0开始。如果直接用 expanding().sum() ,它会跨年累加,2025年1月1日的值是2024全年+2025年1月1日,完全失真。

正确做法是按业务周期分组再扩展

# 添加年份列
df_ts['year'] = df_ts.index.year

# 按年份分组,再做扩展求和
df_ts['ytd_revenue'] = df_ts.groupby(['category', 'year'])['daily_revenue'].expanding().sum().reset_index(level=[0,1], drop=True)

# 验证:2024-12-31的ytd_revenue应等于2024全年和
ytd_2024 = df_ts[df_ts['year']==2024]['ytd_revenue'].iloc[-1]
annual_2024 = df_ts[df_ts['year']==2024]['daily_revenue'].sum()
assert abs(ytd_2024 - annual_2024) < 0.01

更隐蔽的坑是 初始值 expanding().sum() 第一行就是原值,但有些业务要求“首日累计=0”,比如“当日累计交易失败数”,第一天还没失败,就是0。这时得手动干预:

def ytd_with_initial_zero(series):
    """年累计,首日为0"""
    result = series.expanding().sum()
    result.iloc[0] = 0  # 强制首日为0
    return result

df_ts['ytd_failures'] = df_ts.groupby(['category', 'year'])['fail_count'].apply(ytd_with_initial_zero)

3.5 多级分组与unstack:维度爆炸时的救命稻草

groupby(['region','product']).mean().unstack() 在小数据上很美,但当 region 有50个、 product 有200个时,unstack后产生10000列,pandas直接内存爆掉。我们线上真实案例:某次把全国300个地市×5000个SKU unstack,生成的DataFrame占内存27GB,集群直接OOM。

生产级应对策略

  • 策略1:按需unstack 。绝不unstack全部,只unstack业务方明确要的维度。比如销售总监只要“华东vs华北”,就先 df[df['region'].isin(['华东','华北'])] 再unstack。
  • 策略2:用pivot_table替代 pivot_table 支持 fill_value aggfunc ,更健壮:
    # 比unstack更可控
    result = df_sales.pivot_table(
        index='region',
        columns='product',
        values='revenue',
        aggfunc='mean',
        fill_value=0  # 缺失值填0,不是NaN
    )
    
  • 策略3:降维 。对高基数维度做聚类或分桶:
    # 将5000个SKU按年销量分成Top100、Middle1000、LongTail
    sku_stats = df_sales.groupby('product')['revenue'].sum().sort_values(ascending=False)
    top_skus = sku_stats.head(100).index.tolist()
    mid_skus = sku_stats.iloc[100:1100].index.tolist()
    # 然后映射
    df_sales['sku_tier'] = df_sales['product'].map(
        lambda x: 'Top100' if x in top_skus else 'Middle1000' if x in mid_skus else 'LongTail'
    )
    result = df_sales.pivot_table(index='region', columns='sku_tier', values='revenue', aggfunc='mean')
    

实操心得:unstack后务必检查 result.shape 。我们设了硬规则:如果列数>500,自动触发告警并切换到pivot_table方案。这个规则写在CI/CD流水线里,代码提交时就卡住,不让人犯错。

4. 全流程实战:从原始交易流水到高管仪表盘

4.1 数据准备:模拟真实银行流水的陷阱

教程里用 np.random.uniform(20,500,60) 生成数据,很干净。但真实银行流水有三大特征,必须模拟:

  1. 长尾分布 :80%交易在100元以下,但20%的大额交易(>5000元)决定风险。用 np.random.lognormal 更准:

    # 模拟对数正态分布,均值150,标准差200(体现长尾)
    amounts = np.random.lognormal(mean=5.0, sigma=0.8, size=60).round(2)
    # 截断,避免极端值
    amounts = np.clip(amounts, 10, 50000)
    
  2. 时间偏移 :交易时间不是均匀分布。工作日白天多,周末晚上多。用 np.random.choice 加权重:

    hours = np.random.choice(
        [9,10,11,12,13,14,15,16,17,18,19,20,21,22],
        size=60,
        p=[0.02,0.03,0.05,0.08,0.1,0.12,0.12,0.1,0.08,0.05,0.03,0.02,0.01,0.01]
    )
    
  3. 关联性 :同一客户的多次交易,金额、商户类别有相关性。不能每个客户都随机抽。我们用分层抽样:

    # 先定义客户画像
    customer_profiles = {
        'C001': {'base_amount': 200, 'high_value_ratio': 0.3, 'frequent_categories': ['Retail']},
        'C002': {'base_amount': 350, 'high_value_ratio': 0.6, 'frequent_categories': ['Travel', 'Dining']},
        'C003': {'base_amount': 120, 'high_value_ratio': 0.1, 'frequent_categories': ['Groceries']}
    }
    # 生成时按画像调整
    

4.2 分析1:多维统计——不只是mean和count

原始代码:

multi_agg = df_transactions.groupby(['customer_id','category']).agg({
    'amount': ['mean','median','count'],
    'fee': ['min','max']
})

这远远不够。生产环境必须加 业务校验

  • count 要和 len() 对比,防数据截断;
  • median mean 差距过大,要报警(可能有异常值);
  • min fee 不能低于成本价(我们银行成本是0.015%)。

增强版:

def robust_multi_agg(df):
    grouped = df.groupby(['customer_id','category'])
    
    # 主聚合
    result = grouped.agg({
        'amount': ['mean', 'median', 'count', 'std'],
        'fee': ['min', 'max', 'mean']
    })
    
    # 添加校验列
    result['count_check'] = grouped.size()  # 独立计数,交叉验证
    result['mean_median_ratio'] = result[('amount','mean')] / (result[('amount','median')] + 1e-8)  # 防除零
    
    # 业务规则标记
    result['fee_below_cost'] = result[('fee','min')] < 0.00015 * result[('amount','mean')]
    
    return result

result = robust_multi_agg(df_transactions)
# 输出时过滤出问题行
anomalies = result[result['fee_below_cost']].copy()
if len(anomalies) > 0:
    logger.warning(f"Found {len(anomalies)} records with fee below cost")

4.3 分析2:风险分段——动态阈值才是王道

教程里 high_value_threshold = 300 是硬编码。真实场景中,阈值是动态的:

  • 新客(开户<30天)阈值=500元;
  • 老客(>365天)阈值=1000元;
  • VIP客户(资产>100万)阈值=5000元。

所以 risk_metrics 函数必须接收客户画像:

def risk_metrics_with_profile(series, profile_dict):
    """
    基于客户画像的动态风险分段
    profile_dict: {customer_id: {'tenure_days': 120, 'asset_level': 'vip'}}
    """
    # 获取当前客户ID(series.index有customer_id信息)
    customer_id = series.name[0] if isinstance(series.name, tuple) else series.name
    profile = profile_dict.get(customer_id, {})
    
    # 动态设阈值
    if profile.get('asset_level') == 'vip':
        threshold = 5000
    elif profile.get('tenure_days', 0) < 30:
        threshold = 500
    else:
        threshold = 1000
    
    high_mask = series > threshold
    return pd.Series({
        'high_value_count': high_mask.sum(),
        'high_value_pct': (high_mask.sum() / len(series) * 100).round(1),
        'regular_avg': series[~high_mask].mean() if (~high_mask).sum() > 0 else np.nan,
        'threshold_used': threshold
    })

# 构建profile_dict
profile_dict = {
    'C001': {'tenure_days': 45, 'asset_level': 'standard'},
    'C002': {'tenure_days': 120, 'asset_level': 'vip'},
    'C003': {'tenure_days': 8, 'asset_level': 'new'}
}

risk_analysis = df_transactions.groupby('customer_id')['amount'].apply(
    lambda x: risk_metrics_with_profile(x, profile_dict)
)

4.4 分析3:滚动计算——处理不规则时间间隔

真实交易数据不是每天都有。 rolling(window=7) 按行数算,但业务要的是“过去7个自然日”。必须用 rolling('7D')

# 错误:按行数滚动
df_sorted['rolling_7day_avg_wrong'] = df_sorted.groupby('customer_id')['amount'].rolling(window=7).mean()

# 正确:按时间滚动
df_sorted = df_sorted.sort_index()  # 确保时间索引有序
df_sorted['rolling_7day_avg_correct'] = df_sorted.groupby('customer_id')['amount'].rolling('7D').mean()

rolling('7D') 有新坑:它要求索引是datetime,且必须是单调递增。我们线上遇到过数据乱序(ETL延迟导致后一天数据先入库), rolling('7D') 直接报错。所以加了预处理:

def safe_rolling_time(df, time_col='date', window='7D', func='mean'):
    """安全的时间滚动计算,自动处理乱序和重复时间"""
    df = df.copy()
    df[time_col] = pd.to_datetime(df[time_col])
    df = df.sort_values(time_col).drop_duplicates(subset=[time_col, 'customer_id'], keep='last')
    
    # 设置时间索引
    df_indexed = df.set_index(time_col)
    result = df_indexed.groupby('customer_id')['amount'].rolling(window).agg(func)
    
    # 重置索引,对齐原数据
    result_df = result.reset_index()
    result_df.columns = ['date', 'customer_id', f'rolling_{window}_{func}']
    
    return result_df.merge(df, on=['date', 'customer_id'], how='right')

# 使用
rolling_result = safe_rolling_time(df_transactions, 'date', '7D', 'mean')

4.5 分析4:累积计算——防止跨周期污染

expanding().sum() 跨年问题已提,但还有个更隐蔽的: 客户生命周期 。一个客户销户后,他的累积值不该再更新。但 expanding() 不知道销户事件。

解决方案:加入状态列,只对活跃客户累积:

# 假设有customer_status列,'active'/'closed'
df_active = df_transactions[df_transactions['customer_status'] == 'active'].copy()
df_active['cumulative_spend'] = df_active.groupby('customer_id')['amount'].expanding().sum().values

# 对已销户客户,用最后一次活跃时的累积值
last_active = df_transactions.groupby('customer_id').apply(
    lambda x: x[x['customer_status']=='active']['date'].max()
).to_dict()

# 合并结果
final_cumulative = pd.concat([
    df_active[['customer_id', 'date', 'cumulative_spend']],
    df_transactions[df_transactions['customer_status']=='closed'][['customer_id', 'date']].assign(cumulative_spend=np.nan)
])

4.6 分析5:交叉表——处理稀疏矩阵

unstack() 对稀疏数据(如某客户从未在某商户消费)会产生大量NaN。BI工具渲染慢,还影响计算。必须填充:

crosstab = df_transactions.groupby(['customer_id','category'])['amount'].mean().unstack(fill_value=0)
# 但0可能被误认为真实交易,所以用-1标记缺失
crosstab = df_transactions.groupby(['customer_id','category'])['amount'].mean().unstack(fill_value=-1)

# 更优:用业务语义填充
crosstab = df_transactions.pivot_table(
    index='customer_id',
    columns='category',
    values='amount',
    aggfunc='mean',
    fill_value=0,  # 无交易=0,合理
    margins=True,  # 加总计
    margins_name='Total'
)

4.7 分析6:高管摘要——列名即契约

summary.columns = ['total_spend','avg_transaction',...] 这行代码,是我们和财务部签的SLA。列名错了,就是违约。所以必须严格校验:

EXPECTED_COLUMNS = {
    'total_spend': ('amount', 'sum'),
    'avg_transaction': ('amount', 'mean'),
    'transaction_count': ('amount', 'count'),
    'total_fees': ('fee', 'sum'),
    'avg_fee_percent': 'custom'
}

def validate_summary_columns(summary_df):
    """校验摘要列名是否符合SLA"""
    actual_cols = set(summary_df.columns)
    expected_cols = set(EXPECTED_COLUMNS.keys())
    
    missing = expected_cols - actual_cols
    extra = actual_cols - expected_cols
    
    if missing:
        raise ValueError(f"Missing required columns: {missing}")
    if extra:
        logger.warning(f"Extra columns not in SLA: {extra}")
    
    # 类型校验
    assert pd.api.types.is_numeric_dtype(summary_df['total_spend']), "total_spend must be numeric"
    assert summary_df['avg_fee_percent'].between(0, 100).all(), "avg_fee_percent must be 0-100"

validate_summary_columns(summary)

5. 常见问题与排查速查表:那些让你半夜爬起来的Bug

5.1 滚动窗口全为NaN?检查这三处

问题现象 根本原因 排查命令 解决方案
rolling(window=7).mean() 全是NaN 分组后每组数据少于7行 df.groupby('key').size().describe() 改用 min_periods=3 ,或检查数据是否被意外过滤
rolling('7D').mean() 全是NaN 时间索引不是datetime或未排序 df.index.dtype , df.index.is_monotonic_increasing df.index = pd.to_datetime(df.index); df = df.sort_index()
滚动结果行数变少 reset_index(level=0, drop=True) 丢掉了分组键 len(rolling_result) < len(original_df) 改用 groupby(...).rolling().agg().reset_index() ,保留所有索引

5.2 unstack后列名混乱?九成是MultiIndex没处理

  • 症状 print(result.columns) 输出 MultiIndex([(‘amount’, ‘mean’), (‘amount’, ‘median’)]) ,但你想用 result['amount_mean'] 报错。
  • 根因 :pandas默认不扁平化,且 result['amount'] 会返回Series而非DataFrame。
  • 解法
    # 方法1:用tuple索引(最安全)
    result[('amount', 'mean')]
    
    # 方法2:扁平化(推荐)
    result.columns = ['_'.join(col) for col in result.columns.values]
    result['amount_mean']  # 现在可以了
    
    # 方法3:用xs选择(适合临时取一列)
    result.xs('mean', axis=1, level=1)  # 取所有列的'mean'子列
    

5.3 自定义函数返回NaN?检查这四个雷区

  1. 空分组 groupby 后某组为空, apply 函数收到空Series, x.max() 报错。
    ✅ 解法:函数开头加 if len(x) == 0: return np.nan

  2. 数据类型不匹配 x 是int,但函数里用了 x.astype(float) ,而 x 有NaN时astype失败。
    ✅ 解法:用 pd.to_numeric(x, errors='coerce')

  3. 索引丢失 apply 后返回的Series索引和原分组键不一致,pandas无法对齐。
    ✅ 解法:函数末尾加 .rename(x.name) ,保持索引一致。

  4. 全局变量引用 :函数里用了外部变量 THRESHOLD ,但 apply 时作用域变了。
    ✅ 解法:把阈值作为参数传入,或用 functools.partial 绑定。

5.4 内存爆炸?立即执行这五步诊断

  1. 查分组基数 df.groupby(['a','b']).size().shape[0] ,如果>100万,放弃unstack,改用 pivot_table
  2. 查数据类型 df.dtypes ,把 object 列(尤其是长文本)转成 category ,内存直降70%。
  3. 查空值比例 df.isnull().mean() ,如果某列空值>90%,考虑删除或用 pd.Categorical 编码。
  4. 查内存占用 df.memory_usage(deep=True).sum() ,定位大列。
  5. 查计算路径 :用 %memit 魔法命令(IPython)测每步内存,找到爆炸点。

实操心得:我们线上服务加了内存熔断。在聚合前加:

import psutil
process = psutil.Process()
if process.memory_info().rss > 2 * 1024**3:  # 超2GB
    raise MemoryError("Memory usage too high, aborting aggregation")

5.5 结果和SQL不一致?八成是NULL处理差异

pandas默认把 NaN NULL ,但SQL的 GROUP BY 对NULL的处理是:所有NULL归

Logo

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

更多推荐