"""批量构建模块:支持并行或串行地处理多个独立换板打包任务。 提供 dataclass 用于封装单个构建任务、失败信息和批量结果, 并包含工作进程数计算及并行/串行执行调度逻辑。 """ from __future__ import annotations import os from concurrent.futures import ProcessPoolExecutor, as_completed from dataclasses import dataclass from typing import Sequence from .builder import build_packed_3mf from .models import BuildOptions, BuildResult, PlateJob # 独立批量构建的默认最大并行工作进程数 INDIVIDUAL_BATCH_MAX_WORKERS = 8 @dataclass(frozen=True) class IndividualBuildTask: """单个独立构建任务,包含一个板作业和对应的构建选项。""" job: PlateJob """要处理的板作业""" options: BuildOptions """构建选项配置""" @dataclass(frozen=True) class IndividualBuildFailure: """记录单个构建任务失败时的详细信息。""" index: int """任务在原始列表中的索引""" job: PlateJob """失败的板作业""" error: str """错误描述信息""" @dataclass(frozen=True) class IndividualBatchBuildResult: """一批独立构建任务的整体结果汇总。""" results: tuple[BuildResult, ...] """成功完成的所有构建结果""" failures: tuple[IndividualBuildFailure, ...] """所有失败的构建任务记录""" worker_count: int """本次批量构建使用的并行工作进程数""" def individual_batch_worker_count(task_count: int, max_workers: int | None = None) -> int: """根据任务数量计算合适的并行工作进程数。 综合考虑用户指定的最大进程数、CPU 核心数和内置上限, 返回一个合理的并行度。 Args: task_count: 待执行的任务总数 max_workers: 用户指定的最大并行工作进程数,为 None 时自动检测 Returns: 应使用的并行工作进程数,任务数 <= 0 时返回 0 """ if task_count <= 0: return 0 if max_workers is not None: # 以用户指定值为上限,最少 1 个进程 return min(task_count, max(1, int(max_workers))) # 自动检测:取任务数、CPU 核心数、内置上限三者中的最小值 detected_cpu_count = os.cpu_count() or 1 return min(task_count, max(1, detected_cpu_count), INDIVIDUAL_BATCH_MAX_WORKERS) def _build_individual_task(task: IndividualBuildTask) -> BuildResult: """执行单个构建任务,调用核心构建函数。 Args: task: 单个独立构建任务 Returns: BuildResult: 构建结果 """ return build_packed_3mf([task.job], task.options) def _run_individual_batch_serial( tasks: Sequence[IndividualBuildTask], worker_count: int, ) -> IndividualBatchBuildResult: """串行方式执行一批独立构建任务。 逐个执行任务,捕获每个任务的异常并记录为失败。 Args: tasks: 要执行的任务序列 worker_count: 工作进程数(此处仅用于结果记录) Returns: IndividualBatchBuildResult: 包含成功和失败记录的批量结果 """ results: list[BuildResult] = [] failures: list[IndividualBuildFailure] = [] for index, task in enumerate(tasks): try: # 逐个同步执行构建 results.append(_build_individual_task(task)) except Exception as exc: # 捕获任意异常,记录为失败项 failures.append(IndividualBuildFailure(index, task.job, str(exc))) return IndividualBatchBuildResult(tuple(results), tuple(failures), worker_count) def run_individual_batch_builds( tasks: Sequence[IndividualBuildTask], max_workers: int | None = None, ) -> IndividualBatchBuildResult: """并行执行一批独立构建任务。 当计算出的工作进程数 > 1 时,使用 ProcessPoolExecutor 并行处理; 否则回退到串行执行。并行模式下按原始索引位置收集结果。 Args: tasks: 要执行的任务序列 max_workers: 最大并行工作进程数,为 None 时自动检测 Returns: IndividualBatchBuildResult: 包含成功和失败记录的批量结果 """ task_list = list(tasks) worker_count = individual_batch_worker_count(len(task_list), max_workers) # 单进程或无任务时,使用串行路径 if worker_count <= 1: return _run_individual_batch_serial(task_list, worker_count) # 预分配结果槽位,保证按原始顺序输出 ordered_results: list[BuildResult | None] = [None] * len(task_list) failures: list[IndividualBuildFailure] = [] # 使用进程池并行执行 with ProcessPoolExecutor(max_workers=worker_count) as executor: # 提交所有任务并记录 future 到索引的映射 future_to_index = { executor.submit(_build_individual_task, task): index for index, task in enumerate(task_list) } # 按完成顺序收集结果 for future in as_completed(future_to_index): index = future_to_index[future] task = task_list[index] try: # 获取结果,放入对应位置 ordered_results[index] = future.result() except Exception as exc: # 捕获任务执行异常,记录为失败 failures.append(IndividualBuildFailure(index, task.job, str(exc))) # 过滤掉未成功的位置(None),保留有效结果 results = tuple(result for result in ordered_results if result is not None) return IndividualBatchBuildResult( tuple(results), # 按原始索引排序失败记录,便于阅读 tuple(sorted(failures, key=lambda item: item.index)), worker_count, )