fix: improve concurrency re-scheduling at scale

This commit is contained in:
Ruben Fiszel
2024-05-25 15:56:56 +02:00
parent af5c31e6d4
commit e187fa6263
5 changed files with 43 additions and 14 deletions
@@ -0,0 +1,15 @@
{
"db_name": "PostgreSQL",
"query": "UPDATE queue\n SET running = false\n , started_at = null\n , scheduled_for = $1\n WHERE id = $2",
"describe": {
"columns": [],
"parameters": {
"Left": [
"Timestamptz",
"Uuid"
]
},
"nullable": []
},
"hash": "c0977384d5400e41f1d2aab6dea42c202bad922a6829c4768e5cf1fb6d90f810"
}
@@ -1,6 +1,6 @@
{
"db_name": "PostgreSQL",
"query": "SELECT CAST(ROUND(AVG(duration_ms) / 1000, 0) AS BIGINT) AS avg_duration_s FROM\n (SELECT duration_ms FROM completed_job WHERE script_path = $1\n ORDER BY started_at\n DESC LIMIT 10) AS t",
"query": "SELECT CAST(ROUND(AVG(duration_ms) / 1000, 0) AS BIGINT) AS avg_duration_s FROM\n (SELECT duration_ms FROM concurrency_key LEFT JOIN completed_job ON completed_job.id = concurrency_key.job_id WHERE key = $1 AND ended_at IS NOT NULL\n ORDER BY ended_at\n DESC LIMIT 10) AS t",
"describe": {
"columns": [
{
@@ -18,5 +18,5 @@
null
]
},
"hash": "a0ce703def7e976947513029874fb571893c75730e0a8feeb609852423404cf0"
"hash": "effb8e396c5811d3c45f880a6bf33ca564d3109901c93d4f312397e21656ce27"
}
+15 -5
View File
@@ -504,7 +504,7 @@ async fn update_flow(
nf.visible_to_runner_only.unwrap_or(false),
)
.execute(&mut tx)
.await?;
.await.map_err(|e| error::Error::InternalErr(format!("Error updating flow due to flow update: {e}")))?;
if nf.path != flow_path {
check_schedule_conflict(tx.transaction_mut(), &w_id, &nf.path).await?;
@@ -520,7 +520,7 @@ async fn update_flow(
.bind(&flow_path)
.bind(&w_id)
.fetch_all(&mut tx)
.await?;
.await.map_err(|e| error::Error::InternalErr(format!("Error updating flow due to related schedules update: {e}")))?;
let schedule = sqlx::query_as::<_, Schedule>(
"UPDATE schedule SET path = $1, script_path = $1 WHERE path = $2 AND workspace_id = $3 AND is_flow IS true RETURNING *")
@@ -528,7 +528,7 @@ async fn update_flow(
.bind(&flow_path)
.bind(&w_id)
.fetch_optional(&mut tx)
.await?;
.await.map_err(|e| error::Error::InternalErr(format!("Error updating flow due to related schedule update: {e}")))?;
if let Some(schedule) = schedule {
clear_schedule(tx.transaction_mut(), &flow_path, &w_id).await?;
@@ -618,14 +618,24 @@ async fn update_flow(
w_id
)
.execute(&mut new_tx)
.await?;
.await
.map_err(|e| {
error::Error::InternalErr(format!(
"Error updating flow due to updating dependency job field: {e}"
))
})?;
if let Some(old_dep_job) = old_dep_job {
sqlx::query!(
"UPDATE queue SET canceled = true WHERE id = $1",
old_dep_job
)
.execute(&mut new_tx)
.await?;
.await
.map_err(|e| {
error::Error::InternalErr(format!(
"Error updating flow due to cancelling dependency job: {e}"
))
})?;
}
new_tx.commit().await?;
+1
View File
@@ -793,6 +793,7 @@ pub async fn clear_schedule<'c>(
path: &str,
w_id: &str,
) -> Result<()> {
tracing::info!("Clearing schedule {}", path);
sqlx::query!(
"DELETE FROM queue WHERE schedule_path = $1 AND running = false AND workspace_id = $2 AND is_flow_step = false",
path,
+10 -7
View File
@@ -1773,13 +1773,14 @@ pub async fn pull<R: rsmq_async::RsmqConnection + Send + Clone>(
let job_uuid: Uuid = pulled_job.id;
let avg_script_duration: Option<i64> = sqlx::query_scalar!(
"SELECT CAST(ROUND(AVG(duration_ms) / 1000, 0) AS BIGINT) AS avg_duration_s FROM
(SELECT duration_ms FROM completed_job WHERE script_path = $1
ORDER BY started_at
(SELECT duration_ms FROM concurrency_key LEFT JOIN completed_job ON completed_job.id = concurrency_key.job_id WHERE key = $1 AND ended_at IS NOT NULL
ORDER BY ended_at
DESC LIMIT 10) AS t",
job_script_path
job_concurrency_key
)
.fetch_one(&mut tx)
.await?;
tracing::info!("avg script duration computed: {:?}", avg_script_duration);
// optimal scheduling is: 'older_job_in_concurrency_time_window_started_timestamp + script_avg_duration + concurrency_time_window_s'
let estimated_next_schedule_timestamp = ((min_started_at
@@ -1827,13 +1828,15 @@ pub async fn pull<R: rsmq_async::RsmqConnection + Send + Clone>(
tx.commit().await?;
} else {
// if using posgtres, then we're able to re-queue the entire batch of scheduled job for this script_path, so we do it
sqlx::query(&format!(
sqlx::query!(
"UPDATE queue
SET running = false
, started_at = null
, scheduled_for = '{estimated_next_schedule_timestamp}'
WHERE (id = '{job_uuid}') OR (script_path = '{job_script_path}' AND running = false AND scheduled_for <= now())"
))
, scheduled_for = $1
WHERE id = $2",
estimated_next_schedule_timestamp,
job_uuid,
)
.fetch_all(&mut tx)
.await
.map_err(|e| Error::InternalErr(format!("Could not update and re-queue job {job_uuid}. The job will be marked as running but it is not running: {e}")))?;