Loading csst_dag/dag/_dispatcher.py +2 −10 Original line number Diff line number Diff line Loading @@ -118,7 +118,7 @@ def extract_basis_table(dlist: list[dict], basis_keys: tuple) -> table.Table: def split_data_basis(data_basis: table.Table, n_split: int = 1) -> list[table.Table]: """Split data basis into n_split parts.""" """Split data basis into n_split parts via obs_id""" assert ( np.unique(data_basis["dataset"]).size == 1 ), "Only one dataset is allowed for splitting." Loading Loading @@ -546,16 +546,8 @@ class Dispatcher: def dispatch_obsgroup_detector( plan_basis: table.Table, data_basis: table.Table, n_jobs: int = 1, # n_jobs: int = 1, ): # parallel if n_jobs != 1: task_list = joblib.Parallel(n_jobs=n_jobs)( joblib.delayed(Dispatcher.dispatch_obsid)(plan_basis, _) for _ in split_data_basis(data_basis, n_split=n_jobs) ) return sum(task_list, []) # unique obsgroup basis (using group_by) obsgroup_basis = plan_basis.group_by( keys=[ Loading Loading
csst_dag/dag/_dispatcher.py +2 −10 Original line number Diff line number Diff line Loading @@ -118,7 +118,7 @@ def extract_basis_table(dlist: list[dict], basis_keys: tuple) -> table.Table: def split_data_basis(data_basis: table.Table, n_split: int = 1) -> list[table.Table]: """Split data basis into n_split parts.""" """Split data basis into n_split parts via obs_id""" assert ( np.unique(data_basis["dataset"]).size == 1 ), "Only one dataset is allowed for splitting." Loading Loading @@ -546,16 +546,8 @@ class Dispatcher: def dispatch_obsgroup_detector( plan_basis: table.Table, data_basis: table.Table, n_jobs: int = 1, # n_jobs: int = 1, ): # parallel if n_jobs != 1: task_list = joblib.Parallel(n_jobs=n_jobs)( joblib.delayed(Dispatcher.dispatch_obsid)(plan_basis, _) for _ in split_data_basis(data_basis, n_split=n_jobs) ) return sum(task_list, []) # unique obsgroup basis (using group_by) obsgroup_basis = plan_basis.group_by( keys=[ Loading