You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在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

遇到三个问题:

  1. 如何记录当前行的索引以过滤/删除该行?
  2. 能否加速该操作或实现更高效的处理?曾考虑分配二维numpy数组一次性解析数值,但不确定是否有效。
  3. 使用无全局行索引的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.25 21:23:11