From d68cc8ba3000a2dc5d242e7da2ac23abaf0f8e6b Mon Sep 17 00:00:00 2001 From: jianguotian Date: Wed, 16 Sep 2026 11:23:16 +0800 Subject: [PATCH 1/4] perf(mosaic): prefetch row groups with a byte budget --- crates/paimon/src/arrow/format/mosaic.rs | 315 +++++++++++++++++++++-- 1 file changed, 289 insertions(+), 26 deletions(-) diff --git a/crates/paimon/src/arrow/format/mosaic.rs b/crates/paimon/src/arrow/format/mosaic.rs index b10c756c0..11f8483a6 100644 --- a/crates/paimon/src/arrow/format/mosaic.rs +++ b/crates/paimon/src/arrow/format/mosaic.rs @@ -36,7 +36,7 @@ use paimon_mosaic_core::reader::{InputFile, MosaicReader, ReaderAccess}; use paimon_mosaic_core::schema::MosaicSchema; use paimon_mosaic_core::stats::ColumnStats; use paimon_mosaic_core::values::Value as MosaicValue; -use std::collections::{HashMap, HashSet}; +use std::collections::{HashMap, HashSet, VecDeque}; use std::io; use std::ops::Range; use std::sync::Arc; @@ -44,6 +44,8 @@ use std::sync::Arc; pub(crate) struct MosaicFormatReader; const DEFAULT_BATCH_SIZE: usize = 8192; +const DEFAULT_PREFETCH_ROW_GROUPS: usize = 8; +const DEFAULT_PREFETCH_MAX_BYTES: usize = 64 * 1024 * 1024; #[async_trait] impl FormatFileReader for MosaicFormatReader { @@ -84,6 +86,8 @@ impl FormatFileReader for MosaicFormatReader { batch_size, row_selection, handle, + prefetch_row_groups: DEFAULT_PREFETCH_ROW_GROUPS, + prefetch_max_bytes: DEFAULT_PREFETCH_MAX_BYTES, }, |batch| futures::executor::block_on(batch_tx.send(Ok(batch))).is_ok(), ); @@ -113,6 +117,20 @@ struct MosaicReadRequest { batch_size: usize, row_selection: Option>, handle: tokio::runtime::Handle, + prefetch_row_groups: usize, + prefetch_max_bytes: usize, +} + +struct PlannedRowGroup { + index: usize, + rows: usize, + selected_slices: Option>, + estimated_bytes: usize, +} + +struct PendingRowGroup { + estimated_bytes: usize, + handle: std::thread::JoinHandle>, } fn read_mosaic_batches_blocking( @@ -127,9 +145,13 @@ fn read_mosaic_batches_blocking( batch_size, row_selection, handle, + prefetch_row_groups, + prefetch_max_bytes, } = request; - let mosaic_reader = MosaicReader::new(FileReadInputFile::new(reader, handle), file_size) - .map_err(mosaic_read_error)?; + let mosaic_reader = Arc::new( + MosaicReader::new(FileReadInputFile::new(reader, handle), file_size) + .map_err(mosaic_read_error)?, + ); let file_column_names = mosaic_reader .schema() @@ -169,6 +191,8 @@ fn read_mosaic_batches_blocking( build_file_column_indices(mosaic_reader.schema(), &predicates.file_fields) }); + let estimated_row_bytes = estimated_row_bytes(&existing_scan_fields); + let mut planned = VecDeque::new(); let mut row_group_start = 0usize; for row_group_index in 0..mosaic_reader.num_row_groups() { let row_group_rows = mosaic_reader @@ -210,36 +234,188 @@ fn read_mosaic_batches_blocking( } } - let batch = if all_scan_columns_missing { - let row_count = selected_slices - .as_ref() - .map_or(row_group_rows, |slices| selected_row_count(slices)); - empty_batch(read_schema.clone(), row_count)? - } else { - let names = projected_names - .iter() - .map(String::as_str) - .collect::>(); - let mut row_group_reader = mosaic_reader - .row_group_reader_by_names(row_group_index, &names) - .map_err(mosaic_read_error)?; + planned.push_back(PlannedRowGroup { + index: row_group_index, + rows: row_group_rows, + selected_slices, + estimated_bytes: row_group_rows.saturating_mul(estimated_row_bytes), + }); + } + + if all_scan_columns_missing || prefetch_row_groups == 0 { + while let Some(plan) = planned.pop_front() { + let batch = if all_scan_columns_missing { + let row_count = plan + .selected_slices + .as_ref() + .map_or(plan.rows, |slices| selected_row_count(slices)); + empty_batch(read_schema.clone(), row_count)? + } else { + read_planned_row_group(&mosaic_reader, &plan, &projected_names, &read_schema)? + }; + let batch = apply_mosaic_residual(batch, predicates.as_ref(), &existing_scan_fields)?; + if !send_mosaic_batch_chunks(batch, batch_size, &mut send_batch) { + return Ok(()); + } + } + return Ok(()); + } + + let projected_names = Arc::new(projected_names); + let mut pending = VecDeque::::new(); + let mut pending_bytes = 0usize; + let mut outcome = Ok(()); + + while !planned.is_empty() || !pending.is_empty() { + while pending.len() < prefetch_row_groups { + let Some(next) = planned.front() else { + break; + }; + // Always allow the head row group so a too-small budget cannot + // deadlock progress. Every additional prefetched group must fit. + if !pending.is_empty() + && pending_bytes.saturating_add(next.estimated_bytes) > prefetch_max_bytes + { + break; + } + let plan = planned.pop_front().unwrap(); + let estimated_bytes = plan.estimated_bytes; + let reader = Arc::clone(&mosaic_reader); + let names = Arc::clone(&projected_names); + let schema = Arc::clone(&read_schema); + let handle = + std::thread::spawn(move || read_planned_row_group(&reader, &plan, &names, &schema)); + pending_bytes = pending_bytes.saturating_add(estimated_bytes); + pending.push_back(PendingRowGroup { + estimated_bytes, + handle, + }); + } - let batch = row_group_reader.read_columns().map_err(mosaic_read_error)?; - take_row_slices(batch, selected_slices.as_deref(), &read_schema)? + let Some(next) = pending.pop_front() else { + outcome = Err(Error::UnexpectedError { + message: "Mosaic row-group prefetch made no progress".to_string(), + source: None, + }); + break; }; - let batch = match predicates.as_ref() { - Some(predicates) => { - filter_record_batch_by_predicates(batch, predicates, &existing_scan_fields)? + pending_bytes = pending_bytes.saturating_sub(next.estimated_bytes); + let batch = match join_prefetched_row_group(next.handle) { + Ok(batch) => batch, + Err(error) => { + outcome = Err(error); + break; } - None => batch, }; - for chunk in split_batch(batch, batch_size) { - if !send_batch(chunk) { - return Ok(()); + let batch = match apply_mosaic_residual(batch, predicates.as_ref(), &existing_scan_fields) { + Ok(batch) => batch, + Err(error) => { + outcome = Err(error); + break; } + }; + if !send_mosaic_batch_chunks(batch, batch_size, &mut send_batch) { + planned.clear(); + break; } } - Ok(()) + + // A detached row-group task would keep the reader and its async storage + // work alive after cancellation or an error. Join every scheduled task + // before returning; preserve the first user-visible failure. + while let Some(next) = pending.pop_front() { + if let Err(error) = join_prefetched_row_group(next.handle) { + if outcome.is_ok() { + outcome = Err(error); + } + } + } + outcome +} + +fn read_planned_row_group( + mosaic_reader: &MosaicReader, + plan: &PlannedRowGroup, + projected_names: &[String], + read_schema: &SchemaRef, +) -> crate::Result { + let names = projected_names + .iter() + .map(String::as_str) + .collect::>(); + let mut row_group_reader = mosaic_reader + .row_group_reader_by_names(plan.index, &names) + .map_err(mosaic_read_error)?; + let batch = row_group_reader.read_columns().map_err(mosaic_read_error)?; + take_row_slices(batch, plan.selected_slices.as_deref(), read_schema) +} + +fn join_prefetched_row_group( + handle: std::thread::JoinHandle>, +) -> crate::Result { + handle.join().map_err(|panic| Error::UnexpectedError { + message: format!( + "Mosaic row-group prefetch task panicked: {}", + panic_payload_message(panic) + ), + source: None, + })? +} + +fn panic_payload_message(panic: Box) -> String { + if let Some(message) = panic.downcast_ref::<&str>() { + (*message).to_string() + } else if let Some(message) = panic.downcast_ref::() { + message.clone() + } else { + "unknown panic payload".to_string() + } +} + +fn apply_mosaic_residual( + batch: RecordBatch, + predicates: Option<&FilePredicates>, + existing_scan_fields: &[DataField], +) -> crate::Result { + match predicates { + Some(predicates) => { + filter_record_batch_by_predicates(batch, predicates, existing_scan_fields) + } + None => Ok(batch), + } +} + +fn send_mosaic_batch_chunks( + batch: RecordBatch, + batch_size: usize, + send_batch: &mut impl FnMut(RecordBatch) -> bool, +) -> bool { + split_batch(batch, batch_size).into_iter().all(send_batch) +} + +fn estimated_row_bytes(fields: &[DataField]) -> usize { + fields + .iter() + .map(|field| match field.data_type() { + PaimonDataType::Boolean(_) | PaimonDataType::TinyInt(_) => 2, + PaimonDataType::SmallInt(_) => 3, + PaimonDataType::Int(_) + | PaimonDataType::Float(_) + | PaimonDataType::Date(_) + | PaimonDataType::Time(_) => 5, + PaimonDataType::BigInt(_) + | PaimonDataType::Double(_) + | PaimonDataType::Timestamp(_) + | PaimonDataType::LocalZonedTimestamp(_) => 9, + PaimonDataType::Decimal(_) => 17, + PaimonDataType::Char(_) + | PaimonDataType::VarChar(_) + | PaimonDataType::Binary(_) + | PaimonDataType::VarBinary(_) => 40, + _ => 128, + }) + .sum::() + .max(1) } struct MosaicRowGroupStats<'a> { @@ -637,6 +813,7 @@ mod tests { use paimon_mosaic_core::spec::COMPRESSION_NONE; use paimon_mosaic_core::writer::{MosaicWriter, OutputFile, WriterOptions}; use std::ops::Range; + use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::{Arc, Mutex}; struct TestFileRead { @@ -672,6 +849,26 @@ mod tests { runtime_thread: std::thread::ThreadId, } + struct ConcurrentTrackingFileRead { + data: Bytes, + active: Arc, + max_active: Arc, + } + + #[async_trait] + impl FileRead for ConcurrentTrackingFileRead { + async fn read(&self, range: Range) -> crate::Result { + let active = self.active.fetch_add(1, Ordering::SeqCst) + 1; + self.max_active.fetch_max(active, Ordering::SeqCst); + tokio::time::sleep(std::time::Duration::from_millis(20)).await; + let start = usize::try_from(range.start).unwrap(); + let end = usize::try_from(range.end).unwrap(); + let result = self.data.slice(start..end); + self.active.fetch_sub(1, Ordering::SeqCst); + Ok(result) + } + } + #[async_trait] impl FileRead for RuntimeThreadFileRead { async fn read(&self, range: Range) -> crate::Result { @@ -967,6 +1164,48 @@ mod tests { ) } + async fn read_with_prefetch_tracking( + data: Bytes, + prefetch_row_groups: usize, + prefetch_max_bytes: usize, + ) -> (Vec, usize) { + let active = Arc::new(AtomicUsize::new(0)); + let max_active = Arc::new(AtomicUsize::new(0)); + let file_size = data.len() as u64; + let fields = vec![data_fields()[0].clone()]; + let read = ConcurrentTrackingFileRead { + data, + active, + max_active: Arc::clone(&max_active), + }; + let handle = tokio::runtime::Handle::current(); + let batches = tokio::task::spawn_blocking(move || { + let mut batches = Vec::new(); + read_mosaic_batches_blocking( + MosaicReadRequest { + reader: Box::new(read), + file_size, + read_fields: fields, + predicates: None, + batch_size: DEFAULT_BATCH_SIZE, + row_selection: None, + handle, + prefetch_row_groups, + prefetch_max_bytes, + }, + |batch| { + batches.push(batch); + true + }, + )?; + Ok::<_, Error>(batches) + }) + .await + .unwrap() + .unwrap(); + (batches, max_active.load(Ordering::SeqCst)) + } + fn timestamp_fields() -> Vec { vec![ DataField::new( @@ -1038,6 +1277,30 @@ mod tests { assert_eq!(ids.value(4), 5); } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn test_prefetch_overlaps_row_group_reads_and_preserves_order() { + let data = multi_row_group_mosaic(Vec::new()); + let (batches, max_active) = read_with_prefetch_tracking(data, 3, usize::MAX).await; + + assert_eq!(collect_i32_column(&batches, 0), vec![1, 2, 10, 11, 20, 21]); + assert!( + max_active >= 2, + "row-group prefetch should overlap storage reads, max_active={max_active}" + ); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn test_prefetch_byte_budget_limits_groups_ahead() { + let data = multi_row_group_mosaic(Vec::new()); + let (batches, max_active) = read_with_prefetch_tracking(data, 3, 1).await; + + assert_eq!(collect_i32_column(&batches, 0), vec![1, 2, 10, 11, 20, 21]); + assert_eq!( + max_active, 1, + "a one-byte budget should admit only the mandatory head group" + ); + } + #[tokio::test] async fn test_row_id_predicate_is_rejected() { let data = write_mosaic(&sample_batch()); From 95bae46eda293c7f72417feffb8fc4edccf41a4c Mon Sep 17 00:00:00 2001 From: jianguotian Date: Wed, 16 Sep 2026 13:46:38 +0800 Subject: [PATCH 2/4] fix(mosaic): bound row group concurrency per scan --- crates/paimon/src/arrow/format/mod.rs | 11 +- crates/paimon/src/arrow/format/mosaic.rs | 551 +++++++++++++++--- .../paimon/src/arrow/parquet_read_budget.rs | 102 +++- crates/paimon/src/spec/core_options.rs | 58 ++ crates/paimon/src/table/table_read.rs | 3 +- 5 files changed, 632 insertions(+), 93 deletions(-) diff --git a/crates/paimon/src/arrow/format/mod.rs b/crates/paimon/src/arrow/format/mod.rs index f67ae5651..ceec05171 100644 --- a/crates/paimon/src/arrow/format/mod.rs +++ b/crates/paimon/src/arrow/format/mod.rs @@ -193,7 +193,7 @@ pub(crate) fn create_format_reader( create_format_reader_with_budget(path, blob_as_descriptor, read_fields, None) } -/// Create a format reader with a scan-shared Parquet resource budget. +/// Create a format reader with scan-shared row-group resource budgets. pub(crate) fn create_format_reader_with_budget( path: &str, blob_as_descriptor: bool, @@ -219,10 +219,11 @@ pub(crate) fn create_format_reader_with_budget( Box::new(row::RowFormatReader) } else { if lower.ends_with(".mosaic") { - return Ok(shredding::maybe_wrap_reader( - Box::new(mosaic::MosaicFormatReader), - read_fields, - )); + let reader: Box = Box::new(match parquet_read_budget { + Some(read_budget) => mosaic::MosaicFormatReader::with_read_budget(read_budget), + None => mosaic::MosaicFormatReader::default(), + }); + return Ok(shredding::maybe_wrap_reader(reader, read_fields)); } #[cfg(feature = "vortex")] if lower.ends_with(".vortex") { diff --git a/crates/paimon/src/arrow/format/mosaic.rs b/crates/paimon/src/arrow/format/mosaic.rs index 11f8483a6..33912bb68 100644 --- a/crates/paimon/src/arrow/format/mosaic.rs +++ b/crates/paimon/src/arrow/format/mosaic.rs @@ -20,7 +20,9 @@ use crate::arrow::build_target_arrow_schema; use crate::arrow::filtering::{ predicates_may_match_with_schema, remap_predicates_to_file, StatsAccessor, }; +use crate::arrow::parquet_read_budget::MosaicReadPermit; use crate::arrow::residual::{filter_record_batch_by_predicates, widen_scan_fields}; +use crate::arrow::ParquetReadBudget; use crate::io::FileRead; use crate::spec::{DataField, DataType as PaimonDataType, Datum, Predicate}; use crate::table::{ArrowRecordBatchStream, RowRange}; @@ -31,6 +33,7 @@ use arrow_schema::{DataType as ArrowDataType, SchemaRef}; use async_stream::try_stream; use async_trait::async_trait; use bytes::Bytes; +use crossbeam_channel::{Receiver, Sender}; use futures::{SinkExt, StreamExt}; use paimon_mosaic_core::reader::{InputFile, MosaicReader, ReaderAccess}; use paimon_mosaic_core::schema::MosaicSchema; @@ -39,13 +42,29 @@ use paimon_mosaic_core::values::Value as MosaicValue; use std::collections::{HashMap, HashSet, VecDeque}; use std::io; use std::ops::Range; -use std::sync::Arc; +use std::panic::{catch_unwind, AssertUnwindSafe}; +use std::sync::{Arc, Condvar, Mutex, OnceLock}; -pub(crate) struct MosaicFormatReader; +pub(crate) struct MosaicFormatReader { + read_budget: Arc, +} const DEFAULT_BATCH_SIZE: usize = 8192; const DEFAULT_PREFETCH_ROW_GROUPS: usize = 8; const DEFAULT_PREFETCH_MAX_BYTES: usize = 64 * 1024 * 1024; +const MAX_EXECUTOR_WORKERS: usize = 32; + +impl MosaicFormatReader { + pub(crate) fn with_read_budget(read_budget: Arc) -> Self { + Self { read_budget } + } +} + +impl Default for MosaicFormatReader { + fn default() -> Self { + Self::with_read_budget(Arc::new(ParquetReadBudget::default())) + } +} #[async_trait] impl FormatFileReader for MosaicFormatReader { @@ -74,6 +93,7 @@ impl FormatFileReader for MosaicFormatReader { file_fields: predicates.file_fields.clone(), }); let batch_size = batch_size.unwrap_or(DEFAULT_BATCH_SIZE); + let read_budget = Arc::clone(&self.read_budget); let (mut batch_tx, mut batch_rx) = futures::channel::mpsc::channel(1); let read_task = tokio::task::spawn_blocking(move || { @@ -88,6 +108,7 @@ impl FormatFileReader for MosaicFormatReader { handle, prefetch_row_groups: DEFAULT_PREFETCH_ROW_GROUPS, prefetch_max_bytes: DEFAULT_PREFETCH_MAX_BYTES, + read_budget, }, |batch| futures::executor::block_on(batch_tx.send(Ok(batch))).is_ok(), ); @@ -98,7 +119,12 @@ impl FormatFileReader for MosaicFormatReader { Ok(try_stream! { while let Some(batch) = batch_rx.next().await { - yield batch?; + let BudgetedBatch { batch, permit } = batch?; + yield batch; + // Keep both the byte and row-group permits while the decoded + // batch is queued or yielded. Advancing or dropping the + // stream releases them. + drop(permit); } read_task.await.map_err(|e| Error::DataInvalid { message: format!("Mosaic read task failed: {e}"), @@ -119,23 +145,238 @@ struct MosaicReadRequest { handle: tokio::runtime::Handle, prefetch_row_groups: usize, prefetch_max_bytes: usize, + read_budget: Arc, } struct PlannedRowGroup { index: usize, rows: usize, selected_slices: Option>, - estimated_bytes: usize, + budget: RowGroupBudget, } struct PendingRowGroup { - estimated_bytes: usize, - handle: std::thread::JoinHandle>, + permit: MosaicPrefetchPermit, + task: MosaicTask, +} + +struct BudgetedBatch { + batch: RecordBatch, + permit: Option>, +} + +#[derive(Clone, Copy, Debug, PartialEq, Eq)] +enum RowGroupBudget { + Bytes(usize), + /// Mosaic 0.2 does not expose projected physical sizes. Variable-width + /// and nested projections therefore run one row group at a time. + Exclusive, +} + +#[derive(Default)] +struct ByteBudgetState { + bytes: usize, + groups: usize, + exclusive: bool, +} + +struct MosaicByteBudgetInner { + max_bytes: usize, + max_groups: usize, + state: Mutex, + released: Condvar, +} + +#[derive(Clone)] +struct MosaicByteBudget { + inner: Arc, +} + +impl MosaicByteBudget { + fn new(max_bytes: usize, max_groups: usize) -> Self { + Self { + inner: Arc::new(MosaicByteBudgetInner { + max_bytes, + max_groups, + state: Mutex::new(ByteBudgetState::default()), + released: Condvar::new(), + }), + } + } + + fn try_acquire(&self, requirement: RowGroupBudget) -> Option { + let mut state = self + .inner + .state + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + if !self.can_acquire(&state, requirement) { + return None; + } + Self::charge(&mut state, requirement); + Some(MosaicBytePermit { + inner: Arc::clone(&self.inner), + requirement, + }) + } + + fn acquire(&self, requirement: RowGroupBudget) -> MosaicBytePermit { + let mut state = self + .inner + .state + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + while !self.can_acquire(&state, requirement) { + state = self + .inner + .released + .wait(state) + .unwrap_or_else(|poisoned| poisoned.into_inner()); + } + Self::charge(&mut state, requirement); + MosaicBytePermit { + inner: Arc::clone(&self.inner), + requirement, + } + } + + fn can_acquire(&self, state: &ByteBudgetState, requirement: RowGroupBudget) -> bool { + if state.groups >= self.inner.max_groups || state.exclusive { + return false; + } + match requirement { + RowGroupBudget::Exclusive => state.groups == 0, + RowGroupBudget::Bytes(bytes) => { + // An oversized head group may consume the whole budget so the + // reader always makes progress, but it cannot overlap peers. + state.groups == 0 || state.bytes.saturating_add(bytes) <= self.inner.max_bytes + } + } + } + + fn charge(state: &mut ByteBudgetState, requirement: RowGroupBudget) { + state.groups += 1; + match requirement { + RowGroupBudget::Bytes(bytes) => state.bytes = state.bytes.saturating_add(bytes), + RowGroupBudget::Exclusive => state.exclusive = true, + } + } +} + +struct MosaicBytePermit { + inner: Arc, + requirement: RowGroupBudget, +} + +impl Drop for MosaicBytePermit { + fn drop(&mut self) { + let mut state = self + .inner + .state + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + state.groups -= 1; + match self.requirement { + RowGroupBudget::Bytes(bytes) => state.bytes = state.bytes.saturating_sub(bytes), + RowGroupBudget::Exclusive => state.exclusive = false, + } + self.inner.released.notify_all(); + } +} + +struct MosaicPrefetchPermit { + _bytes: MosaicBytePermit, + _row_group: MosaicReadPermit, +} + +type MosaicJob = Box; + +struct MosaicExecutor { + sender: Sender, +} + +struct MosaicTask { + receiver: Receiver>, +} + +impl MosaicExecutor { + fn new() -> crate::Result { + let workers = std::thread::available_parallelism() + .map(|parallelism| parallelism.get().saturating_mul(2)) + .unwrap_or(DEFAULT_PREFETCH_ROW_GROUPS) + .clamp(DEFAULT_PREFETCH_ROW_GROUPS, MAX_EXECUTOR_WORKERS); + let (sender, receiver): (Sender, Receiver) = + crossbeam_channel::bounded(workers.saturating_mul(2)); + for worker in 0..workers { + let receiver = receiver.clone(); + std::thread::Builder::new() + .name(format!("paimon-mosaic-{worker}")) + .spawn(move || { + while let Ok(job) = receiver.recv() { + job(); + } + }) + .map_err(|error| Error::UnexpectedError { + message: "failed to start Mosaic row-group executor".to_string(), + source: Some(Box::new(error)), + })?; + } + Ok(Self { sender }) + } + + fn submit(&self, task: F) -> crate::Result> + where + T: Send + 'static, + F: FnOnce() -> crate::Result + Send + 'static, + { + let (sender, receiver) = crossbeam_channel::bounded(1); + self.sender + .send(Box::new(move || { + let result = match catch_unwind(AssertUnwindSafe(task)) { + Ok(result) => result, + Err(panic) => Err(Error::UnexpectedError { + message: format!( + "Mosaic row-group prefetch task panicked: {}", + panic_payload_message(panic) + ), + source: None, + }), + }; + let _ = sender.send(result); + })) + .map_err(|_| Error::UnexpectedError { + message: "Mosaic row-group executor stopped before accepting a task".to_string(), + source: None, + })?; + Ok(MosaicTask { receiver }) + } +} + +impl MosaicTask { + fn join(self) -> crate::Result { + self.receiver + .recv() + .map_err(|error| Error::UnexpectedError { + message: "Mosaic row-group executor dropped a task result".to_string(), + source: Some(Box::new(error)), + })? + } +} + +fn mosaic_executor() -> crate::Result<&'static MosaicExecutor> { + static EXECUTOR: OnceLock> = OnceLock::new(); + match EXECUTOR.get_or_init(|| MosaicExecutor::new().map_err(|error| error.to_string())) { + Ok(executor) => Ok(executor), + Err(message) => Err(Error::UnexpectedError { + message: message.clone(), + source: None, + }), + } } fn read_mosaic_batches_blocking( request: MosaicReadRequest, - mut send_batch: impl FnMut(RecordBatch) -> bool, + mut send_batch: impl FnMut(BudgetedBatch) -> bool, ) -> crate::Result<()> { let MosaicReadRequest { reader, @@ -147,9 +388,10 @@ fn read_mosaic_batches_blocking( handle, prefetch_row_groups, prefetch_max_bytes, + read_budget, } = request; let mosaic_reader = Arc::new( - MosaicReader::new(FileReadInputFile::new(reader, handle), file_size) + MosaicReader::new(FileReadInputFile::new(reader, handle.clone()), file_size) .map_err(mosaic_read_error)?, ); @@ -191,7 +433,6 @@ fn read_mosaic_batches_blocking( build_file_column_indices(mosaic_reader.schema(), &predicates.file_fields) }); - let estimated_row_bytes = estimated_row_bytes(&existing_scan_fields); let mut planned = VecDeque::new(); let mut row_group_start = 0usize; for row_group_index in 0..mosaic_reader.num_row_groups() { @@ -238,7 +479,7 @@ fn read_mosaic_batches_blocking( index: row_group_index, rows: row_group_rows, selected_slices, - estimated_bytes: row_group_rows.saturating_mul(estimated_row_bytes), + budget: projected_row_group_budget(&existing_scan_fields, row_group_rows), }); } @@ -254,7 +495,7 @@ fn read_mosaic_batches_blocking( read_planned_row_group(&mosaic_reader, &plan, &projected_names, &read_schema)? }; let batch = apply_mosaic_residual(batch, predicates.as_ref(), &existing_scan_fields)?; - if !send_mosaic_batch_chunks(batch, batch_size, &mut send_batch) { + if !send_mosaic_batch_chunks(batch, batch_size, None, &mut send_batch) { return Ok(()); } } @@ -262,45 +503,83 @@ fn read_mosaic_batches_blocking( } let projected_names = Arc::new(projected_names); + let executor = mosaic_executor()?; let mut pending = VecDeque::::new(); - let mut pending_bytes = 0usize; + let byte_budget = MosaicByteBudget::new(prefetch_max_bytes, prefetch_row_groups); let mut outcome = Ok(()); - while !planned.is_empty() || !pending.is_empty() { - while pending.len() < prefetch_row_groups { - let Some(next) = planned.front() else { + 'read: while !planned.is_empty() || !pending.is_empty() { + while let Some(next) = planned.front() { + let Some(bytes) = byte_budget.try_acquire(next.budget) else { break; }; - // Always allow the head row group so a too-small budget cannot - // deadlock progress. Every additional prefetched group must fit. - if !pending.is_empty() - && pending_bytes.saturating_add(next.estimated_bytes) > prefetch_max_bytes + let row_group = match read_budget.try_acquire_mosaic() { + Ok(Some(permit)) => permit, + Ok(None) => { + drop(bytes); + break; + } + Err(error) => { + outcome = Err(error); + break 'read; + } + }; + let plan = planned.pop_front().unwrap(); + let permit = MosaicPrefetchPermit { + _bytes: bytes, + _row_group: row_group, + }; + let reader = Arc::clone(&mosaic_reader); + let names = Arc::clone(&projected_names); + let schema = Arc::clone(&read_schema); + let task = match executor + .submit(move || read_planned_row_group(&reader, &plan, &names, &schema)) { + Ok(task) => task, + Err(error) => { + outcome = Err(error); + break 'read; + } + }; + pending.push_back(PendingRowGroup { permit, task }); + } + + if pending.is_empty() { + let Some(next) = planned.front() else { break; - } + }; + // Earlier decoded batches can still own the budget downstream. + // With no task left to join, wait for their permits to be released. + let bytes = byte_budget.acquire(next.budget); + let row_group = match handle.block_on(read_budget.acquire_mosaic()) { + Ok(permit) => permit, + Err(error) => { + outcome = Err(error); + break 'read; + } + }; + let permit = MosaicPrefetchPermit { + _bytes: bytes, + _row_group: row_group, + }; let plan = planned.pop_front().unwrap(); - let estimated_bytes = plan.estimated_bytes; let reader = Arc::clone(&mosaic_reader); let names = Arc::clone(&projected_names); let schema = Arc::clone(&read_schema); - let handle = - std::thread::spawn(move || read_planned_row_group(&reader, &plan, &names, &schema)); - pending_bytes = pending_bytes.saturating_add(estimated_bytes); - pending.push_back(PendingRowGroup { - estimated_bytes, - handle, - }); + let task = match executor + .submit(move || read_planned_row_group(&reader, &plan, &names, &schema)) + { + Ok(task) => task, + Err(error) => { + outcome = Err(error); + break 'read; + } + }; + pending.push_back(PendingRowGroup { permit, task }); } - let Some(next) = pending.pop_front() else { - outcome = Err(Error::UnexpectedError { - message: "Mosaic row-group prefetch made no progress".to_string(), - source: None, - }); - break; - }; - pending_bytes = pending_bytes.saturating_sub(next.estimated_bytes); - let batch = match join_prefetched_row_group(next.handle) { + let PendingRowGroup { permit, task } = pending.pop_front().unwrap(); + let batch = match task.join() { Ok(batch) => batch, Err(error) => { outcome = Err(error); @@ -314,7 +593,7 @@ fn read_mosaic_batches_blocking( break; } }; - if !send_mosaic_batch_chunks(batch, batch_size, &mut send_batch) { + if !send_mosaic_batch_chunks(batch, batch_size, Some(Arc::new(permit)), &mut send_batch) { planned.clear(); break; } @@ -324,7 +603,7 @@ fn read_mosaic_batches_blocking( // work alive after cancellation or an error. Join every scheduled task // before returning; preserve the first user-visible failure. while let Some(next) = pending.pop_front() { - if let Err(error) = join_prefetched_row_group(next.handle) { + if let Err(error) = next.task.join() { if outcome.is_ok() { outcome = Err(error); } @@ -350,18 +629,6 @@ fn read_planned_row_group( take_row_slices(batch, plan.selected_slices.as_deref(), read_schema) } -fn join_prefetched_row_group( - handle: std::thread::JoinHandle>, -) -> crate::Result { - handle.join().map_err(|panic| Error::UnexpectedError { - message: format!( - "Mosaic row-group prefetch task panicked: {}", - panic_payload_message(panic) - ), - source: None, - })? -} - fn panic_payload_message(panic: Box) -> String { if let Some(message) = panic.downcast_ref::<&str>() { (*message).to_string() @@ -388,34 +655,44 @@ fn apply_mosaic_residual( fn send_mosaic_batch_chunks( batch: RecordBatch, batch_size: usize, - send_batch: &mut impl FnMut(RecordBatch) -> bool, + permit: Option>, + send_batch: &mut impl FnMut(BudgetedBatch) -> bool, ) -> bool { - split_batch(batch, batch_size).into_iter().all(send_batch) + split_batch(batch, batch_size).into_iter().all(|batch| { + send_batch(BudgetedBatch { + batch, + permit: permit.clone(), + }) + }) } -fn estimated_row_bytes(fields: &[DataField]) -> usize { - fields - .iter() - .map(|field| match field.data_type() { - PaimonDataType::Boolean(_) | PaimonDataType::TinyInt(_) => 2, - PaimonDataType::SmallInt(_) => 3, - PaimonDataType::Int(_) - | PaimonDataType::Float(_) - | PaimonDataType::Date(_) - | PaimonDataType::Time(_) => 5, - PaimonDataType::BigInt(_) - | PaimonDataType::Double(_) - | PaimonDataType::Timestamp(_) - | PaimonDataType::LocalZonedTimestamp(_) => 9, - PaimonDataType::Decimal(_) => 17, - PaimonDataType::Char(_) - | PaimonDataType::VarChar(_) - | PaimonDataType::Binary(_) - | PaimonDataType::VarBinary(_) => 40, - _ => 128, - }) - .sum::() - .max(1) +fn projected_row_group_budget(fields: &[DataField], rows: usize) -> RowGroupBudget { + let Some(row_bytes) = fields.iter().try_fold(0usize, |total, field| { + fixed_width_row_bytes(field.data_type()).map(|bytes| total.saturating_add(bytes)) + }) else { + return RowGroupBudget::Exclusive; + }; + RowGroupBudget::Bytes(rows.saturating_mul(row_bytes.max(1))) +} + +fn fixed_width_row_bytes(data_type: &PaimonDataType) -> Option { + // One byte per row for validity is deliberately conservative (Arrow uses + // one bit). These are upper bounds for decoded fixed-width array buffers, + // unlike a guessed average for variable-width values. + match data_type { + PaimonDataType::Boolean(_) | PaimonDataType::TinyInt(_) => Some(2), + PaimonDataType::SmallInt(_) => Some(3), + PaimonDataType::Int(_) + | PaimonDataType::Float(_) + | PaimonDataType::Date(_) + | PaimonDataType::Time(_) => Some(5), + PaimonDataType::BigInt(_) + | PaimonDataType::Double(_) + | PaimonDataType::Timestamp(_) + | PaimonDataType::LocalZonedTimestamp(_) => Some(9), + PaimonDataType::Decimal(_) => Some(17), + _ => None, + } } struct MosaicRowGroupStats<'a> { @@ -1060,7 +1337,7 @@ mod tests { row_selection: Option>, ) -> crate::Result> { let file_size = data.len() as u64; - MosaicFormatReader + MosaicFormatReader::default() .read_batch_stream( Box::new(TestFileRead { data }), file_size, @@ -1082,7 +1359,7 @@ mod tests { ) -> crate::Result>> { let file_size = data.len() as u64; let calls = Arc::new(Mutex::new(Vec::new())); - let _: Vec = MosaicFormatReader + let _: Vec = MosaicFormatReader::default() .read_batch_stream( Box::new(TrackingFileRead { data, @@ -1108,7 +1385,7 @@ mod tests { ) -> crate::Result>> { let file_size = data.len() as u64; let calls = Arc::new(Mutex::new(Vec::new())); - let _: Vec = MosaicFormatReader + let _: Vec = MosaicFormatReader::default() .read_batch_stream( Box::new(TrackingFileRead { data, @@ -1166,19 +1443,27 @@ mod tests { async fn read_with_prefetch_tracking( data: Bytes, + fields: Vec, prefetch_row_groups: usize, prefetch_max_bytes: usize, ) -> (Vec, usize) { let active = Arc::new(AtomicUsize::new(0)); let max_active = Arc::new(AtomicUsize::new(0)); let file_size = data.len() as u64; - let fields = vec![data_fields()[0].clone()]; let read = ConcurrentTrackingFileRead { data, active, max_active: Arc::clone(&max_active), }; let handle = tokio::runtime::Handle::current(); + let read_budget = Arc::new( + ParquetReadBudget::new_with_mosaic_parallelism( + 8, + 256 * 1024 * 1024, + prefetch_row_groups.max(1), + ) + .unwrap(), + ); let batches = tokio::task::spawn_blocking(move || { let mut batches = Vec::new(); read_mosaic_batches_blocking( @@ -1192,8 +1477,9 @@ mod tests { handle, prefetch_row_groups, prefetch_max_bytes, + read_budget, }, - |batch| { + |BudgetedBatch { batch, .. }| { batches.push(batch); true }, @@ -1280,7 +1566,8 @@ mod tests { #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_prefetch_overlaps_row_group_reads_and_preserves_order() { let data = multi_row_group_mosaic(Vec::new()); - let (batches, max_active) = read_with_prefetch_tracking(data, 3, usize::MAX).await; + let fields = vec![data_fields()[0].clone()]; + let (batches, max_active) = read_with_prefetch_tracking(data, fields, 3, usize::MAX).await; assert_eq!(collect_i32_column(&batches, 0), vec![1, 2, 10, 11, 20, 21]); assert!( @@ -1292,7 +1579,8 @@ mod tests { #[tokio::test(flavor = "multi_thread", worker_threads = 4)] async fn test_prefetch_byte_budget_limits_groups_ahead() { let data = multi_row_group_mosaic(Vec::new()); - let (batches, max_active) = read_with_prefetch_tracking(data, 3, 1).await; + let fields = vec![data_fields()[0].clone()]; + let (batches, max_active) = read_with_prefetch_tracking(data, fields, 3, 1).await; assert_eq!(collect_i32_column(&batches, 0), vec![1, 2, 10, 11, 20, 21]); assert_eq!( @@ -1301,6 +1589,99 @@ mod tests { ); } + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn test_prefetch_byte_budget_is_cumulative() { + let data = multi_row_group_mosaic(Vec::new()); + let fields = vec![data_fields()[0].clone()]; + // Each two-row Int group is conservatively charged 10 bytes, so this + // admits two groups but not all three. + let (batches, max_active) = read_with_prefetch_tracking(data, fields, 3, 20).await; + + assert_eq!(collect_i32_column(&batches, 0), vec![1, 2, 10, 11, 20, 21]); + assert_eq!(max_active, 2); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn test_variable_width_prefetch_is_exclusive() { + let data = multi_row_group_mosaic(Vec::new()); + let fields = vec![data_fields()[1].clone()]; + let (batches, max_active) = read_with_prefetch_tracking(data, fields, 3, usize::MAX).await; + + assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 6); + assert_eq!( + max_active, 1, + "variable-width row groups must not use a fixed per-row guess" + ); + } + + #[test] + fn test_fixed_width_budget_is_cumulative_and_oversized_head_progresses() { + assert_eq!( + projected_row_group_budget(&[data_fields()[0].clone()], 2), + RowGroupBudget::Bytes(10) + ); + assert_eq!( + projected_row_group_budget(&[data_fields()[1].clone()], 2), + RowGroupBudget::Exclusive + ); + + let budget = MosaicByteBudget::new(20, 8); + let first = budget.try_acquire(RowGroupBudget::Bytes(10)).unwrap(); + let second = budget.try_acquire(RowGroupBudget::Bytes(10)).unwrap(); + assert!(budget.try_acquire(RowGroupBudget::Bytes(10)).is_none()); + drop(first); + let third = budget.try_acquire(RowGroupBudget::Bytes(10)).unwrap(); + drop((second, third)); + + let oversized = budget.try_acquire(RowGroupBudget::Bytes(21)).unwrap(); + assert!(budget.try_acquire(RowGroupBudget::Bytes(1)).is_none()); + drop(oversized); + assert!(budget.try_acquire(RowGroupBudget::Bytes(1)).is_some()); + } + + #[tokio::test(flavor = "multi_thread", worker_threads = 4)] + async fn test_prefetch_budget_is_shared_across_files() { + let data = multi_row_group_mosaic(Vec::new()); + let fields = vec![data_fields()[0].clone()]; + let budget = Arc::new( + ParquetReadBudget::new_with_mosaic_parallelism(8, 256 * 1024 * 1024, 2).unwrap(), + ); + + let read_file = |data: Bytes| { + let fields = fields.clone(); + let budget = Arc::clone(&budget); + async move { + MosaicFormatReader::with_read_budget(budget) + .read_batch_stream( + Box::new(ConcurrentTrackingFileRead { + data: data.clone(), + active: Arc::new(AtomicUsize::new(0)), + max_active: Arc::new(AtomicUsize::new(0)), + }), + data.len() as u64, + &fields, + None, + None, + None, + ) + .await? + .try_collect::>() + .await + } + }; + + let (first, second) = tokio::join!(read_file(data.clone()), read_file(data)); + let first = first.unwrap(); + let second = second.unwrap(); + assert_eq!(collect_i32_column(&first, 0), vec![1, 2, 10, 11, 20, 21]); + assert_eq!(collect_i32_column(&second, 0), vec![1, 2, 10, 11, 20, 21]); + assert_eq!( + budget.mosaic_peak_inflight(), + 2, + "concurrent files must share the scan-level Mosaic task cap" + ); + } + #[tokio::test] async fn test_row_id_predicate_is_rejected() { let data = write_mosaic(&sample_batch()); @@ -1351,7 +1732,7 @@ mod tests { let calls = Arc::new(Mutex::new(Vec::new())); assert!(file_size > 64 * 1024); - let batches = MosaicFormatReader + let batches = MosaicFormatReader::default() .read_batch_stream( Box::new(TrackingFileRead { data, @@ -1407,7 +1788,7 @@ mod tests { .unwrap(); let batches = runtime.block_on(async move { - MosaicFormatReader + MosaicFormatReader::default() .read_batch_stream( Box::new(RuntimeThreadFileRead { data, diff --git a/crates/paimon/src/arrow/parquet_read_budget.rs b/crates/paimon/src/arrow/parquet_read_budget.rs index bef9c6ebd..e06387791 100644 --- a/crates/paimon/src/arrow/parquet_read_budget.rs +++ b/crates/paimon/src/arrow/parquet_read_budget.rs @@ -18,17 +18,24 @@ use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}; use std::sync::Arc; -use tokio::sync::{OwnedSemaphorePermit, Semaphore}; +use tokio::sync::{OwnedSemaphorePermit, Semaphore, TryAcquireError}; const BYTE_PERMIT_UNIT: u64 = 1024 * 1024; const DEFAULT_PARALLELISM: usize = 8; const DEFAULT_MAX_INFLIGHT_BYTES: u64 = 256 * 1024 * 1024; -/// Shared resource budget for concurrent Parquet row-group reads. +/// Scan-shared resource budgets for concurrent Parquet and Mosaic row-group reads. +/// +/// The public name is retained for API compatibility with callers that inject a +/// Parquet budget; the same scan object now also carries Mosaic's concurrency +/// semaphore so all files in that scan share one cap. #[derive(Debug)] pub struct ParquetReadBudget { parallelism: usize, row_groups: Arc, + mosaic_row_groups: Arc, + #[cfg(test)] + mosaic_diagnostics: Arc, bytes: Arc, byte_permits: u32, max_inflight_bytes: u64, @@ -47,6 +54,13 @@ struct ParquetReadDiagnostics { peak_inflight: AtomicUsize, } +#[cfg(test)] +#[derive(Debug, Default)] +struct MosaicReadDiagnostics { + current_inflight: AtomicUsize, + peak_inflight: AtomicUsize, +} + impl Default for ParquetReadDiagnostics { fn default() -> Self { Self { @@ -73,6 +87,15 @@ pub(crate) struct ParquetReadDiagnosticsSnapshot { impl ParquetReadBudget { pub fn new(parallelism: usize, max_inflight_bytes: u64) -> crate::Result { + Self::new_with_mosaic_parallelism(parallelism, max_inflight_bytes, 8) + } + + /// Create scan-shared row-group budgets for Parquet and Mosaic readers. + pub fn new_with_mosaic_parallelism( + parallelism: usize, + max_inflight_bytes: u64, + mosaic_parallelism: usize, + ) -> crate::Result { if parallelism == 0 || parallelism > Semaphore::MAX_PERMITS { return Err(crate::Error::DataInvalid { message: format!( @@ -88,6 +111,15 @@ impl ParquetReadBudget { source: None, }); } + if mosaic_parallelism == 0 || mosaic_parallelism > Semaphore::MAX_PERMITS { + return Err(crate::Error::DataInvalid { + message: format!( + "Mosaic row-group parallelism must be between 1 and {}, got {mosaic_parallelism}", + Semaphore::MAX_PERMITS + ), + source: None, + }); + } let max_byte_permits = Semaphore::MAX_PERMITS.min(u32::MAX as usize) as u32; let byte_permits = max_inflight_bytes .div_ceil(BYTE_PERMIT_UNIT) @@ -96,6 +128,9 @@ impl ParquetReadBudget { Ok(Self { parallelism, row_groups: Arc::new(Semaphore::new(parallelism)), + mosaic_row_groups: Arc::new(Semaphore::new(mosaic_parallelism)), + #[cfg(test)] + mosaic_diagnostics: Arc::new(MosaicReadDiagnostics::default()), bytes: Arc::new(Semaphore::new(byte_permits as usize)), byte_permits, max_inflight_bytes, @@ -209,6 +244,53 @@ impl ParquetReadBudget { diagnostics, }) } + + pub(crate) async fn acquire_mosaic(&self) -> crate::Result { + let row_group = Arc::clone(&self.mosaic_row_groups) + .acquire_owned() + .await + .map_err(|error| crate::Error::UnexpectedError { + message: "Mosaic row-group read budget was closed".to_string(), + source: Some(Box::new(error)), + })?; + Ok(self.mosaic_permit(row_group)) + } + + pub(crate) fn try_acquire_mosaic(&self) -> crate::Result> { + match Arc::clone(&self.mosaic_row_groups).try_acquire_owned() { + Ok(row_group) => Ok(Some(self.mosaic_permit(row_group))), + Err(TryAcquireError::NoPermits) => Ok(None), + Err(error) => Err(crate::Error::UnexpectedError { + message: "Mosaic row-group read budget was closed".to_string(), + source: Some(Box::new(error)), + }), + } + } + + fn mosaic_permit(&self, row_group: OwnedSemaphorePermit) -> MosaicReadPermit { + #[cfg(test)] + let diagnostics = { + let current = self + .mosaic_diagnostics + .current_inflight + .fetch_add(1, Ordering::SeqCst) + + 1; + self.mosaic_diagnostics + .peak_inflight + .fetch_max(current, Ordering::SeqCst); + Arc::clone(&self.mosaic_diagnostics) + }; + MosaicReadPermit { + _row_group: row_group, + #[cfg(test)] + diagnostics, + } + } + + #[cfg(test)] + pub(crate) fn mosaic_peak_inflight(&self) -> usize { + self.mosaic_diagnostics.peak_inflight.load(Ordering::SeqCst) + } } impl Default for ParquetReadBudget { @@ -225,6 +307,22 @@ pub(crate) struct ParquetReadPermit { diagnostics: Option>, } +#[derive(Debug)] +pub(crate) struct MosaicReadPermit { + _row_group: OwnedSemaphorePermit, + #[cfg(test)] + diagnostics: Arc, +} + +#[cfg(test)] +impl Drop for MosaicReadPermit { + fn drop(&mut self) { + self.diagnostics + .current_inflight + .fetch_sub(1, Ordering::SeqCst); + } +} + impl Drop for ParquetReadPermit { fn drop(&mut self) { if let Some(diagnostics) = &self.diagnostics { diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index d3114fd4c..7f19c0ab4 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -93,6 +93,7 @@ const READ_BATCH_SIZE_OPTION: &str = "read.batch-size"; const PARQUET_ROW_GROUP_PARALLELISM_OPTION: &str = "read.parquet.row-group.parallelism"; const PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION: &str = "read.parquet.row-group.max-inflight-bytes"; +const MOSAIC_ROW_GROUP_PARALLELISM_OPTION: &str = "read.mosaic.row-group.parallelism"; pub(crate) const TABLE_READ_SEQUENCE_NUMBER_ENABLED_OPTION: &str = "table-read.sequence-number.enabled"; pub(crate) const SEQUENCE_FIELD_OPTION: &str = "sequence.field"; @@ -140,6 +141,7 @@ const DEFAULT_WRITE_PARQUET_BUFFER_SIZE: i64 = 256 * 1024 * 1024; const DEFAULT_READ_BATCH_SIZE: usize = 1024; const DEFAULT_PARQUET_ROW_GROUP_PARALLELISM: usize = 8; const DEFAULT_PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES: i64 = 256 * 1024 * 1024; +const DEFAULT_MOSAIC_ROW_GROUP_PARALLELISM: usize = 8; const DYNAMIC_BUCKET_TARGET_ROW_NUM_OPTION: &str = "dynamic-bucket.target-row-num"; const DEFAULT_DYNAMIC_BUCKET_TARGET_ROW_NUM: i64 = 200_000; const DEFAULT_GLOBAL_INDEX_ROW_COUNT_PER_SHARD: i64 = 100_000; @@ -402,6 +404,30 @@ impl<'a> CoreOptions<'a> { Ok(value) } + /// Maximum concurrent Mosaic row-group reads shared by every file in a scan. + pub fn mosaic_row_group_parallelism(&self) -> crate::Result { + let Some(raw) = self.options.get(MOSAIC_ROW_GROUP_PARALLELISM_OPTION) else { + return Ok(DEFAULT_MOSAIC_ROW_GROUP_PARALLELISM); + }; + let value = raw + .parse::() + .map_err(|error| crate::Error::DataInvalid { + message: format!( + "Option '{MOSAIC_ROW_GROUP_PARALLELISM_OPTION}' must be a positive integer, got: {raw}" + ), + source: Some(Box::new(error)), + })?; + if value == 0 { + return Err(crate::Error::DataInvalid { + message: format!( + "Option '{MOSAIC_ROW_GROUP_PARALLELISM_OPTION}' must be greater than 0" + ), + source: None, + }); + } + Ok(value) + } + /// Scan-wide projected uncompressed bytes for concurrent Parquet row groups. pub fn parquet_row_group_max_inflight_bytes(&self) -> crate::Result { let value = match self @@ -1678,6 +1704,38 @@ mod tests { } } + #[test] + fn test_mosaic_row_group_parallelism_option() { + let options = HashMap::new(); + assert_eq!( + CoreOptions::new(&options) + .mosaic_row_group_parallelism() + .unwrap(), + 8 + ); + + let options = HashMap::from([( + MOSAIC_ROW_GROUP_PARALLELISM_OPTION.to_string(), + "3".to_string(), + )]); + assert_eq!( + CoreOptions::new(&options) + .mosaic_row_group_parallelism() + .unwrap(), + 3 + ); + + for value in ["0", "-1", "invalid"] { + let options = HashMap::from([( + MOSAIC_ROW_GROUP_PARALLELISM_OPTION.to_string(), + value.to_string(), + )]); + assert!(CoreOptions::new(&options) + .mosaic_row_group_parallelism() + .is_err()); + } + } + #[test] fn test_source_split_defaults() { let options = HashMap::new(); diff --git a/crates/paimon/src/table/table_read.rs b/crates/paimon/src/table/table_read.rs index b684d4111..dac5bd33b 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -63,9 +63,10 @@ pub(super) fn configured_parquet_read_budget( table: &Table, ) -> crate::Result> { let options = table.schema().core_options(); - Ok(Arc::new(ParquetReadBudget::new( + Ok(Arc::new(ParquetReadBudget::new_with_mosaic_parallelism( options.parquet_row_group_parallelism()?, options.parquet_row_group_max_inflight_bytes()?, + options.mosaic_row_group_parallelism()?, )?)) } From 68ad0c73623c98770f0a5a306b41ed84c84bbd87 Mon Sep 17 00:00:00 2001 From: mingfeng Date: Thu, 17 Sep 2026 01:48:50 -0700 Subject: [PATCH 3/4] refactor(mosaic): own row-group prefetch budget --- crates/paimon/src/arrow/format/mod.rs | 11 +- crates/paimon/src/arrow/format/mosaic.rs | 322 ++++++------------ .../paimon/src/arrow/parquet_read_budget.rs | 102 +----- crates/paimon/src/spec/core_options.rs | 58 ---- crates/paimon/src/table/table_read.rs | 3 +- 5 files changed, 108 insertions(+), 388 deletions(-) diff --git a/crates/paimon/src/arrow/format/mod.rs b/crates/paimon/src/arrow/format/mod.rs index ceec05171..f67ae5651 100644 --- a/crates/paimon/src/arrow/format/mod.rs +++ b/crates/paimon/src/arrow/format/mod.rs @@ -193,7 +193,7 @@ pub(crate) fn create_format_reader( create_format_reader_with_budget(path, blob_as_descriptor, read_fields, None) } -/// Create a format reader with scan-shared row-group resource budgets. +/// Create a format reader with a scan-shared Parquet resource budget. pub(crate) fn create_format_reader_with_budget( path: &str, blob_as_descriptor: bool, @@ -219,11 +219,10 @@ pub(crate) fn create_format_reader_with_budget( Box::new(row::RowFormatReader) } else { if lower.ends_with(".mosaic") { - let reader: Box = Box::new(match parquet_read_budget { - Some(read_budget) => mosaic::MosaicFormatReader::with_read_budget(read_budget), - None => mosaic::MosaicFormatReader::default(), - }); - return Ok(shredding::maybe_wrap_reader(reader, read_fields)); + return Ok(shredding::maybe_wrap_reader( + Box::new(mosaic::MosaicFormatReader), + read_fields, + )); } #[cfg(feature = "vortex")] if lower.ends_with(".vortex") { diff --git a/crates/paimon/src/arrow/format/mosaic.rs b/crates/paimon/src/arrow/format/mosaic.rs index 33912bb68..e0a124ad4 100644 --- a/crates/paimon/src/arrow/format/mosaic.rs +++ b/crates/paimon/src/arrow/format/mosaic.rs @@ -20,9 +20,7 @@ use crate::arrow::build_target_arrow_schema; use crate::arrow::filtering::{ predicates_may_match_with_schema, remap_predicates_to_file, StatsAccessor, }; -use crate::arrow::parquet_read_budget::MosaicReadPermit; use crate::arrow::residual::{filter_record_batch_by_predicates, widen_scan_fields}; -use crate::arrow::ParquetReadBudget; use crate::io::FileRead; use crate::spec::{DataField, DataType as PaimonDataType, Datum, Predicate}; use crate::table::{ArrowRecordBatchStream, RowRange}; @@ -45,27 +43,13 @@ use std::ops::Range; use std::panic::{catch_unwind, AssertUnwindSafe}; use std::sync::{Arc, Condvar, Mutex, OnceLock}; -pub(crate) struct MosaicFormatReader { - read_budget: Arc, -} +pub(crate) struct MosaicFormatReader; const DEFAULT_BATCH_SIZE: usize = 8192; const DEFAULT_PREFETCH_ROW_GROUPS: usize = 8; const DEFAULT_PREFETCH_MAX_BYTES: usize = 64 * 1024 * 1024; const MAX_EXECUTOR_WORKERS: usize = 32; -impl MosaicFormatReader { - pub(crate) fn with_read_budget(read_budget: Arc) -> Self { - Self { read_budget } - } -} - -impl Default for MosaicFormatReader { - fn default() -> Self { - Self::with_read_budget(Arc::new(ParquetReadBudget::default())) - } -} - #[async_trait] impl FormatFileReader for MosaicFormatReader { async fn read_batch_stream( @@ -93,8 +77,6 @@ impl FormatFileReader for MosaicFormatReader { file_fields: predicates.file_fields.clone(), }); let batch_size = batch_size.unwrap_or(DEFAULT_BATCH_SIZE); - let read_budget = Arc::clone(&self.read_budget); - let (mut batch_tx, mut batch_rx) = futures::channel::mpsc::channel(1); let read_task = tokio::task::spawn_blocking(move || { let result = read_mosaic_batches_blocking( @@ -108,7 +90,6 @@ impl FormatFileReader for MosaicFormatReader { handle, prefetch_row_groups: DEFAULT_PREFETCH_ROW_GROUPS, prefetch_max_bytes: DEFAULT_PREFETCH_MAX_BYTES, - read_budget, }, |batch| futures::executor::block_on(batch_tx.send(Ok(batch))).is_ok(), ); @@ -121,9 +102,8 @@ impl FormatFileReader for MosaicFormatReader { while let Some(batch) = batch_rx.next().await { let BudgetedBatch { batch, permit } = batch?; yield batch; - // Keep both the byte and row-group permits while the decoded - // batch is queued or yielded. Advancing or dropping the - // stream releases them. + // Keep the Mosaic prefetch permit while the decoded batch is + // queued or yielded. Advancing or dropping the stream releases it. drop(permit); } read_task.await.map_err(|e| Error::DataInvalid { @@ -145,14 +125,13 @@ struct MosaicReadRequest { handle: tokio::runtime::Handle, prefetch_row_groups: usize, prefetch_max_bytes: usize, - read_budget: Arc, } struct PlannedRowGroup { index: usize, rows: usize, selected_slices: Option>, - budget: RowGroupBudget, + estimated_bytes: usize, } struct PendingRowGroup { @@ -165,110 +144,93 @@ struct BudgetedBatch { permit: Option>, } -#[derive(Clone, Copy, Debug, PartialEq, Eq)] -enum RowGroupBudget { - Bytes(usize), - /// Mosaic 0.2 does not expose projected physical sizes. Variable-width - /// and nested projections therefore run one row group at a time. - Exclusive, -} - #[derive(Default)] -struct ByteBudgetState { +struct MosaicPrefetchBudgetState { bytes: usize, groups: usize, - exclusive: bool, } -struct MosaicByteBudgetInner { +struct MosaicPrefetchBudgetInner { max_bytes: usize, max_groups: usize, - state: Mutex, + state: Mutex, released: Condvar, } #[derive(Clone)] -struct MosaicByteBudget { - inner: Arc, +struct MosaicPrefetchBudget { + inner: Arc, } -impl MosaicByteBudget { +impl MosaicPrefetchBudget { fn new(max_bytes: usize, max_groups: usize) -> Self { Self { - inner: Arc::new(MosaicByteBudgetInner { + inner: Arc::new(MosaicPrefetchBudgetInner { max_bytes, max_groups, - state: Mutex::new(ByteBudgetState::default()), + state: Mutex::new(MosaicPrefetchBudgetState::default()), released: Condvar::new(), }), } } - fn try_acquire(&self, requirement: RowGroupBudget) -> Option { + fn try_acquire(&self, estimated_bytes: usize) -> Option { let mut state = self .inner .state .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); - if !self.can_acquire(&state, requirement) { + if !self.can_acquire(&state, estimated_bytes) { return None; } - Self::charge(&mut state, requirement); - Some(MosaicBytePermit { + Self::charge(&mut state, estimated_bytes); + Some(MosaicPrefetchPermit { inner: Arc::clone(&self.inner), - requirement, + estimated_bytes, }) } - fn acquire(&self, requirement: RowGroupBudget) -> MosaicBytePermit { + fn acquire(&self, estimated_bytes: usize) -> MosaicPrefetchPermit { let mut state = self .inner .state .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); - while !self.can_acquire(&state, requirement) { + while !self.can_acquire(&state, estimated_bytes) { state = self .inner .released .wait(state) .unwrap_or_else(|poisoned| poisoned.into_inner()); } - Self::charge(&mut state, requirement); - MosaicBytePermit { + Self::charge(&mut state, estimated_bytes); + MosaicPrefetchPermit { inner: Arc::clone(&self.inner), - requirement, + estimated_bytes, } } - fn can_acquire(&self, state: &ByteBudgetState, requirement: RowGroupBudget) -> bool { - if state.groups >= self.inner.max_groups || state.exclusive { + fn can_acquire(&self, state: &MosaicPrefetchBudgetState, estimated_bytes: usize) -> bool { + if state.groups >= self.inner.max_groups { return false; } - match requirement { - RowGroupBudget::Exclusive => state.groups == 0, - RowGroupBudget::Bytes(bytes) => { - // An oversized head group may consume the whole budget so the - // reader always makes progress, but it cannot overlap peers. - state.groups == 0 || state.bytes.saturating_add(bytes) <= self.inner.max_bytes - } - } + // An oversized head group may consume the whole budget so the reader + // always makes progress, but it cannot overlap peers. + state.groups == 0 || state.bytes.saturating_add(estimated_bytes) <= self.inner.max_bytes } - fn charge(state: &mut ByteBudgetState, requirement: RowGroupBudget) { + fn charge(state: &mut MosaicPrefetchBudgetState, estimated_bytes: usize) { state.groups += 1; - match requirement { - RowGroupBudget::Bytes(bytes) => state.bytes = state.bytes.saturating_add(bytes), - RowGroupBudget::Exclusive => state.exclusive = true, - } + state.bytes = state.bytes.saturating_add(estimated_bytes); } } -struct MosaicBytePermit { - inner: Arc, - requirement: RowGroupBudget, +struct MosaicPrefetchPermit { + inner: Arc, + estimated_bytes: usize, } -impl Drop for MosaicBytePermit { +impl Drop for MosaicPrefetchPermit { fn drop(&mut self) { let mut state = self .inner @@ -276,19 +238,11 @@ impl Drop for MosaicBytePermit { .lock() .unwrap_or_else(|poisoned| poisoned.into_inner()); state.groups -= 1; - match self.requirement { - RowGroupBudget::Bytes(bytes) => state.bytes = state.bytes.saturating_sub(bytes), - RowGroupBudget::Exclusive => state.exclusive = false, - } + state.bytes = state.bytes.saturating_sub(self.estimated_bytes); self.inner.released.notify_all(); } } -struct MosaicPrefetchPermit { - _bytes: MosaicBytePermit, - _row_group: MosaicReadPermit, -} - type MosaicJob = Box; struct MosaicExecutor { @@ -388,7 +342,6 @@ fn read_mosaic_batches_blocking( handle, prefetch_row_groups, prefetch_max_bytes, - read_budget, } = request; let mosaic_reader = Arc::new( MosaicReader::new(FileReadInputFile::new(reader, handle.clone()), file_size) @@ -479,7 +432,8 @@ fn read_mosaic_batches_blocking( index: row_group_index, rows: row_group_rows, selected_slices, - budget: projected_row_group_budget(&existing_scan_fields, row_group_rows), + estimated_bytes: row_group_rows + .saturating_mul(estimated_row_bytes(&existing_scan_fields)), }); } @@ -505,30 +459,15 @@ fn read_mosaic_batches_blocking( let projected_names = Arc::new(projected_names); let executor = mosaic_executor()?; let mut pending = VecDeque::::new(); - let byte_budget = MosaicByteBudget::new(prefetch_max_bytes, prefetch_row_groups); + let prefetch_budget = MosaicPrefetchBudget::new(prefetch_max_bytes, prefetch_row_groups); let mut outcome = Ok(()); 'read: while !planned.is_empty() || !pending.is_empty() { while let Some(next) = planned.front() { - let Some(bytes) = byte_budget.try_acquire(next.budget) else { + let Some(permit) = prefetch_budget.try_acquire(next.estimated_bytes) else { break; }; - let row_group = match read_budget.try_acquire_mosaic() { - Ok(Some(permit)) => permit, - Ok(None) => { - drop(bytes); - break; - } - Err(error) => { - outcome = Err(error); - break 'read; - } - }; let plan = planned.pop_front().unwrap(); - let permit = MosaicPrefetchPermit { - _bytes: bytes, - _row_group: row_group, - }; let reader = Arc::clone(&mosaic_reader); let names = Arc::clone(&projected_names); let schema = Arc::clone(&read_schema); @@ -550,18 +489,7 @@ fn read_mosaic_batches_blocking( }; // Earlier decoded batches can still own the budget downstream. // With no task left to join, wait for their permits to be released. - let bytes = byte_budget.acquire(next.budget); - let row_group = match handle.block_on(read_budget.acquire_mosaic()) { - Ok(permit) => permit, - Err(error) => { - outcome = Err(error); - break 'read; - } - }; - let permit = MosaicPrefetchPermit { - _bytes: bytes, - _row_group: row_group, - }; + let permit = prefetch_budget.acquire(next.estimated_bytes); let plan = planned.pop_front().unwrap(); let reader = Arc::clone(&mosaic_reader); let names = Arc::clone(&projected_names); @@ -666,33 +594,30 @@ fn send_mosaic_batch_chunks( }) } -fn projected_row_group_budget(fields: &[DataField], rows: usize) -> RowGroupBudget { - let Some(row_bytes) = fields.iter().try_fold(0usize, |total, field| { - fixed_width_row_bytes(field.data_type()).map(|bytes| total.saturating_add(bytes)) - }) else { - return RowGroupBudget::Exclusive; - }; - RowGroupBudget::Bytes(rows.saturating_mul(row_bytes.max(1))) -} - -fn fixed_width_row_bytes(data_type: &PaimonDataType) -> Option { - // One byte per row for validity is deliberately conservative (Arrow uses - // one bit). These are upper bounds for decoded fixed-width array buffers, - // unlike a guessed average for variable-width values. - match data_type { - PaimonDataType::Boolean(_) | PaimonDataType::TinyInt(_) => Some(2), - PaimonDataType::SmallInt(_) => Some(3), - PaimonDataType::Int(_) - | PaimonDataType::Float(_) - | PaimonDataType::Date(_) - | PaimonDataType::Time(_) => Some(5), - PaimonDataType::BigInt(_) - | PaimonDataType::Double(_) - | PaimonDataType::Timestamp(_) - | PaimonDataType::LocalZonedTimestamp(_) => Some(9), - PaimonDataType::Decimal(_) => Some(17), - _ => None, - } +/// Rough decoded size of one projected row, matching Java Mosaic's prefetch budget. +fn estimated_row_bytes(fields: &[DataField]) -> usize { + fields + .iter() + .map(|field| match field.data_type() { + PaimonDataType::Boolean(_) | PaimonDataType::TinyInt(_) => 2, + PaimonDataType::SmallInt(_) => 3, + PaimonDataType::Int(_) + | PaimonDataType::Float(_) + | PaimonDataType::Date(_) + | PaimonDataType::Time(_) => 5, + PaimonDataType::BigInt(_) + | PaimonDataType::Double(_) + | PaimonDataType::Timestamp(_) + | PaimonDataType::LocalZonedTimestamp(_) => 9, + PaimonDataType::Decimal(_) => 17, + PaimonDataType::Char(_) + | PaimonDataType::VarChar(_) + | PaimonDataType::Binary(_) + | PaimonDataType::VarBinary(_) => 40, + _ => 128, + }) + .sum::() + .max(1) } struct MosaicRowGroupStats<'a> { @@ -1337,7 +1262,7 @@ mod tests { row_selection: Option>, ) -> crate::Result> { let file_size = data.len() as u64; - MosaicFormatReader::default() + MosaicFormatReader .read_batch_stream( Box::new(TestFileRead { data }), file_size, @@ -1359,7 +1284,7 @@ mod tests { ) -> crate::Result>> { let file_size = data.len() as u64; let calls = Arc::new(Mutex::new(Vec::new())); - let _: Vec = MosaicFormatReader::default() + let _: Vec = MosaicFormatReader .read_batch_stream( Box::new(TrackingFileRead { data, @@ -1385,7 +1310,7 @@ mod tests { ) -> crate::Result>> { let file_size = data.len() as u64; let calls = Arc::new(Mutex::new(Vec::new())); - let _: Vec = MosaicFormatReader::default() + let _: Vec = MosaicFormatReader .read_batch_stream( Box::new(TrackingFileRead { data, @@ -1456,14 +1381,6 @@ mod tests { max_active: Arc::clone(&max_active), }; let handle = tokio::runtime::Handle::current(); - let read_budget = Arc::new( - ParquetReadBudget::new_with_mosaic_parallelism( - 8, - 256 * 1024 * 1024, - prefetch_row_groups.max(1), - ) - .unwrap(), - ); let batches = tokio::task::spawn_blocking(move || { let mut batches = Vec::new(); read_mosaic_batches_blocking( @@ -1477,7 +1394,6 @@ mod tests { handle, prefetch_row_groups, prefetch_max_bytes, - read_budget, }, |BudgetedBatch { batch, .. }| { batches.push(batch); @@ -1590,96 +1506,58 @@ mod tests { } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] - async fn test_prefetch_byte_budget_is_cumulative() { + async fn test_disabled_prefetch_reads_row_groups_on_demand() { let data = multi_row_group_mosaic(Vec::new()); let fields = vec![data_fields()[0].clone()]; - // Each two-row Int group is conservatively charged 10 bytes, so this - // admits two groups but not all three. - let (batches, max_active) = read_with_prefetch_tracking(data, fields, 3, 20).await; + let (batches, max_active) = read_with_prefetch_tracking(data, fields, 0, usize::MAX).await; assert_eq!(collect_i32_column(&batches, 0), vec![1, 2, 10, 11, 20, 21]); - assert_eq!(max_active, 2); + assert_eq!(max_active, 1); } #[tokio::test(flavor = "multi_thread", worker_threads = 4)] - async fn test_variable_width_prefetch_is_exclusive() { + async fn test_prefetch_byte_budget_is_cumulative() { let data = multi_row_group_mosaic(Vec::new()); - let fields = vec![data_fields()[1].clone()]; - let (batches, max_active) = read_with_prefetch_tracking(data, fields, 3, usize::MAX).await; + let fields = vec![data_fields()[0].clone()]; + // Each two-row Int group is conservatively charged 10 bytes, so this + // admits two groups but not all three. + let (batches, max_active) = read_with_prefetch_tracking(data, fields, 3, 20).await; - assert_eq!(batches.iter().map(RecordBatch::num_rows).sum::(), 6); - assert_eq!( - max_active, 1, - "variable-width row groups must not use a fixed per-row guess" - ); + assert_eq!(collect_i32_column(&batches, 0), vec![1, 2, 10, 11, 20, 21]); + assert_eq!(max_active, 2); } #[test] - fn test_fixed_width_budget_is_cumulative_and_oversized_head_progresses() { + fn test_mosaic_prefetch_budget_is_cumulative_and_oversized_head_progresses() { + assert_eq!(estimated_row_bytes(&[data_fields()[0].clone()]), 5); + assert_eq!(estimated_row_bytes(&[data_fields()[1].clone()]), 40); assert_eq!( - projected_row_group_budget(&[data_fields()[0].clone()], 2), - RowGroupBudget::Bytes(10) - ); - assert_eq!( - projected_row_group_budget(&[data_fields()[1].clone()], 2), - RowGroupBudget::Exclusive + estimated_row_bytes(&[field( + 3, + "nested", + DataType::Array(ArrayType::new(DataType::Int(IntType::new()))) + )]), + 128 ); - let budget = MosaicByteBudget::new(20, 8); - let first = budget.try_acquire(RowGroupBudget::Bytes(10)).unwrap(); - let second = budget.try_acquire(RowGroupBudget::Bytes(10)).unwrap(); - assert!(budget.try_acquire(RowGroupBudget::Bytes(10)).is_none()); + let budget = MosaicPrefetchBudget::new(20, 8); + let first = budget.try_acquire(10).unwrap(); + let second = budget.try_acquire(10).unwrap(); + assert!(budget.try_acquire(10).is_none()); drop(first); - let third = budget.try_acquire(RowGroupBudget::Bytes(10)).unwrap(); + let third = budget.try_acquire(10).unwrap(); drop((second, third)); - let oversized = budget.try_acquire(RowGroupBudget::Bytes(21)).unwrap(); - assert!(budget.try_acquire(RowGroupBudget::Bytes(1)).is_none()); + let oversized = budget.try_acquire(21).unwrap(); + assert!(budget.try_acquire(1).is_none()); drop(oversized); - assert!(budget.try_acquire(RowGroupBudget::Bytes(1)).is_some()); - } - - #[tokio::test(flavor = "multi_thread", worker_threads = 4)] - async fn test_prefetch_budget_is_shared_across_files() { - let data = multi_row_group_mosaic(Vec::new()); - let fields = vec![data_fields()[0].clone()]; - let budget = Arc::new( - ParquetReadBudget::new_with_mosaic_parallelism(8, 256 * 1024 * 1024, 2).unwrap(), - ); + assert!(budget.try_acquire(1).is_some()); - let read_file = |data: Bytes| { - let fields = fields.clone(); - let budget = Arc::clone(&budget); - async move { - MosaicFormatReader::with_read_budget(budget) - .read_batch_stream( - Box::new(ConcurrentTrackingFileRead { - data: data.clone(), - active: Arc::new(AtomicUsize::new(0)), - max_active: Arc::new(AtomicUsize::new(0)), - }), - data.len() as u64, - &fields, - None, - None, - None, - ) - .await? - .try_collect::>() - .await - } - }; - - let (first, second) = tokio::join!(read_file(data.clone()), read_file(data)); - let first = first.unwrap(); - let second = second.unwrap(); - assert_eq!(collect_i32_column(&first, 0), vec![1, 2, 10, 11, 20, 21]); - assert_eq!(collect_i32_column(&second, 0), vec![1, 2, 10, 11, 20, 21]); - assert_eq!( - budget.mosaic_peak_inflight(), - 2, - "concurrent files must share the scan-level Mosaic task cap" - ); + let group_limited = MosaicPrefetchBudget::new(usize::MAX, 2); + let first = group_limited.try_acquire(1).unwrap(); + let second = group_limited.try_acquire(1).unwrap(); + assert!(group_limited.try_acquire(1).is_none()); + drop((first, second)); } #[tokio::test] @@ -1732,7 +1610,7 @@ mod tests { let calls = Arc::new(Mutex::new(Vec::new())); assert!(file_size > 64 * 1024); - let batches = MosaicFormatReader::default() + let batches = MosaicFormatReader .read_batch_stream( Box::new(TrackingFileRead { data, @@ -1788,7 +1666,7 @@ mod tests { .unwrap(); let batches = runtime.block_on(async move { - MosaicFormatReader::default() + MosaicFormatReader .read_batch_stream( Box::new(RuntimeThreadFileRead { data, diff --git a/crates/paimon/src/arrow/parquet_read_budget.rs b/crates/paimon/src/arrow/parquet_read_budget.rs index e06387791..bef9c6ebd 100644 --- a/crates/paimon/src/arrow/parquet_read_budget.rs +++ b/crates/paimon/src/arrow/parquet_read_budget.rs @@ -18,24 +18,17 @@ use std::sync::atomic::{AtomicBool, AtomicU64, AtomicUsize, Ordering}; use std::sync::Arc; -use tokio::sync::{OwnedSemaphorePermit, Semaphore, TryAcquireError}; +use tokio::sync::{OwnedSemaphorePermit, Semaphore}; const BYTE_PERMIT_UNIT: u64 = 1024 * 1024; const DEFAULT_PARALLELISM: usize = 8; const DEFAULT_MAX_INFLIGHT_BYTES: u64 = 256 * 1024 * 1024; -/// Scan-shared resource budgets for concurrent Parquet and Mosaic row-group reads. -/// -/// The public name is retained for API compatibility with callers that inject a -/// Parquet budget; the same scan object now also carries Mosaic's concurrency -/// semaphore so all files in that scan share one cap. +/// Shared resource budget for concurrent Parquet row-group reads. #[derive(Debug)] pub struct ParquetReadBudget { parallelism: usize, row_groups: Arc, - mosaic_row_groups: Arc, - #[cfg(test)] - mosaic_diagnostics: Arc, bytes: Arc, byte_permits: u32, max_inflight_bytes: u64, @@ -54,13 +47,6 @@ struct ParquetReadDiagnostics { peak_inflight: AtomicUsize, } -#[cfg(test)] -#[derive(Debug, Default)] -struct MosaicReadDiagnostics { - current_inflight: AtomicUsize, - peak_inflight: AtomicUsize, -} - impl Default for ParquetReadDiagnostics { fn default() -> Self { Self { @@ -87,15 +73,6 @@ pub(crate) struct ParquetReadDiagnosticsSnapshot { impl ParquetReadBudget { pub fn new(parallelism: usize, max_inflight_bytes: u64) -> crate::Result { - Self::new_with_mosaic_parallelism(parallelism, max_inflight_bytes, 8) - } - - /// Create scan-shared row-group budgets for Parquet and Mosaic readers. - pub fn new_with_mosaic_parallelism( - parallelism: usize, - max_inflight_bytes: u64, - mosaic_parallelism: usize, - ) -> crate::Result { if parallelism == 0 || parallelism > Semaphore::MAX_PERMITS { return Err(crate::Error::DataInvalid { message: format!( @@ -111,15 +88,6 @@ impl ParquetReadBudget { source: None, }); } - if mosaic_parallelism == 0 || mosaic_parallelism > Semaphore::MAX_PERMITS { - return Err(crate::Error::DataInvalid { - message: format!( - "Mosaic row-group parallelism must be between 1 and {}, got {mosaic_parallelism}", - Semaphore::MAX_PERMITS - ), - source: None, - }); - } let max_byte_permits = Semaphore::MAX_PERMITS.min(u32::MAX as usize) as u32; let byte_permits = max_inflight_bytes .div_ceil(BYTE_PERMIT_UNIT) @@ -128,9 +96,6 @@ impl ParquetReadBudget { Ok(Self { parallelism, row_groups: Arc::new(Semaphore::new(parallelism)), - mosaic_row_groups: Arc::new(Semaphore::new(mosaic_parallelism)), - #[cfg(test)] - mosaic_diagnostics: Arc::new(MosaicReadDiagnostics::default()), bytes: Arc::new(Semaphore::new(byte_permits as usize)), byte_permits, max_inflight_bytes, @@ -244,53 +209,6 @@ impl ParquetReadBudget { diagnostics, }) } - - pub(crate) async fn acquire_mosaic(&self) -> crate::Result { - let row_group = Arc::clone(&self.mosaic_row_groups) - .acquire_owned() - .await - .map_err(|error| crate::Error::UnexpectedError { - message: "Mosaic row-group read budget was closed".to_string(), - source: Some(Box::new(error)), - })?; - Ok(self.mosaic_permit(row_group)) - } - - pub(crate) fn try_acquire_mosaic(&self) -> crate::Result> { - match Arc::clone(&self.mosaic_row_groups).try_acquire_owned() { - Ok(row_group) => Ok(Some(self.mosaic_permit(row_group))), - Err(TryAcquireError::NoPermits) => Ok(None), - Err(error) => Err(crate::Error::UnexpectedError { - message: "Mosaic row-group read budget was closed".to_string(), - source: Some(Box::new(error)), - }), - } - } - - fn mosaic_permit(&self, row_group: OwnedSemaphorePermit) -> MosaicReadPermit { - #[cfg(test)] - let diagnostics = { - let current = self - .mosaic_diagnostics - .current_inflight - .fetch_add(1, Ordering::SeqCst) - + 1; - self.mosaic_diagnostics - .peak_inflight - .fetch_max(current, Ordering::SeqCst); - Arc::clone(&self.mosaic_diagnostics) - }; - MosaicReadPermit { - _row_group: row_group, - #[cfg(test)] - diagnostics, - } - } - - #[cfg(test)] - pub(crate) fn mosaic_peak_inflight(&self) -> usize { - self.mosaic_diagnostics.peak_inflight.load(Ordering::SeqCst) - } } impl Default for ParquetReadBudget { @@ -307,22 +225,6 @@ pub(crate) struct ParquetReadPermit { diagnostics: Option>, } -#[derive(Debug)] -pub(crate) struct MosaicReadPermit { - _row_group: OwnedSemaphorePermit, - #[cfg(test)] - diagnostics: Arc, -} - -#[cfg(test)] -impl Drop for MosaicReadPermit { - fn drop(&mut self) { - self.diagnostics - .current_inflight - .fetch_sub(1, Ordering::SeqCst); - } -} - impl Drop for ParquetReadPermit { fn drop(&mut self) { if let Some(diagnostics) = &self.diagnostics { diff --git a/crates/paimon/src/spec/core_options.rs b/crates/paimon/src/spec/core_options.rs index 7f19c0ab4..d3114fd4c 100644 --- a/crates/paimon/src/spec/core_options.rs +++ b/crates/paimon/src/spec/core_options.rs @@ -93,7 +93,6 @@ const READ_BATCH_SIZE_OPTION: &str = "read.batch-size"; const PARQUET_ROW_GROUP_PARALLELISM_OPTION: &str = "read.parquet.row-group.parallelism"; const PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION: &str = "read.parquet.row-group.max-inflight-bytes"; -const MOSAIC_ROW_GROUP_PARALLELISM_OPTION: &str = "read.mosaic.row-group.parallelism"; pub(crate) const TABLE_READ_SEQUENCE_NUMBER_ENABLED_OPTION: &str = "table-read.sequence-number.enabled"; pub(crate) const SEQUENCE_FIELD_OPTION: &str = "sequence.field"; @@ -141,7 +140,6 @@ const DEFAULT_WRITE_PARQUET_BUFFER_SIZE: i64 = 256 * 1024 * 1024; const DEFAULT_READ_BATCH_SIZE: usize = 1024; const DEFAULT_PARQUET_ROW_GROUP_PARALLELISM: usize = 8; const DEFAULT_PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES: i64 = 256 * 1024 * 1024; -const DEFAULT_MOSAIC_ROW_GROUP_PARALLELISM: usize = 8; const DYNAMIC_BUCKET_TARGET_ROW_NUM_OPTION: &str = "dynamic-bucket.target-row-num"; const DEFAULT_DYNAMIC_BUCKET_TARGET_ROW_NUM: i64 = 200_000; const DEFAULT_GLOBAL_INDEX_ROW_COUNT_PER_SHARD: i64 = 100_000; @@ -404,30 +402,6 @@ impl<'a> CoreOptions<'a> { Ok(value) } - /// Maximum concurrent Mosaic row-group reads shared by every file in a scan. - pub fn mosaic_row_group_parallelism(&self) -> crate::Result { - let Some(raw) = self.options.get(MOSAIC_ROW_GROUP_PARALLELISM_OPTION) else { - return Ok(DEFAULT_MOSAIC_ROW_GROUP_PARALLELISM); - }; - let value = raw - .parse::() - .map_err(|error| crate::Error::DataInvalid { - message: format!( - "Option '{MOSAIC_ROW_GROUP_PARALLELISM_OPTION}' must be a positive integer, got: {raw}" - ), - source: Some(Box::new(error)), - })?; - if value == 0 { - return Err(crate::Error::DataInvalid { - message: format!( - "Option '{MOSAIC_ROW_GROUP_PARALLELISM_OPTION}' must be greater than 0" - ), - source: None, - }); - } - Ok(value) - } - /// Scan-wide projected uncompressed bytes for concurrent Parquet row groups. pub fn parquet_row_group_max_inflight_bytes(&self) -> crate::Result { let value = match self @@ -1704,38 +1678,6 @@ mod tests { } } - #[test] - fn test_mosaic_row_group_parallelism_option() { - let options = HashMap::new(); - assert_eq!( - CoreOptions::new(&options) - .mosaic_row_group_parallelism() - .unwrap(), - 8 - ); - - let options = HashMap::from([( - MOSAIC_ROW_GROUP_PARALLELISM_OPTION.to_string(), - "3".to_string(), - )]); - assert_eq!( - CoreOptions::new(&options) - .mosaic_row_group_parallelism() - .unwrap(), - 3 - ); - - for value in ["0", "-1", "invalid"] { - let options = HashMap::from([( - MOSAIC_ROW_GROUP_PARALLELISM_OPTION.to_string(), - value.to_string(), - )]); - assert!(CoreOptions::new(&options) - .mosaic_row_group_parallelism() - .is_err()); - } - } - #[test] fn test_source_split_defaults() { let options = HashMap::new(); diff --git a/crates/paimon/src/table/table_read.rs b/crates/paimon/src/table/table_read.rs index dac5bd33b..b684d4111 100644 --- a/crates/paimon/src/table/table_read.rs +++ b/crates/paimon/src/table/table_read.rs @@ -63,10 +63,9 @@ pub(super) fn configured_parquet_read_budget( table: &Table, ) -> crate::Result> { let options = table.schema().core_options(); - Ok(Arc::new(ParquetReadBudget::new_with_mosaic_parallelism( + Ok(Arc::new(ParquetReadBudget::new( options.parquet_row_group_parallelism()?, options.parquet_row_group_max_inflight_bytes()?, - options.mosaic_row_group_parallelism()?, )?)) } From 26bebe80a5b2a3d28e7c1365cc96fb3103a79f5a Mon Sep 17 00:00:00 2001 From: mingfeng Date: Thu, 17 Sep 2026 03:38:10 -0700 Subject: [PATCH 4/4] fix(mosaic): retry executor startup after transient failure --- crates/paimon/src/arrow/format/mosaic.rs | 64 ++++++++++++++++++++---- 1 file changed, 55 insertions(+), 9 deletions(-) diff --git a/crates/paimon/src/arrow/format/mosaic.rs b/crates/paimon/src/arrow/format/mosaic.rs index e0a124ad4..ae7c5011b 100644 --- a/crates/paimon/src/arrow/format/mosaic.rs +++ b/crates/paimon/src/arrow/format/mosaic.rs @@ -317,15 +317,28 @@ impl MosaicTask { } } +fn get_or_try_init<'a, T, E>( + cell: &'a OnceLock, + init_lock: &Mutex<()>, + initialize: impl FnOnce() -> Result, +) -> Result<&'a T, E> { + if let Some(value) = cell.get() { + return Ok(value); + } + let _guard = init_lock + .lock() + .unwrap_or_else(|poisoned| poisoned.into_inner()); + if let Some(value) = cell.get() { + return Ok(value); + } + let value = initialize()?; + Ok(cell.get_or_init(|| value)) +} + fn mosaic_executor() -> crate::Result<&'static MosaicExecutor> { - static EXECUTOR: OnceLock> = OnceLock::new(); - match EXECUTOR.get_or_init(|| MosaicExecutor::new().map_err(|error| error.to_string())) { - Ok(executor) => Ok(executor), - Err(message) => Err(Error::UnexpectedError { - message: message.clone(), - source: None, - }), - } + static EXECUTOR: OnceLock = OnceLock::new(); + static INIT_LOCK: Mutex<()> = Mutex::new(()); + get_or_try_init(&EXECUTOR, &INIT_LOCK, MosaicExecutor::new) } fn read_mosaic_batches_blocking( @@ -456,6 +469,10 @@ fn read_mosaic_batches_blocking( return Ok(()); } + if planned.is_empty() { + return Ok(()); + } + let projected_names = Arc::new(projected_names); let executor = mosaic_executor()?; let mut pending = VecDeque::::new(); @@ -1016,7 +1033,36 @@ mod tests { use paimon_mosaic_core::writer::{MosaicWriter, OutputFile, WriterOptions}; use std::ops::Range; use std::sync::atomic::{AtomicUsize, Ordering}; - use std::sync::{Arc, Mutex}; + use std::sync::{Arc, Mutex, OnceLock}; + + #[test] + fn test_mosaic_executor_init_retries_after_transient_failure() { + let executor = OnceLock::new(); + let init_lock = Mutex::new(()); + let attempts = AtomicUsize::new(0); + let initialize = || { + if attempts.fetch_add(1, Ordering::SeqCst) == 0 { + Err("transient thread creation failure") + } else { + Ok(42) + } + }; + + assert_eq!( + get_or_try_init(&executor, &init_lock, initialize), + Err("transient thread creation failure") + ); + assert!(executor.get().is_none(), "a failed init must not be cached"); + assert_eq!( + *get_or_try_init(&executor, &init_lock, initialize).unwrap(), + 42 + ); + assert_eq!( + *get_or_try_init(&executor, &init_lock, initialize).unwrap(), + 42 + ); + assert_eq!(attempts.load(Ordering::SeqCst), 2); + } struct TestFileRead { data: Bytes,