使用ProcessPoolExecutor并行执行代码时部分print语句无输出的原因排查
Let's break down why your print statement stops working after the assignment, and how to fix it:
What's Causing the Block?
The core issue here is how multiprocessing.Manager.dict() handles nested mutable objects like pandas DataFrames. Here's the breakdown:
- Manager dictionaries are proxy objects that only track changes to top-level keys. When you modify a nested object (like adding a column to the DataFrame stored under
dict_w[key]), the proxy doesn't properly sync this change across processes. - This modification can trigger process-level deadlocks or blocking: when your code runs
dict_w[key][c + "_lowess"] = lw, the sub-process gets stuck trying to modify the shared DataFrame, and never reaches the subsequentprint(lw)statement. - When you put
print(lw)before the assignment, the sub-process outputs the result immediately before hitting that blocking operation, which is why you see output in that case.
Fixes to Get Your Code Working
The best approach is to avoid modifying shared mutable objects directly in sub-processes. Instead, let each sub-process work on a local copy of the data, then return the results to the main process for aggregation.
Step 1: Rewrite the Feature Extraction Function
Modify extract_features to work on a local DataFrame copy and return all necessary results instead of modifying shared objects:
def extract_features(key, df_local): lowess = sm.nonparametric.lowess f_v = [] lowess_cols = {} # Work on a local copy to avoid shared state issues df = df_local.copy() # Fix the np.zeros call (original had a shape typo) zero_vals = np.zeros(5) for c in np.setdiff1d(df.columns, ["t", "ts"]): lw = lowess( df[c], df.iloc[:, ~df.columns.duplicated()]["ts"], is_sorted=True, frac=0.05, )[:, 1] lowess_cols[c] = lw f_v.extend(zero_vals) # Return everything needed to update the shared dicts later return key, df, lowess_cols, np.array(f_v)
Step 2: Update the Parallelization Logic
Have the main process handle all updates to shared dictionaries after sub-processes finish their work:
def parallelize(dict_w, max_workers, shared_dic): # Convert shared dict to a regular dict with local DataFrame copies local_dict = {k: v.copy() for k, v in dict_w.items()} with concurrent.futures.ProcessPoolExecutor(max_workers=max_workers) as pool: start = time.perf_counter() # Submit all tasks futures = [pool.submit(extract_features, key, local_dict[key]) for key in local_dict] # Process results as they complete for future in concurrent.futures.as_completed(futures): key, df, lowess_cols, f_v = future.result() # Add lowess columns to the DataFrame for col_name, lw_data in lowess_cols.items(): df[col_name + "_lowess"] = lw_data # Update shared dicts from the main process (safe!) dict_w[key] = df shared_dic[key] = f_v # Now print will work reliably print(lw_data) logger.info("Execution time (parallel) = {}".format(time.perf_counter() - start))
Why This Works
- Each sub-process operates on its own copy of the DataFrame, eliminating cross-process blocking from shared mutable state.
- All updates to the shared
dict_wandshared_dichappen in the main process, which is safe because Manager proxies handle top-level key modifications correctly. - The
printstatement runs after the sub-process has finished its work, so there's no blocking to prevent it from executing.
Quick Note on Your Original Code
You had a small typo in np.zeros(1,5) — this creates a 2D array, but extend expects a 1D sequence. Changing it to np.zeros(5) will give you the 5 zero values you're likely aiming for.
内容的提问来源于stack exchange,提问作者bib

