feat(backend): reduce memory allocation for big forloops of flows

This commit is contained in:
Ruben Fiszel
2023-03-30 11:52:54 +02:00
parent 4598194f08
commit f3f29ec24c
+111 -92
View File
@@ -7,7 +7,6 @@
*/
use std::collections::HashMap;
use std::iter;
use std::time::Duration;
use crate::jobs::{add_completed_job, add_completed_job_error, schedule_again_if_scheduled};
@@ -15,7 +14,6 @@ use crate::js_eval::{eval_timeout, IdContext};
use crate::{worker, AuthedClient, KEEP_JOB_DIR};
use anyhow::Context;
use async_recursion::async_recursion;
use dyn_iter::DynIter;
use serde_json::{json, Map, Value};
use tokio::sync::mpsc::Sender;
use tracing::instrument;
@@ -1207,45 +1205,62 @@ async fn push_next_flow_job(
let continue_on_same_worker =
flow.same_worker && module.suspend.is_none() && module.sleep.is_none();
let zipped = {
let all_args = match &next_status {
NextStatus::NextLoopIteration(NextIteration { new_args, .. }) => {
let args = args.as_ref().map(|args| {
let mut args = args.clone();
args.extend(new_args.clone());
args
});
DynIter::new(iter::once(args))
}
NextStatus::AllFlowJobs {
branchall: Some(BranchAllStatus { len, .. }),
iterator: None,
..
} => DynIter::new((0..*len).map(|_| args.as_ref().map(|args| args.clone()))),
NextStatus::AllFlowJobs {
branchall: None,
iterator: Some(Iterator { itered, .. }),
..
} => DynIter::new(itered.into_iter().enumerate().map(|(i, v)| {
args.as_ref().map(|args| {
let mut new_args = args.clone();
new_args.insert("iter".to_string(), json!({ "index": i, "value": v }));
new_args
})
})),
_ => DynIter::new(iter::once(args.as_ref().map(|m| m.clone()))),
};
job_payloads.into_iter().zip(all_args).collect::<Vec<_>>()
};
/* Finally, push the job into the queue */
let mut tx = db.begin().await?;
let mut uuids = vec![];
for (payload, args) in zipped {
let len = match &job_payloads {
ContinuePayload::SingleJob(_) => 1,
ContinuePayload::BranchAllJobs(payloads) => payloads.len(),
ContinuePayload::ForloopJobs { n, .. } => *n,
};
for i in (0..len).into_iter() {
let payload = match &job_payloads {
ContinuePayload::SingleJob(payload) => payload.clone(),
ContinuePayload::BranchAllJobs(payloads) => payloads[i].clone(),
ContinuePayload::ForloopJobs { modules, .. } => {
let mut fm = flow.failure_module.clone();
if let Some(mut failure_module) = flow.failure_module.clone() {
failure_module.id_append(&format!("{}/{}", status.step, i));
fm = Some(failure_module);
}
JobPayload::RawFlow {
value: FlowValue {
modules: (*modules).clone(),
failure_module: fm.clone(),
same_worker: flow.same_worker,
},
path: Some(format!("{}/loop-{}", flow_job.script_path(), i)),
}
}
};
let args = match &next_status {
NextStatus::AllFlowJobs {
branchall: Some(BranchAllStatus { .. }),
iterator: None,
..
} => args.as_ref().map(|args| args.clone()),
NextStatus::NextLoopIteration(NextIteration { new_args, .. }) => {
args.as_ref().map(|args| {
let mut args = args.clone();
args.extend(new_args.clone());
args
})
}
NextStatus::AllFlowJobs {
branchall: None,
iterator: Some(Iterator { itered, .. }),
..
} => args.as_ref().map(|args| {
let mut new_args = args.clone();
new_args.insert(
"iter".to_string(),
json!({ "index": i, "value": itered[i] }),
);
new_args
}),
_ => args.as_ref().map(|args| args.clone()),
};
let (ok, err) = match args {
Ok(v) => (Some(v), None),
Err(e) => (None, Some(e)),
@@ -1485,9 +1500,15 @@ enum NextStatus {
},
}
enum ContinuePayload {
SingleJob(JobPayload),
ForloopJobs { n: usize, modules: Vec<FlowModule> },
BranchAllJobs(Vec<JobPayload>),
}
enum NextFlowTransform {
EmptyInnerFlows,
Continue(Vec<JobPayload>, NextStatus),
Continue(ContinuePayload, NextStatus),
}
// a similar function exists on the backend
@@ -1521,13 +1542,19 @@ async fn compute_next_flow_transform<'c>(
match &module.value {
FlowModuleValue::Identity => Ok((
tx,
NextFlowTransform::Continue(vec![JobPayload::Identity], NextStatus::NextStep),
NextFlowTransform::Continue(
ContinuePayload::SingleJob(JobPayload::Identity),
NextStatus::NextStep,
),
)),
FlowModuleValue::Flow { path, .. } => {
let payload = JobPayload::Flow(path.to_string());
Ok((
tx,
NextFlowTransform::Continue(vec![payload], NextStatus::NextStep),
NextFlowTransform::Continue(
ContinuePayload::SingleJob(payload),
NextStatus::NextStep,
),
))
}
FlowModuleValue::Script { path: script_path, hash: script_hash, .. } => {
@@ -1541,7 +1568,10 @@ async fn compute_next_flow_transform<'c>(
};
Ok((
tx,
NextFlowTransform::Continue(vec![payload], NextStatus::NextStep),
NextFlowTransform::Continue(
ContinuePayload::SingleJob(payload),
NextStatus::NextStep,
),
))
}
FlowModuleValue::RawScript { path, content, language, lock, .. } => {
@@ -1551,12 +1581,12 @@ async fn compute_next_flow_transform<'c>(
Ok((
tx,
NextFlowTransform::Continue(
vec![JobPayload::Code(RawCode {
ContinuePayload::SingleJob(JobPayload::Code(RawCode {
path,
content: content.clone(),
language: language.clone(),
lock: lock.clone(),
})],
})),
NextStatus::NextStep,
),
))
@@ -1654,14 +1684,14 @@ async fn compute_next_flow_transform<'c>(
Ok((
tx,
NextFlowTransform::Continue(
vec![JobPayload::RawFlow {
ContinuePayload::SingleJob(JobPayload::RawFlow {
value: FlowValue {
modules: (*modules).clone(),
failure_module: fm,
same_worker: flow.same_worker,
},
path: Some(format!("{}/loop-{}", flow_job.script_path(), ns.index)),
}],
}),
NextStatus::NextLoopIteration(ns),
),
))
@@ -1669,23 +1699,10 @@ async fn compute_next_flow_transform<'c>(
LoopStatus::ParallelIteration { itered } => Ok((
tx,
NextFlowTransform::Continue(
(0..itered.len())
.map(|i| {
let mut fm = flow.failure_module.clone();
if let Some(mut failure_module) = flow.failure_module.clone() {
failure_module.id_append(&format!("{}/{}", status.step, i));
fm = Some(failure_module);
}
JobPayload::RawFlow {
value: FlowValue {
modules: (*modules).clone(),
failure_module: fm.clone(),
same_worker: flow.same_worker,
},
path: Some(format!("{}/loop-{}", flow_job.script_path(), i)),
}
})
.collect(),
ContinuePayload::ForloopJobs {
n: itered.len(),
modules: (*modules).clone(),
},
NextStatus::AllFlowJobs {
branchall: None,
iterator: Some(windmill_common::flow_status::Iterator {
@@ -1746,7 +1763,7 @@ async fn compute_next_flow_transform<'c>(
Ok((
tx,
NextFlowTransform::Continue(
vec![JobPayload::RawFlow {
ContinuePayload::SingleJob(JobPayload::RawFlow {
value: FlowValue {
modules,
failure_module: fm,
@@ -1757,7 +1774,7 @@ async fn compute_next_flow_transform<'c>(
flow_job.script_path(),
status.step
)),
}],
}),
NextStatus::BranchChosen(branch),
),
))
@@ -1773,32 +1790,34 @@ async fn compute_next_flow_transform<'c>(
return Ok((
tx,
NextFlowTransform::Continue(
branches
.iter()
.enumerate()
.map(|(i, b)| {
let mut fm = flow.failure_module.clone();
if let Some(mut failure_module) =
flow.failure_module.clone()
{
failure_module
.id_append(&format!("{}/{i}", status.step));
fm = Some(failure_module);
}
JobPayload::RawFlow {
value: FlowValue {
modules: b.modules.clone(),
failure_module: fm.clone(),
same_worker: flow.same_worker,
},
path: Some(format!(
"{}/branchall-{}",
flow_job.script_path(),
i
)),
}
})
.collect(),
ContinuePayload::BranchAllJobs(
branches
.iter()
.enumerate()
.map(|(i, b)| {
let mut fm = flow.failure_module.clone();
if let Some(mut failure_module) =
flow.failure_module.clone()
{
failure_module
.id_append(&format!("{}/{i}", status.step));
fm = Some(failure_module);
}
JobPayload::RawFlow {
value: FlowValue {
modules: b.modules.clone(),
failure_module: fm.clone(),
same_worker: flow.same_worker,
},
path: Some(format!(
"{}/branchall-{}",
flow_job.script_path(),
i
)),
}
})
.collect(),
),
NextStatus::AllFlowJobs {
branchall: Some(BranchAllStatus {
branch: 0,
@@ -1855,7 +1874,7 @@ async fn compute_next_flow_transform<'c>(
Ok((
tx,
NextFlowTransform::Continue(
vec![JobPayload::RawFlow {
ContinuePayload::SingleJob(JobPayload::RawFlow {
value: FlowValue {
modules,
failure_module: fm.clone(),
@@ -1866,7 +1885,7 @@ async fn compute_next_flow_transform<'c>(
flow_job.script_path(),
branch_status.branch
)),
}],
}),
NextStatus::NextBranchStep(NextBranch { status: branch_status, flow_jobs }),
),
))