From 602989867b2736a2676fa93f5efdbe867dc77eae Mon Sep 17 00:00:00 2001 From: Danny <76977155+DNYFZR@users.noreply.github.com> Date: Fri, 19 Jun 2026 11:49:21 +0100 Subject: [PATCH 1/5] implement parallel processing --- src/agg.rs | 114 +++++++++++++++++++++++++++++++++++++---------------- 1 file changed, 81 insertions(+), 33 deletions(-) diff --git a/src/agg.rs b/src/agg.rs index bc71724..0ee4106 100644 --- a/src/agg.rs +++ b/src/agg.rs @@ -100,52 +100,100 @@ pub fn count_values( table: &DataFrame, partition_by: &str, iter_regex: &str, + parallel_limit: i64, ) -> Result { + // Get unique sim IDs let mut sim_ids = col_to_vec_i64(&table, partition_by); sim_ids.dedup(); + let n_sims = sim_ids.len(); - let mut sim_results = vec![]; - for sim_id in sim_ids { - let container = table - .get_column_names() + // Configure parallel loops + let batches = n_sims / parallel_limit as usize; + let remainder = n_sims - batches * parallel_limit as usize; + let loop_batches = match remainder { + 0 => batches, + _ => batches + 1, + }; + + // Get vec of timestep col names + let active_cols = table + .get_column_names() + .into_par_iter() + .filter(|c| c.contains(iter_regex)) + .map(|c| c.to_string()) + .collect::>(); + + // Create simulation profiles within parallel limits + let mut df: Vec = Vec::with_capacity(n_sims); + + for batch in 0..loop_batches { + let batch_size = if n_sims < parallel_limit as usize { + n_sims + } else if batch > batches { + remainder + } else { + parallel_limit as usize + }; + + // Run batches of frame chunks in parallel + let start_idx = batch * batch_size; + let end_idx = start_idx + (batch_size - 1); + let mut sim_res = sim_ids[start_idx..=end_idx] .into_par_iter() - .filter(|c| c.contains(iter_regex)) - .map(|c| { - let val_counts = table - .select(vec![c]) - .unwrap() - .rename(c, PlSmallStr::from_str("value")) - .unwrap() - .column("value") - .unwrap() - .as_series() - .unwrap() - .value_counts(false, false, PlSmallStr::from_str(c), false) - .expect("failed to count values"); - - return val_counts + .map(|sim_id| { + // Get simulation table within wider table + let sim_table = table .clone() .lazy() - .with_column(lit(sim_id).alias(partition_by)) - .select([col(partition_by), col("value"), col(c.to_string())]); + .filter(col(partition_by).eq(*sim_id)) + .collect() + .unwrap(); + + // Create val count for each timestep + let container = active_cols + .clone() + .into_par_iter() + .map(|c| { + let val_counts = sim_table + .select(vec![&c]) + .unwrap() + .rename(&c, PlSmallStr::from_str("value")) + .unwrap() + .column("value") + .unwrap() + .as_series() + .unwrap() + .value_counts(false, false, PlSmallStr::from_str(&c), false) + .expect("failed to count values"); + + return val_counts + .clone() + .lazy() + .with_column(lit(*sim_id).alias(partition_by)) + .select([col(partition_by), col("value"), col(&c)]); + }) + .collect::>(); + + // Join timestep cols into single df + let mut df = container[0].clone(); + for idx in 1..container.len() { + df = df.lazy().join( + container[idx].clone(), + [col(partition_by), col("value")], + [col(partition_by), col("value")], + JoinArgs::new(JoinType::Left), + ); + } + return df; }) .collect::>(); - // update table - let mut df = container[0].clone(); - for idx in 1..container.len() { - df = df.lazy().join( - container[idx].clone(), - [col(partition_by), col("value")], - [col(partition_by), col("value")], - JoinArgs::new(JoinType::Left), - ); - } - sim_results.push(df); + // push results to output container + df.append(&mut sim_res); } // Combine results - let df = concat(sim_results, UnionArgs::default())?; + let df = concat(df, UnionArgs::default())?; // Replace nulls with zero return Ok(df From cce81520c30c627e6398b90abace7fb56bc60658 Mon Sep 17 00:00:00 2001 From: Danny <76977155+DNYFZR@users.noreply.github.com> Date: Fri, 19 Jun 2026 11:49:37 +0100 Subject: [PATCH 2/5] expose parallel controls --- python/react_rs/__init__.py | 6 ++++++ src/lib.rs | 9 +++++++-- 2 files changed, 13 insertions(+), 2 deletions(-) diff --git a/python/react_rs/__init__.py b/python/react_rs/__init__.py index 30ca9c2..09b6a6e 100644 --- a/python/react_rs/__init__.py +++ b/python/react_rs/__init__.py @@ -90,6 +90,8 @@ def constrain( in each timestep (length must match number of timesteps in simulation output) - partition_by : string column name containing the simulation ID - run_method : string trigger for rust run method - options: full / batched / parallel + - para_limit : int value for the maximum number of parallelised simulations to run at any one time + Returns --- @@ -150,6 +152,7 @@ def profile( df: _pl.DataFrame, partition_by: str, iter_regex: str, + parallel_limit: int, ) -> _pl.DataFrame: """ Profile (Rust) @@ -163,6 +166,8 @@ def profile( - partition_by : string column name containing the simulation ID - iter_regex : string pattern for accessing the unique timesteps in the simulation output + - para_limit : int value for the maximum number of parallelised simulations to run at any one time + Returns --- @@ -175,4 +180,5 @@ def profile( df=df, partition_by=partition_by, iter_regex=iter_regex, + para_limit=parallel_limit, ) diff --git a/src/lib.rs b/src/lib.rs index e851268..b5c338a 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -138,7 +138,12 @@ mod react_rs { } #[pyfunction] - fn profile(df: &Bound<'_, PyAny>, partition_by: &str, iter_regex: &str) -> PyResult> { + fn profile( + df: &Bound<'_, PyAny>, + partition_by: &str, + iter_regex: &str, + para_limit: i64, + ) -> PyResult> { // Convert Python dataset to Rust let df = match import_py_dataframe(df) { Ok(df) => df, @@ -146,6 +151,6 @@ mod react_rs { }; // Execute value count on dataframe - return return_py_dataframe(agg::count_values(&df, partition_by, iter_regex)); + return return_py_dataframe(agg::count_values(&df, partition_by, iter_regex, para_limit)); } } From 8682ce6dbe273ff5c8c790fa37e19bda0a710b13 Mon Sep 17 00:00:00 2001 From: Danny <76977155+DNYFZR@users.noreply.github.com> Date: Fri, 19 Jun 2026 11:49:50 +0100 Subject: [PATCH 3/5] update README --- README.md | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/README.md b/README.md index f6bbdf1..27e4506 100644 --- a/README.md +++ b/README.md @@ -77,7 +77,7 @@ sim_result_constrained = react_rs.constrain( parallel_limit=10, ) -# Aggregate simulation outputs) +# Aggregate simulation outputs sim_result_agg = react_rs.aggregate( df=sim_result, partition_by="sim_id", @@ -99,12 +99,14 @@ sim_profile = react_rs.profile( df=sim_result, partition_by="sim_id", iter_regex="step", + parallel_limit=10, ) sim_constrained_profile = react_rs.profile( df=sim_result_constrained, partition_by="sim_id", iter_regex="step", + parallel_limit=10, ) ``` From efcacb6651f5136f4333b4d33e01983818f0fc7d Mon Sep 17 00:00:00 2001 From: Danny <76977155+DNYFZR@users.noreply.github.com> Date: Fri, 19 Jun 2026 11:50:00 +0100 Subject: [PATCH 4/5] update version --- Cargo.toml | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/Cargo.toml b/Cargo.toml index 0e057e8..3d8189f 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -1,5 +1,5 @@ [package] -version = "0.9.3" +version = "0.9.8" name = "react-rs" edition = "2024" From 08eaa7d728c4c065ba2440b8efdd7e722c32f296 Mon Sep 17 00:00:00 2001 From: Danny <76977155+DNYFZR@users.noreply.github.com> Date: Fri, 19 Jun 2026 12:19:07 +0100 Subject: [PATCH 5/5] add parallel limit to test --- tests/test_app.py | 1 + 1 file changed, 1 insertion(+) diff --git a/tests/test_app.py b/tests/test_app.py index 76c428e..c645589 100644 --- a/tests/test_app.py +++ b/tests/test_app.py @@ -169,6 +169,7 @@ def test_profiler(self, test_case): df=sim_result, partition_by="sim_id", iter_regex="step", + parallel_limit=test_case["parallel_limit"], ) if "value" not in sim_profile.columns: