diff --git a/backend/windmill-object-store/src/lib.rs b/backend/windmill-object-store/src/lib.rs index 11235a8533..78a0e22c92 100644 --- a/backend/windmill-object-store/src/lib.rs +++ b/backend/windmill-object-store/src/lib.rs @@ -979,25 +979,67 @@ impl Write for ChannelWriter { } #[cfg(not(feature = "parquet"))] -pub async fn convert_json_line_stream>( - mut _stream: impl futures::TryStreamExt> + Unpin, +#[derive(Debug, Clone, Copy, Default)] +pub struct IngestStats { + pub rows: u64, + pub bytes: u64, + pub elapsed: std::time::Duration, + pub fetch_wait: std::time::Duration, + pub write_time: std::time::Duration, + pub first_row_latency: Option, +} + +#[cfg(not(feature = "parquet"))] +pub async fn convert_json_line_stream( + mut _stream: impl futures::TryStreamExt> + Unpin, _output_format: S3ModeFormat, -) -> anyhow::Result>> { - Ok(async_stream::stream! { + _progress: Option>, +) -> anyhow::Result<( + impl futures::TryStreamExt>, + IngestStats, +)> +where + V: serde::Serialize, + E: Into, +{ + let stream = async_stream::stream! { yield Err(anyhow::anyhow!("Parquet feature is not enabled. Cannot convert JSON line stream.")); - }) + }; + Ok((stream, IngestStats::default())) } #[cfg(feature = "parquet")] -pub async fn convert_json_line_stream>( - mut stream: impl TryStreamExt> + Unpin, +#[derive(Debug, Clone, Copy)] +pub struct IngestStats { + pub rows: u64, + pub bytes: u64, + pub elapsed: std::time::Duration, + pub fetch_wait: std::time::Duration, + pub write_time: std::time::Duration, + pub first_row_latency: Option, +} + +#[cfg(feature = "parquet")] +pub async fn convert_json_line_stream( + mut stream: impl TryStreamExt> + Unpin, output_format: S3ModeFormat, -) -> anyhow::Result>> { + progress: Option>, +) -> anyhow::Result<( + impl TryStreamExt>, + IngestStats, +)> +where + V: serde::Serialize, + E: Into, +{ const MAX_MPSC_SIZE: usize = 1000; + const WRITE_BUF_CAPACITY: usize = 256 * 1024; + const PROGRESS_INTERVAL_SECS: u64 = 10; use datafusion::{execution::context::SessionContext, prelude::NdJsonReadOptions}; use futures::StreamExt; use std::path::PathBuf; + use std::time::{Duration, Instant}; use tokio::io::AsyncWriteExt; let mut path = PathBuf::from(std::env::temp_dir()); @@ -1006,26 +1048,78 @@ pub async fn convert_json_line_stream>( .to_str() .ok_or_else(|| anyhow::anyhow!("Invalid path"))?; - let mut file: tokio::fs::File = tokio::fs::File::create(&path).await.map_err(to_anyhow)?; + let file: tokio::fs::File = tokio::fs::File::create(&path).await.map_err(to_anyhow)?; + let mut file = tokio::io::BufWriter::with_capacity(WRITE_BUF_CAPACITY, file); - while let Some(chunk) = stream.next().await { - match chunk { - Ok(chunk) => { - let b: bytes::Bytes = serde_json::to_string(&chunk)?.into(); - file.write_all(&b).await?; - file.write_all(b"\n").await?; + let ingest_start = Instant::now(); + let mut first_row_latency: Option = None; + let mut row_count: u64 = 0; + let mut bytes_written: u64 = 0; + let mut write_time = Duration::ZERO; + let mut progress_timer = tokio::time::interval(Duration::from_secs(PROGRESS_INTERVAL_SECS)); + progress_timer.tick().await; // drop the immediate tick + let build_stats = |row_count: u64, + bytes_written: u64, + write_time: Duration, + first_row_latency: Option, + elapsed: Duration| + -> IngestStats { + IngestStats { + rows: row_count, + bytes: bytes_written, + elapsed, + fetch_wait: elapsed.saturating_sub(write_time), + write_time, + first_row_latency, + } + }; + + let mut done = false; + while !done { + tokio::select! { + chunk = stream.next() => { + match chunk { + Some(Ok(chunk)) => { + if first_row_latency.is_none() { + first_row_latency = Some(ingest_start.elapsed()); + } + let t_write = Instant::now(); + let mut s = serde_json::to_string(&chunk)?; + s.push('\n'); + bytes_written += s.len() as u64; + file.write_all(s.as_bytes()).await?; + write_time += t_write.elapsed(); + row_count += 1; + } + Some(Err(e)) => { + tokio::fs::remove_file(&path).await?; + return Err(e.into()); + } + None => done = true, + } } - Err(e) => { - tokio::fs::remove_file(&path).await?; - return Err(e.into()); + _ = progress_timer.tick(), if progress.is_some() => { + let stats = build_stats(row_count, bytes_written, write_time, first_row_latency, ingest_start.elapsed()); + if let Some(tx) = &progress { + let _ = tx.try_send(stats); + } } } } file.flush().await?; + let file = file.into_inner(); file.sync_all().await?; drop(file); + let ingest_stats = build_stats( + row_count, + bytes_written, + write_time, + first_row_latency, + ingest_start.elapsed(), + ); + let ctx = SessionContext::new(); ctx.register_json("my_table", path_str, NdJsonReadOptions::default()) .await @@ -1109,7 +1203,10 @@ pub async fn convert_json_line_stream>( Ok::<_, anyhow::Error>(()) }); - Ok(tokio_stream::wrappers::ReceiverStream::new(rx)) + Ok(( + tokio_stream::wrappers::ReceiverStream::new(rx), + ingest_stats, + )) } lazy_static::lazy_static! { diff --git a/backend/windmill-worker/src/bigquery_executor.rs b/backend/windmill-worker/src/bigquery_executor.rs index b092e135b6..c9e803deaf 100644 --- a/backend/windmill-worker/src/bigquery_executor.rs +++ b/backend/windmill-worker/src/bigquery_executor.rs @@ -203,8 +203,8 @@ fn do_bigquery_inner<'a>( } }; - let stream = - convert_json_line_stream(rows_stream.boxed(), s3.format).await?; + let (stream, _) = + convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; s3.upload(stream.boxed()).await?; return Ok(vec![to_raw_value(&s3.to_return_s3_obj())]); diff --git a/backend/windmill-worker/src/mssql_executor.rs b/backend/windmill-worker/src/mssql_executor.rs index c1f8ae2fe7..d7be6b0d1c 100644 --- a/backend/windmill-worker/src/mssql_executor.rs +++ b/backend/windmill-worker/src/mssql_executor.rs @@ -244,19 +244,61 @@ pub async fn do_mssql( // fetching data in an asynchronous manner, if needed. if let Some(s3) = s3 { + let s3_format = s3.format; let rows_stream = async_stream::stream! { let mut stream = prepared_query.query(&mut client).await.map_err(to_anyhow)?.into_row_stream().map(|row| { - let raw_value = row_to_json(row.map_err(to_anyhow)?).map_err(to_anyhow); - let json = raw_value.and_then(|raw_value| serde_json::from_str(raw_value.get()).map_err(to_anyhow)); - json + row_to_json(row.map_err(to_anyhow)?).map_err(to_anyhow) }); while let Some(row) = stream.next().await { yield row; } }; - let stream = convert_json_line_stream(rows_stream.boxed(), s3.format).await?; + let (progress_tx, mut progress_rx) = + tokio::sync::mpsc::channel::(8); + let progress_conn = conn.clone(); + let progress_job_id = job.id; + let progress_workspace = job.workspace_id.clone(); + let progress_task = tokio::spawn(async move { + while let Some(stats) = progress_rx.recv().await { + let logs = format!( + "\nMSSQL s3 stream progress: {} rows, {:.1} MB | elapsed {:.1}s (db-fetch {:.1}s, write {:.1}s)", + stats.rows, + stats.bytes as f64 / 1_048_576.0, + stats.elapsed.as_secs_f64(), + stats.fetch_wait.as_secs_f64(), + stats.write_time.as_secs_f64(), + ); + append_logs(&progress_job_id, &progress_workspace, logs, &progress_conn).await; + } + }); + + let (stream, ingest_stats) = + convert_json_line_stream(rows_stream.boxed(), s3_format, Some(progress_tx)).await?; + let _ = progress_task.await; + + let ingest_log = format!( + "\nMSSQL s3 stream ingest done ({:?}): {} rows, {:.1} MB in {:.2}s (db-fetch {:.2}s, write {:.2}s, first row after {:.2}s)", + s3_format, + ingest_stats.rows, + ingest_stats.bytes as f64 / 1_048_576.0, + ingest_stats.elapsed.as_secs_f64(), + ingest_stats.fetch_wait.as_secs_f64(), + ingest_stats.write_time.as_secs_f64(), + ingest_stats.first_row_latency.unwrap_or_default().as_secs_f64(), + ); + append_logs(&job.id, &job.workspace_id, ingest_log, conn).await; + + let upload_start = std::time::Instant::now(); s3.upload(stream.boxed()).await?; + let upload_elapsed = upload_start.elapsed(); + + let final_log = format!( + "\nMSSQL s3 stream upload+transcode done: {:.2}s | total {:.2}s", + upload_elapsed.as_secs_f64(), + (ingest_stats.elapsed + upload_elapsed).as_secs_f64(), + ); + append_logs(&job.id, &job.workspace_id, final_log, conn).await; Ok(to_raw_value(&s3.to_return_s3_obj())) } else { diff --git a/backend/windmill-worker/src/mysql_executor.rs b/backend/windmill-worker/src/mysql_executor.rs index a3a9468459..aee7de16f8 100644 --- a/backend/windmill-worker/src/mysql_executor.rs +++ b/backend/windmill-worker/src/mysql_executor.rs @@ -107,7 +107,8 @@ fn do_mysql_inner<'a>( } }; - let stream = convert_json_line_stream(rows_stream.boxed(), s3.format).await?; + let (stream, _) = + convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; s3.upload(stream.boxed()).await?; Ok(vec![to_raw_value(&s3.to_return_s3_obj())]) diff --git a/backend/windmill-worker/src/oracledb_executor.rs b/backend/windmill-worker/src/oracledb_executor.rs index bc90ba7811..55ffac467b 100644 --- a/backend/windmill-worker/src/oracledb_executor.rs +++ b/backend/windmill-worker/src/oracledb_executor.rs @@ -162,7 +162,8 @@ pub fn do_oracledb_inner<'a>( } if let Some(s3) = s3 { - let stream = convert_json_line_stream(rows_stream.boxed(), s3.format).await?; + let (stream, _) = + convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; s3.upload(stream.boxed()).await?; return Ok(vec![to_raw_value(&s3.to_return_s3_obj())]); } else { diff --git a/backend/windmill-worker/src/pg_executor.rs b/backend/windmill-worker/src/pg_executor.rs index fbe0cfbcb7..2e577da6a8 100644 --- a/backend/windmill-worker/src/pg_executor.rs +++ b/backend/windmill-worker/src/pg_executor.rs @@ -203,7 +203,8 @@ fn do_postgresql_inner<'a>( row_result.and_then(|row| postgres_row_to_json_value(row).map_err(to_anyhow)) }); - let stream = convert_json_line_stream(rows_stream.boxed(), s3.format).await?; + let (stream, _) = + convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; s3.upload(stream.boxed()).await?; return Ok(vec![to_raw_value(&s3.to_return_s3_obj())]); diff --git a/backend/windmill-worker/src/snowflake_executor.rs b/backend/windmill-worker/src/snowflake_executor.rs index 90d06287bc..0873624f54 100644 --- a/backend/windmill-worker/src/snowflake_executor.rs +++ b/backend/windmill-worker/src/snowflake_executor.rs @@ -501,7 +501,8 @@ fn do_snowflake_inner<'a>( if let Some(s3) = s3 { let rows_stream = rows_stream.map(|r| serde_json::value::to_value(&r?).map_err(to_anyhow)); - let stream = convert_json_line_stream(rows_stream.boxed(), s3.format).await?; + let (stream, _) = + convert_json_line_stream(rows_stream.boxed(), s3.format, None).await?; s3.upload(stream.boxed()).await?; Ok(vec![to_raw_value(&s3.to_return_s3_obj())]) } else {