基于Python Multiprocessing Pool实现批量文本读写循环并行化问询
Hey there! Let's break down how to parallelize that loop in your auto method. First off, since each iteration over range(len(a)) works independently (each creates its own TxtInOut{i} directory and doesn't interfere with other iterations), this is ideal for parallel processing. Here's a step-by-step solution:
Step 1: Fix a Critical Serial-to-Parallel Issue
Your original code writes to a shared sample.txt in the main directory for every iteration. When running in parallel, multiple processes will overwrite this file simultaneously, leading to corrupted data. We'll fix this by writing sample.txt directly into the corresponding TxtInOut{i} directory instead.
Step 2: Refactor Loop Logic into a Reusable Function
We'll extract the inner loop code into a dedicated function that handles one iteration at a time. This makes it easy to feed into a parallel executor.
Step 3: Implement Parallel Execution with concurrent.futures.ProcessPoolExecutor
ProcessPoolExecutor is a clean, high-level way to manage parallel processes in Python. It handles process spawning, task distribution, and cleanup automatically.
Modified Code
Here's your updated auto method with parallelization:
def auto(self, directory, number_samples, number_trajectories, grid_jump): import glob, os from concurrent.futures import ProcessPoolExecutor # Run pre-processing steps once (these don't need parallelization) a = self.generate_samples(number_samples, number_trajectories, grid_jump) b = [i.split(',') for i in self.readlist('rng.txt')] # Convert to absolute path to avoid working directory confusion across processes directory = os.path.abspath(directory) def process_sample(i): """Handles processing for a single sample iteration""" try: # Create the dedicated output directory path = os.path.join(directory, f'TxtInOut{i}') os.mkdir(path) # Write sample.txt directly to the output directory (no shared file!) sample_path = os.path.join(path, 'sample.txt') with open(sample_path, 'w') as sample: for j in range(len(b)): sample.write(f'{b[j][0]},{b[j][1]},{b[j][2]},{a[i][j]}\n') # Process each .cio file for file in glob.glob(os.path.join(directory, '*cio')): e = self.file_cio(file) c = self.fig_fig(e[1][6][1]) self.wwq_write(self.basins_bsn(e[1][46][1])[1][63][1], path) self.basins_bsn_write(e[1][46][1], path) for j in c[0]: d = self.subbasin(c[0][j]) self.pnd_write(d[1][30][1], path) for k in d[3]: self.hru_write(d[3][k][0], path) self.mgt_write(d[3][k][1], path) self.gw_write(d[3][k][4], path) self.sep_write(d[3][k][6], path) self.subbasin_write(c[0][j], path) for j in c[1]: self.rte_write(c[1][j][0], path) self.swq_write(c[1][j][1], path) self.file_cio_write(file, path) except Exception as exc: print(f'Error processing sample {i}: {exc}') raise # Re-raise if you want the executor to catch it # Run processes in parallel # Use max_workers = os.cpu_count() to match your core count, or adjust as needed with ProcessPoolExecutor(max_workers=os.cpu_count()) as executor: executor.map(process_sample, range(len(a)))
Key Notes to Avoid Headaches
- Windows Compatibility: If you're on Windows, make sure your main code is wrapped in an
if __name__ == '__main__':block. This prevents infinite process spawning due to how Windows handles multiprocessing. - Pickle Safety: Your
SWATinstance must be pickle-able (serializable) for it to work across processes. If your class has non-serializable attributes (like open file handles or network connections), refactor those to be created insideprocess_sampleinstead of stored onself. - Resource Limits: If
number_samplesis extremely large, you might want to lowermax_workers(e.g.,os.cpu_count() - 1) to avoid overwhelming your system with too many concurrent processes. - Error Handling: The
try-exceptblock inprocess_samplehelps catch and report issues specific to individual iterations, so you don't lose all progress if one sample fails.
内容的提问来源于stack exchange,提问作者Reza Ehsani

