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

大型Pandas DataFrame正则匹配的多进程优化方案求助

优化大型Pandas正则匹配任务的方案

核心问题分析(为什么ChatGPT的代码更慢)

你的多进程代码效率低下的原因有三个关键问题:

  • 任务粒度极细:每个进程仅处理一行数据,序列化/反序列化每行的通信开销远超过正则匹配本身的计算耗时,导致CPU大部分时间浪费在进程调度和数据传递上
  • 无效且耗时的锁竞争:用multiprocessing.Value加锁更新进度计数器,每次匹配成功都要触发锁竞争,严重拖慢并行效率;且代码中processed_rows变量未被实际更新,进度跟踪完全无效
  • 未合理配置进程数:默认ProcessPoolExecutor的进程数等于CPU核心数,但细粒度任务会导致频繁调度,无法充分利用CPU资源

第一步:正则预编译+向量化匹配(最快的优化,无需多进程)

Pandas的str.contains是C实现的向量化操作,比Python循环快10-100倍。结合正则预编译+合并,可以直接把36小时的任务压缩到几小时甚至更短。

代码实现

import re
import pandas as pd

# 预编译并合并正则列表(用非捕获分组避免不必要的内存开销)
def compile_and_merge_regexes(regex_list):
    # 合并多个正则为一个,用|分隔,非捕获分组保证匹配逻辑正确
    combined_pattern = '|'.join(f'(?:{regex})' for regex in regex_list)
    return re.compile(combined_pattern)

# 预编译三组正则
in_debit_regex = compile_and_merge_regexes(in_debit_regexes)
in_credit_regex = compile_and_merge_regexes(in_credit_regexes)
in_balance_regex = compile_and_merge_regexes(in_balance_regexes)

# 向量化匹配,na=False处理空值
in_debit_clean_df = ZERO_cashflow_IN_combined_df['message_body'].str.contains(in_debit_regex, na=False)
in_credit_clean_df = ZERO_cashflow_IN_combined_df['message_body'].str.contains(in_credit_regex, na=False)
in_balance_clean_df = ZERO_cashflow_IN_combined_df['message_body'].str.contains(in_balance_regex, na=False)

第二步:如果仍需并行(针对超大规模DataFrame)

如果DataFrame大到内存紧张,或者向量化后仍需进一步提速,采用批处理并行(每个进程处理一个数据批次,而非单行),结合你的ml.c5.9xlarge的36vCPU,最大化利用算力。

代码实现

from concurrent.futures import ProcessPoolExecutor
import numpy as np

def process_batch(batch_series, compiled_regex):
    # 每个批次内部用向量化匹配
    return batch_series.str.contains(compiled_regex, na=False)

def parallel_contains(df_series, compiled_regex, num_workers=36):
    # 拆分DataFrame为多个批次,数量等于进程数
    batches = np.array_split(df_series, num_workers)
    with ProcessPoolExecutor(max_workers=num_workers) as executor:
        # 每个进程处理一个批次
        results = list(executor.map(process_batch, batches, [compiled_regex]*len(batches)))
    # 合并所有批次的结果
    return pd.concat(results)

# 使用并行版本
in_debit_clean_df = parallel_contains(ZERO_cashflow_IN_combined_df['message_body'], in_debit_regex)
in_credit_clean_df = parallel_contains(ZERO_cashflow_IN_combined_df['message_body'], in_credit_regex)
in_balance_clean_df = parallel_contains(ZERO_cashflow_IN_combined_df['message_body'], in_balance_regex)

额外优化建议

  1. 正则本身优化:检查正则表达式是否有冗余,比如用非捕获分组(?:...)替代捕获分组(...),避免不必要的内存占用;移除无用的转义字符
  2. 内存优化:如果DataFrame过大,先过滤空值或无效行:df = df[df['message_body'].notna()],减少处理量
  3. 分布式处理:如果内存无法容纳整个DataFrame,用Dask DataFrame替代Pandas,自动拆分数据并分布式计算,适配SageMaker的算力

内容的提问来源于stack exchange,提问作者elusive17

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.28 13:29:50