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 b10c756c0..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,19 +33,38 @@ 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; 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; +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 { @@ -72,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 || { @@ -84,6 +106,9 @@ impl FormatFileReader for MosaicFormatReader { batch_size, row_selection, 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(), ); @@ -94,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}"), @@ -113,11 +143,240 @@ struct MosaicReadRequest { batch_size: usize, row_selection: Option>, 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, +} + +struct PendingRowGroup { + 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, @@ -127,9 +386,14 @@ fn read_mosaic_batches_blocking( batch_size, row_selection, handle, + prefetch_row_groups, + prefetch_max_bytes, + read_budget, } = 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.clone()), file_size) + .map_err(mosaic_read_error)?, + ); let file_column_names = mosaic_reader .schema() @@ -169,6 +433,7 @@ fn read_mosaic_batches_blocking( build_file_column_indices(mosaic_reader.schema(), &predicates.file_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 +475,224 @@ 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, + budget: projected_row_group_budget(&existing_scan_fields, row_group_rows), + }); + } + + 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, None, &mut send_batch) { + return Ok(()); + } + } + return Ok(()); + } + + 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 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 { + 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); + 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 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 }); + } - let batch = row_group_reader.read_columns().map_err(mosaic_read_error)?; - take_row_slices(batch, selected_slices.as_deref(), &read_schema)? + let PendingRowGroup { permit, task } = pending.pop_front().unwrap(); + let batch = match task.join() { + Ok(batch) => batch, + Err(error) => { + outcome = Err(error); + break; + } }; - let batch = match predicates.as_ref() { - Some(predicates) => { - filter_record_batch_by_predicates(batch, predicates, &existing_scan_fields)? + let batch = match apply_mosaic_residual(batch, predicates.as_ref(), &existing_scan_fields) { + 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(()); + if !send_mosaic_batch_chunks(batch, batch_size, Some(Arc::new(permit)), &mut send_batch) { + planned.clear(); + break; + } + } + + // 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) = next.task.join() { + if outcome.is_ok() { + outcome = Err(error); } } } - Ok(()) + 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 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, + permit: Option>, + send_batch: &mut impl FnMut(BudgetedBatch) -> bool, +) -> bool { + split_batch(batch, batch_size).into_iter().all(|batch| { + send_batch(BudgetedBatch { + batch, + permit: permit.clone(), + }) + }) +} + +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> { @@ -637,6 +1090,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 +1126,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 { @@ -863,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, @@ -885,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, @@ -911,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, @@ -967,6 +1441,57 @@ 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 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( + 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, + read_budget, + }, + |BudgetedBatch { 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 +1563,125 @@ 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 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!( + 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 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!( + max_active, 1, + "a one-byte budget should admit only the mandatory head group" + ); + } + + #[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()); @@ -1088,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, @@ -1144,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()?, )?)) }