如何并行化处理含大型复杂对象的Python 3程序?
Hey there! Let's work through how to parallelize your Python 3.6 program effectively—given your 5000+ line codebase, massive dataset of 5000+ files, and reliance on dill for serializing non-pickleable objects, we need a targeted approach that plays to Python's strengths while avoiding common pitfalls.
Your workflow breaks down into three core phases, each with different parallelization sweet spots. Let's tackle them one by one:
1. File Parsing (IO-Bound Workload)
File reading is inherently IO-bound, so thread pools are more efficient than process pools here—they avoid the heavy memory overhead of spawning full processes. Python 3.6 has concurrent.futures.ThreadPoolExecutor built-in, which is perfect for this.
Implementation Snippet:
from concurrent.futures import ThreadPoolExecutor import os def parse_single_file(file_path): # Drop your existing file parsing logic here with open(file_path, 'r') as f: raw_data = f.read() # Add your parsing/validation steps... return parsed_data # Build list of all dataset file paths file_paths = [ os.path.join("your_dataset_dir", fname) for fname in os.listdir("your_dataset_dir") if os.path.isfile(os.path.join("your_dataset_dir", fname)) ] # Use 2-4x your CPU core count for thread pool size (adjust based on IO speed) with ThreadPoolExecutor(max_workers=8) as executor: # Use imap instead of map if you want iterative results (saves memory) parsed_results = list(executor.map(parse_single_file, file_paths))
Pro Tip for Error Handling:
Avoid crashing the whole batch if one file fails. Use submit() + as_completed() to catch exceptions per file:
from concurrent.futures import as_completed with ThreadPoolExecutor(max_workers=8) as executor: futures = {executor.submit(parse_single_file, path): path for path in file_paths} parsed_results = [] for future in as_completed(futures): file_path = futures[future] try: result = future.result() parsed_results.append(result) except Exception as e: print(f"Failed to parse {file_path}: {str(e)}")
2. Internal Representation Generation (CPU-Bound Workload)
If converting parsed data to your custom internal objects involves heavy computation (like complex transformations or calculations), process pools are the way to go—they bypass Python's GIL to utilize multiple CPU cores fully. Since you're using dill for non-pickleable objects, we'll need to configure the process pool to use dill instead of the default pickle.
Implementation Snippet:
import multiprocessing import dill def generate_internal_representation(parsed_data): # Your logic to build non-pickleable internal objects here custom_obj = ... return custom_obj # Configure multiprocessing to use dill for serialization def init_worker(): multiprocessing.reduction.ForkingPickler = dill.Pickler if __name__ == '__main__': # Use all available CPU cores (adjust if you need to reserve resources) with multiprocessing.Pool( processes=multiprocessing.cpu_count(), initializer=init_worker ) as pool: internal_representations = pool.map(generate_internal_representation, parsed_results)
Note: On Windows, Python 3.6 uses the
spawnstart method for processes, so wrap your main logic inif __name__ == '__main__':to avoid infinite process spawning. Linux/macOS useforkby default, which is more forgiving.
3. Statistics & Serialization
Statistics Calculation:
- If you're computing global stats across the entire dataset: Process data in batches with a process pool, compute partial stats per batch, then merge the results.
- If stats are per-object: Merge this step with the internal representation generation phase to avoid redundant data handling.
Dill Serialization:
- Per-object serialization: If you're saving each internal object to a separate file, use a thread pool (IO-bound) to parallelize writes.
- Full dataset serialization: If you're saving one large aggregated object, parallelization won't help much—but you can split the object into chunks, serialize each chunk in parallel, then combine them later.
Serialization Snippet (Per-Object):
def serialize_single_obj(obj, save_path): with open(save_path, 'wb') as f: dill.dump(obj, f) # Generate save paths for each internal object save_paths = [ os.path.join("serialized_dir", f"obj_{idx}.dill") for idx in range(len(internal_representations)) ] with ThreadPoolExecutor(max_workers=4) as executor: executor.map(serialize_single_obj, internal_representations, save_paths)
joblib If you want to avoid manual pool management, joblib (compatible with Python 3.6) simplifies parallel tasks and natively supports dill for serialization. It auto-detects workload type (IO/CPU) and picks the right backend:
from joblib import Parallel, delayed # Parallel file parsing (thread-based for IO) parsed_results = Parallel(n_jobs=8, backend="threading")( delayed(parse_single_file)(path) for path in file_paths ) # Parallel internal object generation (process-based for CPU) internal_representations = Parallel(n_jobs=-1, backend="multiprocessing", verbose=10)( delayed(generate_internal_representation)(data) for data in parsed_results )
n_jobs=-1uses all CPU cores;verboseshows progress to help with debugging.
- Stick to dill version 0.3.x—newer versions may drop Python 3.6 support.
- If your custom objects have edge cases (like closures or dynamically generated classes), register them with dill explicitly using
dill.register()to ensure smooth serialization. - Avoid loading all 5000+ results into memory at once—use iterative methods like
executor.imap()or batch processing to prevent out-of-memory errors.
内容的提问来源于stack exchange,提问作者fortea

