优化Pandas中DataFrame重塑操作:10万行数据提速需求
优化Pandas日志数据重塑与缺失值填充的性能方案
你的现有方案慢的核心原因是循环生成DataFrame列表+多次Merge(时间复杂度随信号数量线性上升),以及缺失值填充的冗余操作。以下是针对10万行数据的高效优化实现:
先明确数据结构(基于你的步骤推断)
假设原始CSV读取后的DataFrame结构:
| Timestamp | Signals | Data_col |
|---|---|---|
| 2024-01-01 00:00:00 | SignalA | col1:1,col2:2 |
| 2024-01-01 00:00:01 | SignalA | col1:3,col2:4 |
| 2024-01-01 00:00:00 | SignalB | colX:5,colY:6 |
目标输出的宽表格式:
| Timestamp | SignalA_col1 | SignalA_col2 | SignalB_colX | SignalB_colY |
|---|---|---|---|---|
| 2024-01-01 00:00:00 | 1 | 2 | 5 | 6 |
| 2024-01-01 00:00:01 | 3 | 4 | NaN | NaN |
1. 高效拆分Data列
替换循环拆分,用矢量化的explode+str.split操作,避免逐行处理:
import pandas as pd # 读取CSV时提前指定类型减少内存占用 raw_df = pd.read_csv( "log_data.csv", dtype={"Timestamp": str, "Signals": "category", "Data_col": str} ) # 转换时间列(必须提前处理,后续排序/填充依赖) raw_df["Timestamp"] = pd.to_datetime(raw_df["Timestamp"]) # 拆分Data_col为key-value对:先按逗号拆成列表,再展开为行,最后拆分键值 raw_df = raw_df.assign(temp=raw_df["Data_col"].str.split(",")).explode("temp") raw_df[["key", "value"]] = raw_df["temp"].str.split(":", expand=True) raw_df = raw_df.drop(columns=["temp", "Data_col"]) # 将value转为数值类型(根据实际数据调整为int/float) raw_df["value"] = pd.to_numeric(raw_df["value"], errors="coerce")
2. 一步重塑为目标宽表(替代循环+Merge)
用pivot_table或unstack完成长表转宽表,这两个都是Pandas内部优化的矢量化操作,比多次Merge快10~100倍:
方法一:pivot_table(更直观)
# 合并Signals和key为最终列名 raw_df["full_colname"] = raw_df["Signals"] + "_" + raw_df["key"] # 透视生成宽表,aggfunc根据需求调整(比如取第一个值/平均值) wide_df = raw_df.pivot_table( index="Timestamp", columns="full_colname", values="value", aggfunc="first" ).reset_index()
方法二:unstack(性能略优)
# 设置多层索引后直接展开 wide_df = raw_df.set_index(["Timestamp", "Signals", "key"])["value"].unstack(["Signals", "key"]) # 重命名列:合并Signals和key wide_df.columns = wide_df.columns.map(lambda x: f"{x[0]}_{x[1]}") wide_df = wide_df.reset_index()
3. 高效填充缺失值
替换冗余的fillna写法,先排序再链式填充,或用时间插值(适合时间序列):
# 必须先按时间排序,否则填充逻辑错误 wide_df = wide_df.sort_values("Timestamp") # 方案1:先向前填充再向后填充(通用场景) wide_df = wide_df.ffill().bfill() # 方案2:时间插值(适合时间间隔规则的日志,性能更优) wide_df = wide_df.set_index("Timestamp").interpolate(method="time").reset_index()
额外性能提升技巧
- 内存优化:将
Signals设为category类型(如果唯一信号数量远小于总行数),数值列用float32/int32代替默认的float64/int64,减少内存占用,直接提升运算速度。 - 并行处理:如果数据量超过内存,用Dask DataFrame替代Pandas,自动并行处理数据:
import dask.dataframe as dd ddf = dd.read_csv("log_data.csv", dtype={"Signals": "category"}) # 执行类似Pandas的操作,最后用compute()得到结果 result = ddf.assign(temp=ddf["Data_col"].str.split(",")).explode("temp").compute()
内容的提问来源于stack exchange,提问作者Hannibal
相关产品推荐
相关产品推荐

