perf: speed up mssql s3 ingest and add phase logs to job output (#8884)

* perf: speed up mssql s3 ingest and add phase logs to job output

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* chore: drop mssql s3 progress interval to 10s for better visibility

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

* chore: drop transient tests that compared against removed code path

Co-Authored-By: Claude Opus 4.7 (1M context) <noreply@anthropic.com>

---------

Co-authored-by: Claude Opus 4.7 (1M context) <noreply@anthropic.com>
This commit is contained in:
hugocasa
2026-04-20 18:20:02 +00:00
committed by GitHub
co-authored by Claude Opus 4.7
parent 94f27af838
commit 43a6b57581
7 changed files with 172 additions and 29 deletions
+116 -19
View File
@@ -979,25 +979,67 @@ impl Write for ChannelWriter {
}
#[cfg(not(feature = "parquet"))]
pub async fn convert_json_line_stream<E: Into<anyhow::Error>>(
mut _stream: impl futures::TryStreamExt<Item = Result<serde_json::Value, E>> + 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<std::time::Duration>,
}
#[cfg(not(feature = "parquet"))]
pub async fn convert_json_line_stream<V, E>(
mut _stream: impl futures::TryStreamExt<Item = Result<V, E>> + Unpin,
_output_format: S3ModeFormat,
) -> anyhow::Result<impl futures::TryStreamExt<Item = anyhow::Result<bytes::Bytes>>> {
Ok(async_stream::stream! {
_progress: Option<tokio::sync::mpsc::Sender<IngestStats>>,
) -> anyhow::Result<(
impl futures::TryStreamExt<Item = anyhow::Result<bytes::Bytes>>,
IngestStats,
)>
where
V: serde::Serialize,
E: Into<anyhow::Error>,
{
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<E: Into<anyhow::Error>>(
mut stream: impl TryStreamExt<Item = Result<serde_json::Value, E>> + 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<std::time::Duration>,
}
#[cfg(feature = "parquet")]
pub async fn convert_json_line_stream<V, E>(
mut stream: impl TryStreamExt<Item = Result<V, E>> + Unpin,
output_format: S3ModeFormat,
) -> anyhow::Result<impl TryStreamExt<Item = anyhow::Result<bytes::Bytes>>> {
progress: Option<tokio::sync::mpsc::Sender<IngestStats>>,
) -> anyhow::Result<(
impl TryStreamExt<Item = anyhow::Result<bytes::Bytes>>,
IngestStats,
)>
where
V: serde::Serialize,
E: Into<anyhow::Error>,
{
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<E: Into<anyhow::Error>>(
.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<Duration> = 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<Duration>,
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<E: Into<anyhow::Error>>(
Ok::<_, anyhow::Error>(())
});
Ok(tokio_stream::wrappers::ReceiverStream::new(rx))
Ok((
tokio_stream::wrappers::ReceiverStream::new(rx),
ingest_stats,
))
}
lazy_static::lazy_static! {
@@ -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())]);
+46 -4
View File
@@ -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::<windmill_object_store::IngestStats>(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 {
@@ -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())])
@@ -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 {
+2 -1
View File
@@ -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())]);
@@ -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 {