拿到一份几百列、上百万行的宽表,如果继续用 df.apply(axis=1) 逐行做计算,执行时间会迅速膨胀。原因并不复杂:Python 的全局解释器锁让默认的多线程无法真正并行执行 CPU 密集操作,而 apply 又频繁调用 Python 函数,开销相当可观。反过来看,DataFrame 的很多列级操作天然互相独立,例如对每一列做标准化、缺失值填充、分位数统计,这些任务完全可以在不同 CPU 核心上同时进行。这篇文章会围绕多列并行处理,给出可以直接落地的实现方式和调优思路,同时把进程开销、序列化成本和跨平台差异讲清楚。

为什么多列计算容易成为性能瓶颈
Pandas 底层大量使用 C 和 NumPy 实现,很多向量化操作速度非常快。例如 df['col'].mean() 这样的计算,主要耗时并不在 Python 循环上。可一旦我们为了提高灵活性而使用 apply、transform 或者自定义函数逐列处理数据,Python 函数调用就成了瓶颈。每一列都会产生一次独立的函数调度,当列数从几十列增加到几百列时,串行执行的总时长会线性甚至超线性增长。
另一个容易被忽略的点是全局解释器锁,也就是 GIL。多线程在遇到需要大量 CPU 计算的 Python 对象操作时,并不能把多个核心真正利用起来。pandas 的部分底层操作会释放 GIL,但多数自定义函数仍然受限制。多进程则通过独立的解释器绕开这个问题,代价是数据需要在进程间复制或序列化。对多列并行处理来说,这个代价通常可以接受,因为每一列只需要序列化一次,换来的却是多个核心同时工作。
还有一个现实因素是机器的核心数。现在开发机或服务器普遍具备 8 核、16 核甚至更多核心,但默认的串行列处理只占用一个核心。与其把希望寄托在更快的单核频率上,不如把相互独立的列任务拆开,让空闲核心参与计算。这个思路适合任何列与列之间没有依赖关系的场景,例如特征标准化、缺失值插补、多列自定义分箱和批量统计。
使用 ProcessPoolExecutor 实现列级并行
标准库中的 concurrent.futures.ProcessPoolExecutor 是实现列级并行最直接的方式之一。它的接口足够简单,而且能自动管理进程池。下面先构造一个稍大的 DataFrame,然后编写一个串行版本作为对照。任务是对 12 个列分别做标准化,每列有一百万条数据。
import pandas as pd
import numpy as np
import time
np.random.seed(42)
df = pd.DataFrame(
np.random.randn(1_000_000, 12),
columns=[f'col_{i}' for i in range(12)]
)
def process_serial(frame):
result = {}
for col in frame.columns:
series = frame[col].copy()
result[col] = (series - series.mean()) / series.std()
return pd.DataFrame(result)
start = time.perf_counter()
res_serial = process_serial(df)
print(f'serial cost: {time.perf_counter() - start:.2f}s')
串行版本的逻辑很清晰:遍历每一列,取出该列数据,计算均值和标准差,再做标准化。问题在于这个循环只在一个核心上运行,且每列都要重复访问 DataFrame 对象。对于 12 列来说差异可能不明显,但如果列数达到上百列,或者每列都包含更复杂的 Python 计算,串行耗时就会非常突出。
并行版本需要把每一列打包成任务交给进程池。这里有一个细节:不要把整个 DataFrame 作为参数传给每个子进程,否则每个任务都会重复序列化其他列的数据。更合理的做法是只传递当前列的数据,例如 df[col].copy()。虽然复制本身也有成本,但相比把完整 DataFrame 序列化多遍要轻得多。
from concurrent.futures import ProcessPoolExecutor
def standardize_col(args):
col_name, series = args
return col_name, (series - series.mean()) / series.std()
def process_parallel(frame):
tasks = [(col, frame[col].copy()) for col in frame.columns]
with ProcessPoolExecutor(max_workers=6) as executor:
results = executor.map(standardize_col, tasks, chunksize=2)
return pd.DataFrame({name: values for name, values in results})
if __name__ == '__main__':
start = time.perf_counter()
res_parallel = process_parallel(df)
print(f'parallel cost: {time.perf_counter() - start:.2f}s')
print(res_parallel.head())
在这个实现中,standardize_col 接收一个元组,包含列名和该列的 Series。进程池启动后,executor.map 会按顺序返回结果,构建新的 DataFrame 时列顺序仍然保持。使用 chunksize 可以减少任务提交和结果回收的通信次数,尤其适合列数较多但单列计算量较小的场景。需要注意的是,并行版本中的函数必须定义在模块级别,否则子进程无法通过 pickle 机制找到它,lambda 表达式和不完整的闭包都可能导致异常。
除了 ProcessPoolExecutor,multiprocessing.Pool 也能完成同样的事情。多数情况下 ProcessPoolExecutor 的 API 更现代,但如果你需要更细粒度的进程控制,或者要在 Jupyter 中反复调试,可以把函数放到独立的 Python 文件里再导入。Jupyter 的交互式单元中直接启动多进程,经常会出现无法序列化或子进程重复导入的问题,这点在做数据分析时需要额外注意。
线程与进程的取舍以及任务拆分细节
并行处理多列数据并不一定都要用多进程。如果每一列的计算主要是等待 IO,例如批量读文件、请求外部接口、写数据库,那么多线程同样可以显著提升吞吐量,而且线程启动更轻量。判断标准很简单:计算是否长时间占用 CPU。如果函数内部是纯 Python 数值计算、字符串处理、复杂特征逻辑,多进程通常比多线程更稳定;如果函数内部调用了释放 GIL 的 C 库,或者大量时间花在网络等待上,多线程反而更合适。
进程并非越多越好。每启动一个进程,操作系统都需要分配资源,任务数据也要经过 pickle 序列化和反序列化。当列对应的数据量很小、任务数量很少时,进程启动成本可能超过并行带来的收益。一个简单的方法是先跑一次串行版本,记录下总耗时,再用 4 到 8 个进程跑并行版本做比较。如果并行版本反而更慢,说明任务粒度太细,建议把多个列合并成一个任务,或者直接放弃并行。通常单列数据量在十万行以上、列数在几十列以上时,多进程的优势会逐渐明显。
任务拆分还涉及内存使用。把大 DataFrame 拆成大量任务,如果一次提交过多,可能会导致进程间通信缓冲区膨胀。可以设置合理的 max_workers,并用 chunksize 控制每个批次的任务数量。对于特别大的列数据,也可以考虑只传列名加 DataFrame,但要把 DataFrame 提前放到共享内存中,只是 Python 多进程的共享内存对 pandas 对象并不友好,所以直接传该列数据仍然是简单可靠的方案。
最后要强调跨平台行为。Linux 和 macOS 默认使用 fork 方式创建子进程,子进程可以继承父进程的内存快照;Windows 则使用 spawn 方式,会重新导入主模块。因此在 Windows 上必须把启动逻辑放在 if __name__ == '__main__' 内,否则子进程会递归创建进程池导致崩溃。无论哪个平台,编写模块级函数、减少闭包引用、保持参数可序列化,都是避免多进程奇奇怪怪报错的基本原则。把这些细节处理好后,多列并行处理可以稳定融入日常的数据预处理和特征工程流程。
Pandas DataFrame并行处理多列数据修改时间:2026-09-27 10:54:08