如何实现支持自定义函数的Pandas DataFrame.apply并行化函数
Got it, let's build that parallelized apply function step by step. Here's a complete implementation that fits your requirements, along with explanations and a test example:
Parallelizing pandas DataFrame.apply() by partitioning rows
Complete Code
import multiprocessing import numpy as np import pandas as pd from nltk.tokenize import word_tokenize # Your custom row-processing function def apply_me(text): # Importing NLTK components inside the function avoids multiprocessing resource conflicts return word_tokenize(text.lower()) # Helper function to apply the custom logic to a single DataFrame chunk def _apply_to_chunk(chunk, func): # Use axis=1 for row-wise processing; switch to axis=0 if you need column-wise apply return chunk.apply(func, axis=1) # The main parallelization function def parallelize_appl(df, func, num_partitions): # Split the DataFrame into the specified number of row-based partitions df_chunks = np.array_split(df, num_partitions) # Create a process pool (size matches partition count for balanced workload) with multiprocessing.Pool(num_partitions) as pool: # Use starmap to pass both the chunk and custom function to our helper processed_chunks = pool.starmap(_apply_to_chunk, [(chunk, func) for chunk in df_chunks]) # Combine all processed chunks back into a single Series/DataFrame return pd.concat(processed_chunks)
How to Use It
Let's test with sample text data:
# Create test DataFrame sample_data = { 'text': [ "The quick brown fox jumps over the lazy dog", "Parallel processing speeds up pandas operations", "NLTK makes text tokenization straightforward", "Multiprocessing leverages multiple CPU cores" ] } df = pd.DataFrame(sample_data) # Run parallel apply with 2 partitions (adjust based on your CPU core count) result = parallelize_appl(df, apply_me, num_partitions=2) print(result)
Sample Output:
0 [the, quick, brown, fox, jumps, over, the, lazy, dog] 1 [parallel, processing, speeds, up, pandas, operations] 2 [nltk, makes, text, tokenization, straightforward] 3 [multiprocessing, leverages, multiple, cpu, cores] Name: text, dtype: object
Key Details Explained
- Data Partitioning:
np.array_split()splits the original DataFrame into roughly equal chunks, ensuring each process gets a manageable slice of rows to process. - Process Pool:
multiprocessing.Pool()creates a set of worker processes. Matching the pool size to the number of partitions keeps each worker busy with one chunk, minimizing idle time. - Helper Function:
_apply_to_chunkwraps your custom logic and uses pandas' nativeapply()on each chunk—this keeps our parallel logic clean and reuses pandas' optimized apply behavior. - Result Merging:
pd.concat()combines all processed chunks back into a single structure that matches the original DataFrame's row order and format.
Important Notes
- Resource Handling: If your custom function uses external resources (like NLTK models, API clients, or file handles), initialize/import them inside the function (as we did with
word_tokenize). This avoids conflicts from shared resources across processes. - Choosing
num_partitions: A good starting point is setting it to your CPU core count (get this withmultiprocessing.cpu_count()). Too many partitions can lead to unnecessary process-switching overhead. - Pickle Compatibility: Ensure your custom function and any input data are pickleable (most basic Python objects and pandas structures are). Multiprocessing uses pickle to transfer data between processes.
- Axis Flexibility: If you need column-wise processing instead of row-wise, just change
axis=1toaxis=0in_apply_to_chunk.
内容的提问来源于stack exchange,提问作者alvas
相关产品推荐
相关产品推荐

