From be1fcc016efacb0409ca390f7ba60688d3597c05 Mon Sep 17 00:00:00 2001 From: Ruben Fiszel Date: Mon, 17 Feb 2025 21:26:46 +0100 Subject: [PATCH] fix: improve que job indices for faster performances --- ...f073a79e300f3dd48f14122d1782eee663cd.json} | 7 +++--- ...cb4bdd64075939a4aa2c117e18372511ea7e0.json | 12 ---------- ...8008a9479bf4b3d7231371ebf26382ecde365.json | 12 ++++++++++ ...05d394a7cbcf0038c72a78add5c7b02ef5927.json | 2 +- backend/windmill-api/src/jobs.rs | 7 +++++- backend/windmill-common/src/lib.rs | 6 ++--- backend/windmill-queue/src/jobs.rs | 3 +++ backend/windmill-worker/src/worker.rs | 2 +- benchmarks/benchmark_oneoff.ts | 22 ++++++++++--------- benchmarks/benchmark_suite.ts | 2 +- 10 files changed, 43 insertions(+), 32 deletions(-) rename backend/.sqlx/{query-19cc8499f682ec34d54bc4f694cb281a9bd7f5431c646c6268513751fff95395.json => query-0cb0e912bc942af2b1ef784455f3f073a79e300f3dd48f14122d1782eee663cd.json} (72%) delete mode 100644 backend/.sqlx/query-47455b0ebaf999ab58b2cba3d74cb4bdd64075939a4aa2c117e18372511ea7e0.json create mode 100644 backend/.sqlx/query-b9b38d63af3670d1f11d5cbb82a8008a9479bf4b3d7231371ebf26382ecde365.json diff --git a/backend/.sqlx/query-19cc8499f682ec34d54bc4f694cb281a9bd7f5431c646c6268513751fff95395.json b/backend/.sqlx/query-0cb0e912bc942af2b1ef784455f3f073a79e300f3dd48f14122d1782eee663cd.json similarity index 72% rename from backend/.sqlx/query-19cc8499f682ec34d54bc4f694cb281a9bd7f5431c646c6268513751fff95395.json rename to backend/.sqlx/query-0cb0e912bc942af2b1ef784455f3f073a79e300f3dd48f14122d1782eee663cd.json index 0c379a7bdb..0a9e91b206 100644 --- a/backend/.sqlx/query-19cc8499f682ec34d54bc4f694cb281a9bd7f5431c646c6268513751fff95395.json +++ b/backend/.sqlx/query-0cb0e912bc942af2b1ef784455f3f073a79e300f3dd48f14122d1782eee663cd.json @@ -1,6 +1,6 @@ { "db_name": "PostgreSQL", - "query": "SELECT coalesce(COUNT(*) FILTER(WHERE suspend = 0 AND running = false), 0) as \"database_length!\", coalesce(COUNT(*) FILTER(WHERE suspend > 0), 0) as \"suspended!\" FROM v2_as_queue WHERE (workspace_id = $1 OR $2) AND scheduled_for <= now()", + "query": "SELECT coalesce(COUNT(*) FILTER(WHERE suspend = 0 AND running = false), 0) as \"database_length!\", coalesce(COUNT(*) FILTER(WHERE suspend > 0), 0) as \"suspended!\" FROM v2_as_queue WHERE (workspace_id = $1 OR $2) AND scheduled_for <= now() AND ($3::text[] IS NULL OR tag = ANY($3))", "describe": { "columns": [ { @@ -17,7 +17,8 @@ "parameters": { "Left": [ "Text", - "Bool" + "Bool", + "TextArray" ] }, "nullable": [ @@ -25,5 +26,5 @@ null ] }, - "hash": "19cc8499f682ec34d54bc4f694cb281a9bd7f5431c646c6268513751fff95395" + "hash": "0cb0e912bc942af2b1ef784455f3f073a79e300f3dd48f14122d1782eee663cd" } diff --git a/backend/.sqlx/query-47455b0ebaf999ab58b2cba3d74cb4bdd64075939a4aa2c117e18372511ea7e0.json b/backend/.sqlx/query-47455b0ebaf999ab58b2cba3d74cb4bdd64075939a4aa2c117e18372511ea7e0.json deleted file mode 100644 index 6ccaeece1c..0000000000 --- a/backend/.sqlx/query-47455b0ebaf999ab58b2cba3d74cb4bdd64075939a4aa2c117e18372511ea7e0.json +++ /dev/null @@ -1,12 +0,0 @@ -{ - "db_name": "PostgreSQL", - "query": "VACUUM (skip_locked) v2_job_queue, v2_job_runtime, v2_job_status", - "describe": { - "columns": [], - "parameters": { - "Left": [] - }, - "nullable": [] - }, - "hash": "47455b0ebaf999ab58b2cba3d74cb4bdd64075939a4aa2c117e18372511ea7e0" -} diff --git a/backend/.sqlx/query-b9b38d63af3670d1f11d5cbb82a8008a9479bf4b3d7231371ebf26382ecde365.json b/backend/.sqlx/query-b9b38d63af3670d1f11d5cbb82a8008a9479bf4b3d7231371ebf26382ecde365.json new file mode 100644 index 0000000000..641c94c555 --- /dev/null +++ b/backend/.sqlx/query-b9b38d63af3670d1f11d5cbb82a8008a9479bf4b3d7231371ebf26382ecde365.json @@ -0,0 +1,12 @@ +{ + "db_name": "PostgreSQL", + "query": "VACUUM v2_job_queue, v2_job_runtime, v2_job_status", + "describe": { + "columns": [], + "parameters": { + "Left": [] + }, + "nullable": [] + }, + "hash": "b9b38d63af3670d1f11d5cbb82a8008a9479bf4b3d7231371ebf26382ecde365" +} diff --git a/backend/.sqlx/query-ddf2eccb78a310ed00c7d8b9c3f05d394a7cbcf0038c72a78add5c7b02ef5927.json b/backend/.sqlx/query-ddf2eccb78a310ed00c7d8b9c3f05d394a7cbcf0038c72a78add5c7b02ef5927.json index c2dfed73a2..5bfff47576 100644 --- a/backend/.sqlx/query-ddf2eccb78a310ed00c7d8b9c3f05d394a7cbcf0038c72a78add5c7b02ef5927.json +++ b/backend/.sqlx/query-ddf2eccb78a310ed00c7d8b9c3f05d394a7cbcf0038c72a78add5c7b02ef5927.json @@ -15,7 +15,7 @@ ] }, "nullable": [ - null + true ] }, "hash": "ddf2eccb78a310ed00c7d8b9c3f05d394a7cbcf0038c72a78add5c7b02ef5927" diff --git a/backend/windmill-api/src/jobs.rs b/backend/windmill-api/src/jobs.rs index c678af5465..b5e1bc5a08 100644 --- a/backend/windmill-api/src/jobs.rs +++ b/backend/windmill-api/src/jobs.rs @@ -1647,6 +1647,7 @@ struct QueueStats { #[derive(Deserialize)] pub struct CountQueueJobsQuery { all_workspaces: Option, + tags: Option, } async fn count_queue_jobs( @@ -1654,12 +1655,16 @@ async fn count_queue_jobs( Path(w_id): Path, Query(cq): Query, ) -> error::JsonResult { + let tags = cq + .tags + .map(|t| t.split(',').map(|s| s.to_string()).collect::>()); Ok(Json( sqlx::query_as!( QueueStats, - "SELECT coalesce(COUNT(*) FILTER(WHERE suspend = 0 AND running = false), 0) as \"database_length!\", coalesce(COUNT(*) FILTER(WHERE suspend > 0), 0) as \"suspended!\" FROM v2_as_queue WHERE (workspace_id = $1 OR $2) AND scheduled_for <= now()", + "SELECT coalesce(COUNT(*) FILTER(WHERE suspend = 0 AND running = false), 0) as \"database_length!\", coalesce(COUNT(*) FILTER(WHERE suspend > 0), 0) as \"suspended!\" FROM v2_as_queue WHERE (workspace_id = $1 OR $2) AND scheduled_for <= now() AND ($3::text[] IS NULL OR tag = ANY($3))", w_id, w_id == "admins" && cq.all_workspaces.unwrap_or(false), + tags.as_ref().map(|v| v.as_slice()) ) .fetch_one(&db) .await?, diff --git a/backend/windmill-common/src/lib.rs b/backend/windmill-common/src/lib.rs index 396eafdf7f..c8996b3595 100644 --- a/backend/windmill-common/src/lib.rs +++ b/backend/windmill-common/src/lib.rs @@ -292,9 +292,9 @@ pub async fn connect( // .after_connect(move |conn, _| { // if worker_mode { // Box::pin(async move { - // sqlx::query("SET enable_seqscan = OFF;") - // .execute(conn) - // .await?; + // // sqlx::query("SET enable_seqscan = OFF;") + // // .execute(conn) + // // .await?; // Ok(()) // }) // } else { diff --git a/backend/windmill-queue/src/jobs.rs b/backend/windmill-queue/src/jobs.rs index 2d2025da92..5d360d3acb 100644 --- a/backend/windmill-queue/src/jobs.rs +++ b/backend/windmill-queue/src/jobs.rs @@ -2108,12 +2108,15 @@ async fn pull_single_job_and_mark_as_running_no_concurrency_limit<'c>( for query in queries.iter() { // tracing::info!("Pulling job with query: {}", query); + // let instant = std::time::Instant::now(); let r = sqlx::query_as::<_, PulledJob>(query) .bind(worker_name) .fetch_optional(db) .await?; if let Some(pulled_job) = r { + // tracing::info!("pulled job: {:?}", instant.elapsed().as_micros()); + highest_priority_job = Some(pulled_job); break; } diff --git a/backend/windmill-worker/src/worker.rs b/backend/windmill-worker/src/worker.rs index 5592e0d722..590b57227c 100644 --- a/backend/windmill-worker/src/worker.rs +++ b/backend/windmill-worker/src/worker.rs @@ -1277,7 +1277,7 @@ pub async fn run_worker( tokio::task::spawn( (async move { tracing::info!(worker = %worker_name, hostname = %hostname, "vacuuming queue"); - if let Err(e) = sqlx::query!("VACUUM (skip_locked) v2_job_queue, v2_job_runtime, v2_job_status") + if let Err(e) = sqlx::query!("VACUUM v2_job_queue, v2_job_runtime, v2_job_status") .execute(&db2) .await { diff --git a/benchmarks/benchmark_oneoff.ts b/benchmarks/benchmark_oneoff.ts index 285d471df8..62a3e3d5f7 100644 --- a/benchmarks/benchmark_oneoff.ts +++ b/benchmarks/benchmark_oneoff.ts @@ -37,6 +37,7 @@ async function verifyOutputs(uuids: string[], workspace: string) { console.log(`Incorrect results: ${incorrectResults}`); } +export const NON_TEST_TAGS = ["deno", "python", "go", "bash", "dedicated", "bun", "nativets", "flow"] export async function main({ host, email, @@ -96,11 +97,11 @@ export async function main({ windmill.setClient(final_token, host); const enc = (s: string) => new TextEncoder().encode(s); - async function getQueueCount() { + async function getQueueCount(tags?: string[]) { return ( await ( await fetch( - config.server + "/api/w/" + config.workspace_id + "/jobs/queue/count", + config.server + "/api/w/" + config.workspace_id + "/jobs/queue/count" + (tags && tags.length > 0 ? "?tags=" + tags.join(",") : ""), { headers: { ["Authorization"]: "Bearer " + config.token } } ) ).json() @@ -132,11 +133,11 @@ export async function main({ } let pastJobs = 0; - async function getCompletedJobsCount(): Promise { + async function getCompletedJobsCount(tags?: string[]): Promise { const completedJobs = ( await ( await fetch( - host + "/api/w/" + config.workspace_id + "/jobs/completed/count", + host + "/api/w/" + config.workspace_id + "/jobs/completed/count" + (tags && tags.length > 0 ? "?tags=" + tags.join(",") : ""), { headers: { ["Authorization"]: "Bearer " + config.token } } ) ).json() @@ -208,8 +209,9 @@ export async function main({ } let testOtherTag = false; - const otherTagTodo = 500000; if (testOtherTag) { + const otherTagTodo = 2000000; + let parsed = JSON.parse(body); parsed.tag = "test"; let nbody = JSON.stringify(parsed); @@ -237,7 +239,7 @@ export async function main({ } } - pastJobs = await getCompletedJobsCount(); + pastJobs = await getCompletedJobsCount(NON_TEST_TAGS); const response = await fetch( config.server + @@ -283,14 +285,14 @@ export async function main({ while (completedJobs < jobsSent) { const loopStart = Date.now(); if (!didStart) { - const actual_queue = await getQueueCount(); - if (actual_queue < jobsSent + (testOtherTag ? otherTagTodo : 0)) { + const actual_queue = await getQueueCount(NON_TEST_TAGS); + if (actual_queue < jobsSent) { start = Date.now(); didStart = true; } } else { const elapsed = start ? Date.now() - start : 0; - completedJobs = await getCompletedJobsCount(); + completedJobs = await getCompletedJobsCount(NON_TEST_TAGS); if (nStepsFlow > 0) { completedJobs = Math.floor(completedJobs / (nStepsFlow + 1)); } @@ -328,7 +330,7 @@ export async function main({ console.log(`avg. throughput (jobs/time): ${jobsSent / total_duration_sec}`); console.log("completed jobs", completedJobs); - console.log("queue length:", await getQueueCount()); + console.log("queue length:", await getQueueCount(NON_TEST_TAGS)); if ( !noVerify && diff --git a/benchmarks/benchmark_suite.ts b/benchmarks/benchmark_suite.ts index 7828c62d9b..e03148ce99 100644 --- a/benchmarks/benchmark_suite.ts +++ b/benchmarks/benchmark_suite.ts @@ -27,7 +27,7 @@ async function warmUp( token, workspace, kind: "noop", - jobs: 100000, + jobs: 50000, }); }