生产级多维聚合:从银行风控实战提炼的7种pandas聚合模式
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)
生成数据,很干净。但真实银行流水有三大特征,必须模拟:
-
长尾分布 :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) -
时间偏移 :交易时间不是均匀分布。工作日白天多,周末晚上多。用
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] ) -
关联性 :同一客户的多次交易,金额、商户类别有相关性。不能每个客户都随机抽。我们用分层抽样:
# 先定义客户画像 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?检查这四个雷区
-
空分组 :
groupby后某组为空,apply函数收到空Series,x.max()报错。
✅ 解法:函数开头加if len(x) == 0: return np.nan -
数据类型不匹配 :
x是int,但函数里用了x.astype(float),而x有NaN时astype失败。
✅ 解法:用pd.to_numeric(x, errors='coerce') -
索引丢失 :
apply后返回的Series索引和原分组键不一致,pandas无法对齐。
✅ 解法:函数末尾加.rename(x.name),保持索引一致。 -
全局变量引用 :函数里用了外部变量
THRESHOLD,但apply时作用域变了。
✅ 解法:把阈值作为参数传入,或用functools.partial绑定。
5.4 内存爆炸?立即执行这五步诊断
-
查分组基数
:
df.groupby(['a','b']).size().shape[0],如果>100万,放弃unstack,改用pivot_table。 -
查数据类型
:
df.dtypes,把object列(尤其是长文本)转成category,内存直降70%。 -
查空值比例
:
df.isnull().mean(),如果某列空值>90%,考虑删除或用pd.Categorical编码。 -
查内存占用
:
df.memory_usage(deep=True).sum(),定位大列。 -
查计算路径
:用
%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归
更多推荐



所有评论(0)