From 9512cab8078d991feac6d0034f394deb3cb774a9 Mon Sep 17 00:00:00 2001 From: duxin Date: Wed, 8 Jul 2026 14:09:54 +0800 Subject: [PATCH] =?UTF-8?q?fix:=20=E7=A7=BB=E9=99=A4=20=5Flocal=5Fkriging?= =?UTF-8?q?=20=E5=86=85=E9=83=A8=20multiprocessing.Pool=20=E9=81=BF?= =?UTF-8?q?=E5=85=8D=20Windows=20=E5=B5=8C=E5=A5=97=20spawn=20=E6=AD=BB?= =?UTF-8?q?=E9=94=81?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit 根因: step11 用 ProcessPoolExecutor 派发 CSV 到子进程, 子进程内 ContentMapper.process_data → _local_kriging 又创建 multiprocessing.Pool。Windows spawn 模式下嵌套 Pool 死锁。 修复: _local_kriging 内部改为顺序执行 16 个块。 每块 ~12s (650K 网格点 × 50 近邻), 总计 ~3 分钟, 完全可接受。 step11 层的 ProcessPoolExecutor 仍提供 CSV 级并行。 --- src/postprocessing/map.py | 20 +++++++++++--------- 1 file changed, 11 insertions(+), 9 deletions(-) diff --git a/src/postprocessing/map.py b/src/postprocessing/map.py index 2c13ee4..fe6ddc4 100644 --- a/src/postprocessing/map.py +++ b/src/postprocessing/map.py @@ -597,14 +597,11 @@ class ContentMapper: grid_y = grid_yy[:, 0] total_cells = len(grid_x) * len(grid_y) - _n_workers = min(os.cpu_count() or 4, 8) print(f"正在使用 局部克里金 (自适应分块 + 40% 重叠缓冲):" - f"网格={total_cells:,} 点, " - f"{_n_workers} workers") + f"网格={total_cells:,} 点") grid_content = self._local_kriging( points, values, grid_x, grid_y, - n_workers=_n_workers, n_closest_points=50, ) @@ -692,7 +689,7 @@ class ContentMapper: return grid_content def _local_kriging(self, points, values, grid_x, grid_y, - n_workers=4, n_closest_points=50): + n_closest_points=50): """局部克里金:自适应分块 + 重叠缓冲区 + 保护性近邻限制 1. 自适应块大小: 根据 extent 自动切分为 ~4×4 块 (16~25块) @@ -700,7 +697,6 @@ class ContentMapper: 3. 保护性近邻: n_closest_points=50,稀释极端异常值 4. 网格点仅使用严格不重叠的块范围(Buffer 仅用于筛选采样点) """ - import multiprocessing x_min, x_max = float(grid_x[0]), float(grid_x[-1]) y_min, y_max = float(grid_y[0]), float(grid_y[-1]) @@ -775,9 +771,15 @@ class ContentMapper: if len(tasks) <= 1: return self._local_krige_block(*tasks[0]) - print(f" 启动 {min(n_workers, len(tasks))} 个 worker 进程...") - with multiprocessing.Pool(processes=min(n_workers, len(tasks))) as pool: - results = pool.map(_local_krige_block_worker, tasks) + # 在 ProcessPoolExecutor 的子进程内(step11 批量模式)顺序执行, + # 避免 Windows spawn 模式下嵌套 multiprocessing 死锁。 + # 16 个块顺序跑 ~1-2 分钟,完全可接受。 + print(f" 顺序执行 {len(tasks)} 个局部 Kriging 块...") + results = [] + for i, task in enumerate(tasks): + if i % max(1, len(tasks) // 4) == 0 or i == len(tasks) - 1: + print(f" [LocalKrige] {i+1}/{len(tasks)} ...") + results.append(_local_krige_block_worker(task)) # 拼接:全 NaN 数组,逐块填回 grid_full = np.full((len(grid_y), len(grid_x)), np.nan, dtype=np.float64)