在Pandas/Dask中高效实现基于逗号分隔字符串求和的行过滤
问题描述
给定带表头的数据集:
Name,Signal,Date MyName,"1,2,3,4,5,6,7,8,9,10",19-04-2024 MyName,"1,2,3,4,5,6,7,8,9,10",19-04-2024
需求是基于Signal列中数组的求和结果过滤行,尝试了以下代码:
df = read_csv("my_csv.csv", dtype={"Signal" : "string"}, parse_dates=True) for i in df["Signal"]: t = np.array([int(x) for x in i.split(",")]) if t.sum() == 100: #etc
遇到三个问题:
- 如何记录当前行的索引以过滤/删除该行?
- 能否加速该操作或实现更高效的处理?曾考虑分配二维numpy数组一次性解析数值,但不确定是否有效。
- 使用无全局行索引的Dask时,是否有无需将所有数据转为numpy数组的高效过滤方法?
解决方案
1. 记录行索引实现过滤
不要用普通for循环遍历,优先用矢量化操作生成布尔掩码,或者用iterrows()获取索引:
方法一:布尔掩码(推荐)
直接对Signal列应用自定义逻辑生成过滤条件,无需手动记录索引:
import pandas as pd df = pd.read_csv("my_csv.csv", dtype={"Signal": "string"}, parse_dates=["Date"]) # 生成过滤掩码:筛选Signal求和等于100的行 filter_mask = df["Signal"].apply(lambda s: sum(int(x) for x in s.split(",")) == 100) # 应用过滤 filtered_df = df[filter_mask] # 若要直接修改原数据框,直接赋值即可 # df = df[filter_mask]
如果需要保留求和结果用于后续分析,可以先生成求和列:
def calc_signal_sum(signal_str): return sum(int(x) for x in signal_str.split(",")) df["Signal_Sum"] = df["Signal"].apply(calc_signal_sum) filtered_df = df[df["Signal_Sum"] == 100]
方法二:用iterrows()记录索引
如果必须用循环,用iterrows()同时获取索引和行数据:
indices_to_keep = [] for idx, row in df.iterrows(): num_list = [int(x) for x in row["Signal"].split(",")] if sum(num_list) == 100: indices_to_keep.append(idx) filtered_df = df.loc[indices_to_keep]
2. 加速操作的高效方案
批量解析为二维数组确实能大幅提升效率,以下是两种实用方案:
方法一:numpy批量解析(适合Signal元素数量固定的场景)
import numpy as np # 将所有Signal字符串转为二维整数数组 signal_arrays = np.array([list(map(int, s.split(","))) for s in df["Signal"]]) # 计算每行的和 signal_sums = signal_arrays.sum(axis=1) # 过滤行 filtered_df = df[signal_sums == 100]
方法二:pandas字符串拆分+分组求和(适合Signal元素数量不固定的场景)
如果不同行的Signal元素个数不一致,批量转数组会报错,用这种方法兼容性更好:
# 拆分Signal为多列,展开为单行单个数值,再按原索引分组求和 signal_sums = df["Signal"].str.split(",", expand=True).stack().astype(int).groupby(level=0).sum() filtered_df = df[signal_sums == 100]
性能对比
- 循环遍历最慢,仅适合极小数据集;
- numpy批量解析速度最快,适合元素数量固定的数据集;
- 字符串拆分+分组求和兼容性最优,性能介于前两者之间。
3. Dask中的高效过滤方法
Dask无需全局行索引,直接用map_partitions对每个分区独立处理,内存占用低,适合超大规模数据集:
import dask.dataframe as dd def filter_partition(partition): # 对单个分区内的行计算Signal求和并过滤 mask = partition["Signal"].apply(lambda s: sum(int(x) for x in s.split(",")) == 100) return partition[mask] # 读取Dask数据框 ddf = dd.read_csv("my_csv.csv", dtype={"Signal": "string"}, parse_dates=["Date"]) # 对每个分区应用过滤逻辑 filtered_ddf = ddf.map_partitions(filter_partition) # 按需触发计算获取结果 result_df = filtered_ddf.compute()
如果需要保留求和列,也可以在分区内生成临时列后过滤:
def filter_partition(partition): partition["Signal_Sum"] = partition["Signal"].apply(lambda s: sum(int(x) for x in s.split(","))) return partition[partition["Signal_Sum"] == 100]
内容的提问来源于stack exchange,提问作者AnthonyML
相关产品推荐
相关产品推荐

