This is an automated email from the ASF dual-hosted git repository.
JingsongLi pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/paimon-rust.git
The following commit(s) were added to refs/heads/main by this push:
new c0694642 perf(parquet): read row groups concurrently (#633)
c0694642 is described below
commit c0694642ae62769fcce6660a0ea4b6b7aabb75f2
Author: XiaoHongbo <[email protected]>
AuthorDate: Mon Aug 3 11:09:45 2026 +0800
perf(parquet): read row groups concurrently (#633)
---
.../datafusion/src/physical_plan/scan.rs | 76 ++++
crates/integrations/datafusion/src/table/mod.rs | 4 +-
.../datafusion/src/variant_pushdown.rs | 4 +-
crates/paimon/src/arrow/format/mod.rs | 17 +-
crates/paimon/src/arrow/format/parquet.rs | 467 ++++++++++++++++++++-
crates/paimon/src/arrow/mod.rs | 2 +
crates/paimon/src/arrow/parquet_read_budget.rs | 144 +++++++
crates/paimon/src/spec/core_options.rs | 101 +++++
crates/paimon/src/table/data_evolution_reader.rs | 35 +-
crates/paimon/src/table/data_file_reader.rs | 31 +-
crates/paimon/src/table/format_read_builder.rs | 17 +-
crates/paimon/src/table/format_table_read.rs | 21 +-
crates/paimon/src/table/kv_file_reader.rs | 150 ++++++-
crates/paimon/src/table/read_builder.rs | 43 +-
crates/paimon/src/table/table_read.rs | 108 ++++-
15 files changed, 1169 insertions(+), 51 deletions(-)
diff --git a/crates/integrations/datafusion/src/physical_plan/scan.rs
b/crates/integrations/datafusion/src/physical_plan/scan.rs
index bf8373dd..fe0a3859 100644
--- a/crates/integrations/datafusion/src/physical_plan/scan.rs
+++ b/crates/integrations/datafusion/src/physical_plan/scan.rs
@@ -49,6 +49,7 @@ use datafusion::physical_plan::filter_pushdown::{
use datafusion::physical_plan::stream::RecordBatchStreamAdapter;
use datafusion::physical_plan::{DisplayAs, ExecutionPlan, Partitioning,
PlanProperties};
use futures::{FutureExt, StreamExt, TryStreamExt};
+use paimon::arrow::ParquetReadBudget;
use paimon::spec::{DataField, Datum, MergeEngine, Predicate, PredicateBuilder,
PredicateOperator};
use paimon::table::{ScanTrace, Table};
use paimon::DataSplit;
@@ -775,9 +776,12 @@ pub struct PaimonTableScan {
/// Static filters already covered by `pushed_predicate` use Paimon's
native
/// Parquet row filter instead, avoiding duplicate decoder evaluation.
decoder_filters: Vec<Arc<dyn PhysicalExpr>>,
+ /// Query-wide budget shared by every DataFusion scan partition.
+ parquet_read_budget: Arc<ParquetReadBudget>,
}
impl PaimonTableScan {
+ #[cfg(test)]
#[allow(clippy::too_many_arguments)]
pub(crate) fn new(
schema: ArrowSchemaRef,
@@ -790,6 +794,75 @@ impl PaimonTableScan {
scan_trace: Option<ScanTrace>,
pushed_variants: Option<String>,
case_sensitive: bool,
+ ) -> Self {
+ Self::new_with_parquet_read_budget(
+ schema,
+ table,
+ read_type,
+ pushed_predicate,
+ planned_partitions,
+ limit,
+ filter_exact,
+ scan_trace,
+ pushed_variants,
+ case_sensitive,
+ Arc::new(ParquetReadBudget::default()),
+ )
+ }
+
+ #[allow(clippy::too_many_arguments)]
+ pub(crate) fn try_new(
+ schema: ArrowSchemaRef,
+ table: Table,
+ read_type: Vec<DataField>,
+ pushed_predicate: Option<Predicate>,
+ planned_partitions: Vec<Arc<[DataSplit]>>,
+ limit: Option<usize>,
+ filter_exact: bool,
+ scan_trace: Option<ScanTrace>,
+ pushed_variants: Option<String>,
+ case_sensitive: bool,
+ ) -> DFResult<Self> {
+ let options = table.schema().core_options();
+ let parquet_read_budget = Arc::new(
+ ParquetReadBudget::new(
+ options
+ .parquet_row_group_parallelism()
+ .map_err(to_datafusion_error)?,
+ options
+ .parquet_row_group_max_inflight_bytes()
+ .map_err(to_datafusion_error)?,
+ )
+ .map_err(to_datafusion_error)?,
+ );
+ Ok(Self::new_with_parquet_read_budget(
+ schema,
+ table,
+ read_type,
+ pushed_predicate,
+ planned_partitions,
+ limit,
+ filter_exact,
+ scan_trace,
+ pushed_variants,
+ case_sensitive,
+ parquet_read_budget,
+ ))
+ }
+
+ #[allow(clippy::too_many_arguments)]
+ fn new_with_parquet_read_budget(
+ schema: ArrowSchemaRef,
+ table: Table,
+ read_type: Vec<DataField>,
+ pushed_predicate: Option<Predicate>,
+ planned_partitions: Vec<Arc<[DataSplit]>>,
+ limit: Option<usize>,
+ filter_exact: bool,
+ scan_trace: Option<ScanTrace>,
+ pushed_variants: Option<String>,
+ case_sensitive: bool,
+ parquet_read_budget: Arc<ParquetReadBudget>,
) -> Self {
let plan_properties = Arc::new(PlanProperties::new(
EquivalenceProperties::new(schema.clone()),
@@ -810,6 +883,7 @@ impl PaimonTableScan {
case_sensitive,
runtime_filters: Vec::new(),
decoder_filters: Vec::new(),
+ parquet_read_budget,
}
}
@@ -997,6 +1071,7 @@ impl ExecutionPlan for PaimonTableScan {
let case_sensitive = self.case_sensitive;
let runtime_filters = self.runtime_filters.clone();
let decoder_filters = self.decoder_filters.clone();
+ let parquet_read_budget = Arc::clone(&self.parquet_read_budget);
let fut = async move {
let mut read_builder = table.new_read_builder();
@@ -1013,6 +1088,7 @@ impl ExecutionPlan for PaimonTableScan {
if !paimon_predicates.is_empty() {
read_builder.with_filter(Predicate::and(paimon_predicates));
}
+ read_builder.with_parquet_read_budget(parquet_read_budget);
let mut read =
read_builder.new_read().map_err(to_datafusion_error)?;
if !runtime_filter_plan.datafusion_filters.is_empty() {
diff --git a/crates/integrations/datafusion/src/table/mod.rs
b/crates/integrations/datafusion/src/table/mod.rs
index 1bfc47fc..113c3bb8 100644
--- a/crates/integrations/datafusion/src/table/mod.rs
+++ b/crates/integrations/datafusion/src/table/mod.rs
@@ -372,7 +372,7 @@ impl PaimonScanBuilder<'_> {
.collect()
};
- Ok(Arc::new(PaimonTableScan::new(
+ Ok(Arc::new(PaimonTableScan::try_new(
projected_schema,
self.table.clone(),
read_type,
@@ -383,7 +383,7 @@ impl PaimonScanBuilder<'_> {
self.scan_trace,
None,
self.case_sensitive,
- )))
+ )?))
}
}
diff --git a/crates/integrations/datafusion/src/variant_pushdown.rs
b/crates/integrations/datafusion/src/variant_pushdown.rs
index a5bb3482..641c6a5a 100644
--- a/crates/integrations/datafusion/src/variant_pushdown.rs
+++ b/crates/integrations/datafusion/src/variant_pushdown.rs
@@ -248,7 +248,7 @@ impl ExtensionPlanner for VariantExtractionExtensionPlanner
{
.as_ref()
.is_none_or(|p| read_builder.is_exact_filter_pushdown(p));
- Ok(Some(Arc::new(PaimonTableScan::new(
+ Ok(Some(Arc::new(PaimonTableScan::try_new(
Arc::clone(&node.arrow_schema),
node.table.clone(),
node.read_type.clone(),
@@ -259,7 +259,7 @@ impl ExtensionPlanner for VariantExtractionExtensionPlanner
{
Some(scan_trace),
Some(node.pushed_variants.clone()),
case_sensitive,
- ))))
+ )?)))
}
}
diff --git a/crates/paimon/src/arrow/format/mod.rs
b/crates/paimon/src/arrow/format/mod.rs
index 4e0990f8..dde85a38 100644
--- a/crates/paimon/src/arrow/format/mod.rs
+++ b/crates/paimon/src/arrow/format/mod.rs
@@ -28,6 +28,7 @@ mod vortex;
#[cfg(test)]
pub(crate) use parquet::ParquetFormatWriter;
+use super::ParquetReadBudget;
use super::RowFilterFactory;
use crate::io::{FileRead, OutputFile};
use crate::spec::stats::BinaryTableStats;
@@ -159,14 +160,28 @@ impl FormatWriteResult {
}
/// Create a format reader based on the file extension.
+#[cfg(test)]
pub(crate) fn create_format_reader(
path: &str,
blob_as_descriptor: bool,
read_fields: &[DataField],
+) -> crate::Result<Box<dyn FormatFileReader>> {
+ create_format_reader_with_budget(path, blob_as_descriptor, read_fields,
None)
+}
+
+/// 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,
+ read_fields: &[DataField],
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
) -> crate::Result<Box<dyn FormatFileReader>> {
let lower = path.to_ascii_lowercase();
let reader: Box<dyn FormatFileReader> = if lower.ends_with(".parquet") {
- Box::new(parquet::ParquetFormatReader)
+ Box::new(match parquet_read_budget {
+ Some(read_budget) =>
parquet::ParquetFormatReader::with_read_budget(read_budget),
+ None => parquet::ParquetFormatReader::default(),
+ })
} else if lower.ends_with(".blob") {
Box::new(blob::BlobFormatReader::new(
path.to_string(),
diff --git a/crates/paimon/src/arrow/format/parquet.rs
b/crates/paimon/src/arrow/format/parquet.rs
index c851299f..58491401 100644
--- a/crates/paimon/src/arrow/format/parquet.rs
+++ b/crates/paimon/src/arrow/format/parquet.rs
@@ -18,9 +18,10 @@
use super::shredding::PhysicalFormatWriterFactory;
use super::{FilePredicates, FormatFileReader, FormatFileWriter,
FormatWriteResult};
use crate::arrow::filtering::{predicates_may_match_with_schema, StatsAccessor};
+use crate::arrow::parquet_read_budget::ParquetReadPermit;
use crate::arrow::shredding::map::MapShreddingReadPlan;
use crate::arrow::shredding::ShreddingReadPlan;
-use crate::arrow::{RowFilter, RowFilterContext};
+use crate::arrow::{ParquetReadBudget, RowFilter, RowFilterContext};
use crate::io::{FileRead, OutputFile};
use crate::spec::stats::BinaryTableStats;
use crate::spec::{
@@ -35,8 +36,8 @@ use bytes::Bytes;
use futures::future::BoxFuture;
use futures::{FutureExt, StreamExt, TryFutureExt, TryStreamExt};
use parquet::arrow::arrow_reader::{
- ArrowPredicate, ArrowPredicateFn, ArrowReaderOptions, RowFilter as
ParquetRowFilter,
- RowSelection, RowSelector,
+ ArrowPredicate, ArrowPredicateFn, ArrowReaderMetadata, ArrowReaderOptions,
+ RowFilter as ParquetRowFilter, RowSelection, RowSelector,
};
use parquet::arrow::async_reader::{AsyncFileReader, MetadataFetch};
use parquet::arrow::{AsyncArrowWriter, ParquetRecordBatchStreamBuilder,
ProjectionMask};
@@ -51,8 +52,26 @@ use std::cmp::Ordering;
use std::collections::HashMap;
use std::ops::Range;
use std::sync::Arc;
+use tokio::sync::mpsc;
-pub(crate) struct ParquetFormatReader;
+#[derive(Default)]
+pub(crate) struct ParquetFormatReader {
+ read_budget: Option<Arc<ParquetReadBudget>>,
+}
+
+impl ParquetFormatReader {
+ pub(crate) fn with_read_budget(read_budget: Arc<ParquetReadBudget>) ->
Self {
+ Self {
+ read_budget: Some(read_budget),
+ }
+ }
+}
+
+enum ParquetRowGroupMessage {
+ Batch(RecordBatch),
+ Error(Error),
+ Done,
+}
struct ParquetRowFilterPredicate {
inner: Box<dyn RowFilter>,
@@ -301,7 +320,8 @@ impl FormatFileReader for ParquetFormatReader {
batch_size: Option<usize>,
row_selection: Option<Vec<RowRange>>,
) -> crate::Result<ArrowRecordBatchStream> {
- let arrow_file_reader = ArrowFileReader::new(file_size, reader);
+ let shared_reader: Arc<dyn FileRead> = reader.into();
+ let arrow_file_reader = ArrowFileReader::new(file_size,
Arc::clone(&shared_reader));
let empty_predicates = Vec::new();
let (preds, file_fields): (&[Predicate], &[DataField]) = match
predicates {
@@ -361,7 +381,7 @@ impl FormatFileReader for ParquetFormatReader {
.collect();
let mask = ProjectionMask::roots(&parquet_schema, root_indices);
- batch_stream_builder = batch_stream_builder.with_projection(mask);
+ batch_stream_builder =
batch_stream_builder.with_projection(mask.clone());
let mut decoder_predicates = build_parquet_row_filter(&parquet_schema,
preds, file_fields)?
.map(ParquetRowFilter::into_predicates)
@@ -460,6 +480,107 @@ impl FormatFileReader for ParquetFormatReader {
MapShreddingReadPlan::create(&scan_fields,
batch_stream_builder.schema())?
.map(Arc::new);
+ // A normal Parquet stream fetches and decodes row groups one by one.
+ // For remote object stores, a full scan of a compacted file can
+ // therefore serialize dozens of independent range requests behind one
+ // DataFusion partition. Build one stream per row group on the
+ // predicate-free path and run a bounded number concurrently.
+ //
+ // Row-group receivers are consumed in order and buffer one batch each,
+ // preserving positional `_ROW_ID`, sort order, and batch
backpressure. Reads
+ // with predicates or an explicit row selection retain the original
+ // single-stream path until their selections are split per row group.
+ let row_group_parallelism = self
+ .read_budget
+ .as_ref()
+ .filter(|_| preds.is_empty() && row_filter_factory.is_none() &&
row_selection.is_none())
+ .map(|budget| {
+ budget
+ .parallelism()
+ .min(batch_stream_builder.metadata().num_row_groups())
+ })
+ .unwrap_or(1);
+ if row_group_parallelism > 1 {
+ let row_group_count =
batch_stream_builder.metadata().num_row_groups();
+ let reader_metadata = ArrowReaderMetadata::try_new(
+ batch_stream_builder.metadata().clone(),
+ ArrowReaderOptions::new(),
+ )?;
+ let projected_bytes = batch_stream_builder
+ .metadata()
+ .row_groups()
+ .iter()
+ .map(|row_group| projected_row_group_bytes(row_group, &mask))
+ .collect::<Vec<_>>();
+ let read_budget =
Arc::clone(self.read_budget.as_ref().expect("checked above"));
+ let (row_group_tx, mut row_group_rx) =
mpsc::channel(row_group_parallelism);
+ tokio::spawn(async move {
+ for (row_group_index, projected_bytes) in
projected_bytes.into_iter().enumerate() {
+ let Ok(slot) = row_group_tx.reserve().await else {
+ return;
+ };
+ let permit = match tokio::select! {
+ _ = row_group_tx.closed() => return,
+ permit = read_budget.acquire(projected_bytes) =>
permit,
+ } {
+ Ok(permit) => permit,
+ Err(error) => {
+ slot.send(Err(error));
+ return;
+ }
+ };
+ let (batch_tx, batch_rx) = mpsc::channel(1);
+ let row_group_reader = Arc::clone(&shared_reader);
+ let row_group_metadata = reader_metadata.clone();
+ let row_group_mask = mask.clone();
+ tokio::spawn(read_row_group(
+ row_group_reader,
+ file_size,
+ row_group_metadata,
+ row_group_mask,
+ row_group_index,
+ batch_size,
+ permit,
+ batch_tx,
+ ));
+ slot.send(Ok(batch_rx));
+ }
+ });
+ let stream = async_stream::try_stream! {
+ for _ in 0..row_group_count {
+ let mut batches = row_group_rx.recv().await.ok_or_else(|| {
+ Error::UnexpectedError {
+ message: "Parquet row-group coordinator stopped
early".to_string(),
+ source: None,
+ }
+ })??;
+ let mut completed = false;
+ while let Some(message) = batches.recv().await {
+ match message {
+ ParquetRowGroupMessage::Batch(batch) => {
+ yield match &map_read_plan {
+ Some(plan) => plan.assemble_batch(&batch)?,
+ None => batch,
+ };
+ }
+ ParquetRowGroupMessage::Error(error) =>
Err(error)?,
+ ParquetRowGroupMessage::Done => {
+ completed = true;
+ break;
+ }
+ }
+ }
+ if !completed {
+ Err(Error::UnexpectedError {
+ message: "Parquet row-group reader stopped
early".to_string(),
+ source: None,
+ })?;
+ }
+ }
+ };
+ return Ok(stream.boxed());
+ }
+
let batch_stream = batch_stream_builder.build()?;
if all_enforced {
@@ -504,6 +625,78 @@ impl FormatFileReader for ParquetFormatReader {
}
}
+fn projected_row_group_bytes(row_group: &RowGroupMetaData, projection:
&ProjectionMask) -> u64 {
+ row_group
+ .columns()
+ .iter()
+ .enumerate()
+ .filter(|(leaf_index, _)| projection.leaf_included(*leaf_index))
+ .filter_map(|(_, column)|
u64::try_from(column.uncompressed_size()).ok())
+ .fold(0u64, u64::saturating_add)
+}
+
+#[allow(clippy::too_many_arguments)]
+async fn read_row_group(
+ reader: Arc<dyn FileRead>,
+ file_size: u64,
+ reader_metadata: ArrowReaderMetadata,
+ projection: ProjectionMask,
+ row_group_index: usize,
+ batch_size: Option<usize>,
+ _permit: ParquetReadPermit,
+ sender: mpsc::Sender<ParquetRowGroupMessage>,
+) {
+ let mut builder = ParquetRecordBatchStreamBuilder::new_with_metadata(
+ ArrowFileReader::new(file_size, reader),
+ reader_metadata,
+ )
+ .with_projection(projection)
+ .with_row_groups(vec![row_group_index]);
+ if let Some(size) = batch_size {
+ builder = builder.with_batch_size(size);
+ }
+ let mut stream = match builder.build() {
+ Ok(stream) => stream,
+ Err(error) => {
+ let _ = sender
+ .send(ParquetRowGroupMessage::Error(error.into()))
+ .await;
+ return;
+ }
+ };
+
+ forward_row_group_batches(&mut stream, sender).await;
+}
+
+async fn forward_row_group_batches<S, E>(
+ mut stream: S,
+ sender: mpsc::Sender<ParquetRowGroupMessage>,
+) where
+ S: futures::Stream<Item = std::result::Result<RecordBatch, E>> + Unpin,
+ E: Into<Error>,
+{
+ loop {
+ let Ok(slot) = sender.reserve().await else {
+ return;
+ };
+ let next = tokio::select! {
+ _ = sender.closed() => return,
+ next = stream.next() => next,
+ };
+ match next {
+ Some(Ok(batch)) => slot.send(ParquetRowGroupMessage::Batch(batch)),
+ Some(Err(error)) => {
+ slot.send(ParquetRowGroupMessage::Error(error.into()));
+ return;
+ }
+ None => {
+ slot.send(ParquetRowGroupMessage::Done);
+ return;
+ }
+ }
+ }
+}
+
// ---------------------------------------------------------------------------
// Parquet row-filter helpers
// ---------------------------------------------------------------------------
@@ -1693,7 +1886,7 @@ fn build_row_ranges_selection(
/// - `preload_offset_index`: Load the Offset Index as part of
[`Self::get_metadata`].
struct ArrowFileReader {
file_size: u64,
- r: Box<dyn FileRead>,
+ r: Arc<dyn FileRead>,
}
/// coalesce threshold: 1 MiB.
@@ -1709,7 +1902,7 @@ const METADATA_SIZE_HINT: usize = 512 * 1024;
const IO_BLOCK_SIZE: u64 = 4 * 1024 * 1024;
impl ArrowFileReader {
- fn new(file_size: u64, r: Box<dyn FileRead>) -> Self {
+ fn new(file_size: u64, r: Arc<dyn FileRead>) -> Self {
Self { file_size, r }
}
@@ -1974,15 +2167,18 @@ fn split_ranges_for_concurrency(merged:
Vec<Range<u64>>, concurrency: usize) ->
#[allow(clippy::type_complexity)] // test row literals use nested
Option<Vec<(&str, Option<i64>)>>
mod tests {
use super::build_parquet_row_filter;
+ use super::{
+ forward_row_group_batches, FilePredicates, ParquetFormatReader,
ParquetFormatWriter,
+ ParquetRowGroupMessage,
+ };
use super::{
AsyncArrowWriter, Bytes, PageIndexPolicy, ParquetMetaDataReader,
Predicate,
PredicateOperator, RowSelection,
};
- use super::{FilePredicates, ParquetFormatReader, ParquetFormatWriter};
use crate::arrow::format::{
create_format_reader, create_format_writer, FormatFileReader,
FormatFileWriter,
};
- use crate::arrow::{build_target_arrow_schema, variant_arrow_type};
+ use crate::arrow::{build_target_arrow_schema, variant_arrow_type,
ParquetReadBudget};
use crate::io::FileIOBuilder;
use crate::spec::{
ArrayType, BigIntType, DataField, DataType, Datum, IntType, MapType,
PredicateBuilder,
@@ -1990,15 +2186,19 @@ mod tests {
};
use crate::table::RowRange;
use crate::variant::GenericVariant;
+ use crate::Error;
use arrow_array::{
Array, BinaryArray, Int32Array, Int64Array, MapArray, RecordBatch,
StringArray, StructArray,
};
use arrow_buffer::{NullBuffer, OffsetBuffer, ScalarBuffer};
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema
as ArrowSchema};
- use futures::TryStreamExt;
+ use futures::{StreamExt, TryStreamExt};
use parquet::schema::{parser::parse_message_type, types::SchemaDescriptor};
use std::collections::HashMap;
+ use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
use std::sync::Arc;
+ use std::time::Duration;
+ use tokio::sync::mpsc;
fn test_fields() -> Vec<DataField> {
vec![
@@ -2433,6 +2633,241 @@ mod tests {
.unwrap()
}
+ #[derive(Clone)]
+ struct ConcurrentTrackingFileRead {
+ data: Bytes,
+ in_flight: Arc<AtomicUsize>,
+ max_in_flight: Arc<AtomicUsize>,
+ }
+
+ #[async_trait::async_trait]
+ impl crate::io::FileRead for ConcurrentTrackingFileRead {
+ async fn read(&self, range: std::ops::Range<u64>) ->
crate::Result<Bytes> {
+ let current = self.in_flight.fetch_add(1, AtomicOrdering::SeqCst)
+ 1;
+ self.max_in_flight
+ .fetch_max(current, AtomicOrdering::SeqCst);
+ tokio::time::sleep(Duration::from_millis(10)).await;
+ self.in_flight.fetch_sub(1, AtomicOrdering::SeqCst);
+ Ok(self.data.slice(range.start as usize..range.end as usize))
+ }
+ }
+
+ #[tokio::test]
+ async fn test_parquet_reader_reads_row_groups_concurrently_in_order() {
+ const ROWS: i32 = 512;
+ let schema = writer_arrow_schema();
+ let props = parquet::file::properties::WriterProperties::builder()
+ .set_max_row_group_row_count(Some(64))
+ .set_dictionary_enabled(false)
+ .build();
+ let mut data = Vec::new();
+ let mut writer = AsyncArrowWriter::try_new(&mut data, schema.clone(),
Some(props)).unwrap();
+ let ids = (0..ROWS).collect::<Vec<_>>();
+ writer
+ .write(&writer_test_batch(
+ &schema,
+ ids.clone(),
+ ids.iter().map(|id| id * 10).collect(),
+ ))
+ .await
+ .unwrap();
+ writer.close().await.unwrap();
+
+ let in_flight = Arc::new(AtomicUsize::new(0));
+ let max_in_flight = Arc::new(AtomicUsize::new(0));
+ let file_reader = ConcurrentTrackingFileRead {
+ data: Bytes::from(data),
+ in_flight,
+ max_in_flight: Arc::clone(&max_in_flight),
+ };
+ let file_size = file_reader.data.len() as u64;
+ let fields = vec![DataField::new(
+ 0,
+ "id".to_string(),
+ DataType::Int(IntType::new()),
+ )];
+ let batches = ParquetFormatReader::with_read_budget(Arc::new(
+ ParquetReadBudget::new(8, 256 * 1024 * 1024).unwrap(),
+ ))
+ .read_batch_stream(
+ Box::new(file_reader),
+ file_size,
+ &fields,
+ None,
+ Some(32),
+ None,
+ )
+ .await
+ .unwrap()
+ .try_collect::<Vec<_>>()
+ .await
+ .unwrap();
+ let actual = batches
+ .iter()
+ .flat_map(|batch| {
+ batch
+ .column(0)
+ .as_any()
+ .downcast_ref::<Int32Array>()
+ .unwrap()
+ .values()
+ .iter()
+ .copied()
+ .collect::<Vec<_>>()
+ })
+ .collect::<Vec<_>>();
+
+ assert_eq!(actual, ids);
+ assert!(
+ max_in_flight.load(AtomicOrdering::SeqCst) > 1,
+ "row-group reads should overlap"
+ );
+ }
+
+ #[tokio::test]
+ async fn test_parquet_read_budget_is_shared_across_readers() {
+ const ROWS: i32 = 256;
+ let schema = writer_arrow_schema();
+ let props = parquet::file::properties::WriterProperties::builder()
+ .set_max_row_group_row_count(Some(32))
+ .set_dictionary_enabled(false)
+ .build();
+ let mut data = Vec::new();
+ let mut writer = AsyncArrowWriter::try_new(&mut data, schema.clone(),
Some(props)).unwrap();
+ let ids = (0..ROWS).collect::<Vec<_>>();
+ writer
+ .write(&writer_test_batch(
+ &schema,
+ ids.clone(),
+ ids.iter().map(|id| id * 10).collect(),
+ ))
+ .await
+ .unwrap();
+ writer.close().await.unwrap();
+
+ let data = Bytes::from(data);
+ let in_flight = Arc::new(AtomicUsize::new(0));
+ let max_in_flight = Arc::new(AtomicUsize::new(0));
+ let new_file_reader = || ConcurrentTrackingFileRead {
+ data: data.clone(),
+ in_flight: Arc::clone(&in_flight),
+ max_in_flight: Arc::clone(&max_in_flight),
+ };
+ let file_size = data.len() as u64;
+ let fields = vec![DataField::new(
+ 0,
+ "id".to_string(),
+ DataType::Int(IntType::new()),
+ )];
+ let budget = Arc::new(ParquetReadBudget::new(2, 256 * 1024 *
1024).unwrap());
+ let first = ParquetFormatReader::with_read_budget(Arc::clone(&budget))
+ .read_batch_stream(
+ Box::new(new_file_reader()),
+ file_size,
+ &fields,
+ None,
+ Some(32),
+ None,
+ )
+ .await
+ .unwrap();
+ let second = ParquetFormatReader::with_read_budget(budget)
+ .read_batch_stream(
+ Box::new(new_file_reader()),
+ file_size,
+ &fields,
+ None,
+ Some(32),
+ None,
+ )
+ .await
+ .unwrap();
+
+ max_in_flight.store(0, AtomicOrdering::SeqCst);
+ let (first, second) = tokio::join!(
+ first.try_collect::<Vec<_>>(),
+ second.try_collect::<Vec<_>>()
+ );
+ assert_eq!(
+ first
+ .unwrap()
+ .iter()
+ .map(RecordBatch::num_rows)
+ .sum::<usize>(),
+ ROWS as usize
+ );
+ assert_eq!(
+ second
+ .unwrap()
+ .iter()
+ .map(RecordBatch::num_rows)
+ .sum::<usize>(),
+ ROWS as usize
+ );
+ let observed = max_in_flight.load(AtomicOrdering::SeqCst);
+ assert_eq!(
+ observed, 2,
+ "two readers must share the same row-group budget"
+ );
+ }
+
+ #[tokio::test]
+ async fn test_row_group_batch_forwarding_applies_backpressure() {
+ let schema = Arc::new(ArrowSchema::empty());
+ let batch = RecordBatch::try_new_with_options(
+ Arc::clone(&schema),
+ Vec::new(),
+ &arrow_array::RecordBatchOptions::new().with_row_count(Some(1)),
+ )
+ .unwrap();
+ let polls = Arc::new(AtomicUsize::new(0));
+ let tracked_polls = Arc::clone(&polls);
+ let stream =
+ futures::stream::iter(vec![Ok::<_, Error>(batch.clone()), Ok::<_,
Error>(batch)])
+ .inspect(move |_| {
+ tracked_polls.fetch_add(1, AtomicOrdering::SeqCst);
+ });
+ let (tx, mut rx) = mpsc::channel(1);
+ let task = tokio::spawn(forward_row_group_batches(stream, tx));
+
+ tokio::time::sleep(Duration::from_millis(20)).await;
+ assert_eq!(
+ polls.load(AtomicOrdering::SeqCst),
+ 1,
+ "a full output channel must stop polling and decoding the source
stream"
+ );
+ assert!(matches!(
+ rx.recv().await,
+ Some(ParquetRowGroupMessage::Batch(_))
+ ));
+
+ tokio::time::sleep(Duration::from_millis(20)).await;
+ assert_eq!(polls.load(AtomicOrdering::SeqCst), 2);
+ drop(rx);
+ task.await.unwrap();
+ }
+
+ #[tokio::test]
+ async fn test_row_group_batch_forwarding_stops_during_pending_io() {
+ let (polled_tx, polled_rx) = tokio::sync::oneshot::channel();
+ let mut polled_tx = Some(polled_tx);
+ let stream = futures::stream::poll_fn(move |_| {
+ if let Some(tx) = polled_tx.take() {
+ let _ = tx.send(());
+ }
+ std::task::Poll::Pending::<Option<Result<RecordBatch, Error>>>
+ });
+ let (tx, rx) = mpsc::channel(1);
+ let task = tokio::spawn(forward_row_group_batches(stream, tx));
+
+ polled_rx.await.unwrap();
+ drop(rx);
+ tokio::time::timeout(Duration::from_millis(100), task)
+ .await
+ .expect("dropping the receiver must cancel a pending row-group
read")
+ .unwrap();
+ }
+
#[tokio::test]
async fn test_parquet_writer_write_and_close() {
let file_io = FileIOBuilder::new("memory").build().unwrap();
@@ -2915,7 +3350,7 @@ mod tests {
"items".to_string(),
DataType::Array(ArrayType::new(DataType::Int(IntType::new()))),
)];
- let stream = ParquetFormatReader
+ let stream = ParquetFormatReader::default()
.read_batch_stream(
Box::new(file_read),
file_size,
@@ -3304,7 +3739,7 @@ mod tests {
file_fields: fields.clone(),
};
- let batches = ParquetFormatReader
+ let batches = ParquetFormatReader::default()
.read_batch_stream(
Box::new(input.reader().await.unwrap()),
file_size,
@@ -3432,7 +3867,7 @@ mod tests {
file_fields: fields.clone(),
};
- ParquetFormatReader
+ ParquetFormatReader::default()
.read_batch_stream(
Box::new(input.reader().await.unwrap()),
file_size,
@@ -3477,7 +3912,7 @@ mod tests {
file_fields: id_name_age_file_fields(),
};
- let reader = ParquetFormatReader;
+ let reader = ParquetFormatReader::default();
let mut stream = reader
.read_batch_stream(
Box::new(reader_input),
@@ -3671,7 +4106,7 @@ mod tests {
row_filter_factory: None,
file_fields,
};
- let reader = ParquetFormatReader;
+ let reader = ParquetFormatReader::default();
let mut stream = reader
.read_batch_stream(
Box::new(reader_input),
diff --git a/crates/paimon/src/arrow/mod.rs b/crates/paimon/src/arrow/mod.rs
index 37fa81f7..2fe1a6e2 100644
--- a/crates/paimon/src/arrow/mod.rs
+++ b/crates/paimon/src/arrow/mod.rs
@@ -17,11 +17,13 @@
pub(crate) mod filtering;
pub(crate) mod format;
+mod parquet_read_budget;
pub(crate) mod residual;
mod row_filter;
pub(crate) mod schema_evolution;
pub(crate) mod shredding;
+pub use parquet_read_budget::ParquetReadBudget;
pub use row_filter::{RowFilter, RowFilterContext, RowFilterFactory};
use crate::spec::{
diff --git a/crates/paimon/src/arrow/parquet_read_budget.rs
b/crates/paimon/src/arrow/parquet_read_budget.rs
new file mode 100644
index 00000000..e0f6e5cc
--- /dev/null
+++ b/crates/paimon/src/arrow/parquet_read_budget.rs
@@ -0,0 +1,144 @@
+// Licensed to the Apache Software Foundation (ASF) under one
+// or more contributor license agreements. See the NOTICE file
+// distributed with this work for additional information
+// regarding copyright ownership. The ASF licenses this file
+// to you under the Apache License, Version 2.0 (the
+// "License"); you may not use this file except in compliance
+// with the License. You may obtain a copy of the License at
+//
+// http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing,
+// software distributed under the License is distributed on an
+// "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+// KIND, either express or implied. See the License for the
+// specific language governing permissions and limitations
+// under the License.
+
+use std::sync::Arc;
+
+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;
+
+/// Shared resource budget for concurrent Parquet row-group reads.
+#[derive(Debug)]
+pub struct ParquetReadBudget {
+ parallelism: usize,
+ row_groups: Arc<Semaphore>,
+ bytes: Arc<Semaphore>,
+ byte_permits: u32,
+}
+
+impl ParquetReadBudget {
+ pub fn new(parallelism: usize, max_inflight_bytes: u64) ->
crate::Result<Self> {
+ if parallelism == 0 || parallelism > Semaphore::MAX_PERMITS {
+ return Err(crate::Error::DataInvalid {
+ message: format!(
+ "Parquet row-group parallelism must be between 1 and {},
got {parallelism}",
+ Semaphore::MAX_PERMITS
+ ),
+ source: None,
+ });
+ }
+ if max_inflight_bytes == 0 {
+ return Err(crate::Error::DataInvalid {
+ message: "Parquet row-group max in-flight bytes must be
greater than 0".to_string(),
+ 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)
+ .min(u64::from(max_byte_permits)) as u32;
+
+ Ok(Self {
+ parallelism,
+ row_groups: Arc::new(Semaphore::new(parallelism)),
+ bytes: Arc::new(Semaphore::new(byte_permits as usize)),
+ byte_permits,
+ })
+ }
+
+ pub fn parallelism(&self) -> usize {
+ self.parallelism
+ }
+
+ pub(crate) async fn acquire(
+ &self,
+ projected_uncompressed_bytes: u64,
+ ) -> crate::Result<ParquetReadPermit> {
+ let row_group = Arc::clone(&self.row_groups)
+ .acquire_owned()
+ .await
+ .map_err(|error| crate::Error::UnexpectedError {
+ message: "Parquet row-group read budget was
closed".to_string(),
+ source: Some(Box::new(error)),
+ })?;
+ let requested = projected_uncompressed_bytes
+ .max(1)
+ .div_ceil(BYTE_PERMIT_UNIT)
+ .min(u64::from(self.byte_permits)) as u32;
+ let bytes = Arc::clone(&self.bytes)
+ .acquire_many_owned(requested)
+ .await
+ .map_err(|error| crate::Error::UnexpectedError {
+ message: "Parquet byte read budget was closed".to_string(),
+ source: Some(Box::new(error)),
+ })?;
+ Ok(ParquetReadPermit {
+ _row_group: row_group,
+ _bytes: bytes,
+ })
+ }
+}
+
+impl Default for ParquetReadBudget {
+ fn default() -> Self {
+ Self::new(DEFAULT_PARALLELISM, DEFAULT_MAX_INFLIGHT_BYTES)
+ .expect("default Parquet read budget is valid")
+ }
+}
+
+#[derive(Debug)]
+pub(crate) struct ParquetReadPermit {
+ _row_group: OwnedSemaphorePermit,
+ _bytes: OwnedSemaphorePermit,
+}
+
+#[cfg(test)]
+mod tests {
+ use super::*;
+ use std::time::Duration;
+
+ #[tokio::test]
+ async fn shared_budget_blocks_until_permits_are_released() {
+ let budget = Arc::new(ParquetReadBudget::new(2,
BYTE_PERMIT_UNIT).unwrap());
+ let first = budget.acquire(2 * BYTE_PERMIT_UNIT).await.unwrap();
+
+ assert!(
+ tokio::time::timeout(Duration::from_millis(20), budget.acquire(1))
+ .await
+ .is_err(),
+ "the projected-byte budget must be shared across readers"
+ );
+
+ drop(first);
+ tokio::time::timeout(Duration::from_secs(1), budget.acquire(1))
+ .await
+ .expect("dropping a read must release its permits")
+ .unwrap();
+ }
+
+ #[test]
+ fn rejects_invalid_limits() {
+ assert!(ParquetReadBudget::new(0, BYTE_PERMIT_UNIT).is_err());
+ assert!(ParquetReadBudget::new(1, 0).is_err());
+ assert!(
+ ParquetReadBudget::new(Semaphore::MAX_PERMITS.saturating_add(1),
BYTE_PERMIT_UNIT)
+ .is_err()
+ );
+ }
+}
diff --git a/crates/paimon/src/spec/core_options.rs
b/crates/paimon/src/spec/core_options.rs
index 507471c8..e26c6bc3 100644
--- a/crates/paimon/src/spec/core_options.rs
+++ b/crates/paimon/src/spec/core_options.rs
@@ -75,6 +75,9 @@ const MANIFEST_TARGET_SIZE_OPTION: &str =
"manifest.target-size";
const MANIFEST_MERGE_MIN_COUNT_OPTION: &str = "manifest.merge-min-count";
const WRITE_PARQUET_BUFFER_SIZE_OPTION: &str = "write.parquet-buffer-size";
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";
pub(crate) const SEQUENCE_FIELD_OPTION: &str = "sequence.field";
pub(crate) const DISABLE_EXPLICIT_TYPE_CASTING_OPTION: &str =
"disable-explicit-type-casting";
pub(crate) const DISABLE_ALTER_COLUMN_NULL_TO_NOT_NULL_OPTION: &str =
@@ -114,6 +117,8 @@ const DEFAULT_CHANGELOG_FILE_PREFIX: &str = "changelog-";
const DEFAULT_TARGET_FILE_SIZE: i64 = 256 * 1024 * 1024;
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 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;
@@ -338,6 +343,55 @@ impl<'a> CoreOptions<'a> {
Ok(value as usize)
}
+ /// Maximum concurrent Parquet row-group reads per scan. Set to `1` to
disable.
+ pub fn parquet_row_group_parallelism(&self) -> crate::Result<usize> {
+ let Some(raw) = self.options.get(PARQUET_ROW_GROUP_PARALLELISM_OPTION)
else {
+ return Ok(DEFAULT_PARQUET_ROW_GROUP_PARALLELISM);
+ };
+ let value = raw
+ .parse::<usize>()
+ .map_err(|error| crate::Error::DataInvalid {
+ message: format!(
+ "Option '{PARQUET_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 '{PARQUET_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<u64> {
+ let value = match self
+ .options
+ .get(PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION)
+ {
+ Some(raw) => parse_memory_size(raw).ok_or_else(||
crate::Error::DataInvalid {
+ message: format!(
+ "Option '{PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION}'
must be a valid memory size, got: {raw}"
+ ),
+ source: None,
+ })?,
+ None => DEFAULT_PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES,
+ };
+ u64::try_from(value)
+ .ok()
+ .filter(|value| *value > 0)
+ .ok_or_else(|| crate::Error::DataInvalid {
+ message: format!(
+ "Option '{PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION}'
must be greater than 0, got: {value}"
+ ),
+ source: None,
+ })
+ }
+
/// Reject scan options whose semantics the Rust core does not yet
implement.
///
/// These are not malformed input — they are unimplemented scan modes — so
@@ -1290,6 +1344,53 @@ mod tests {
}
}
+ #[test]
+ fn test_parquet_row_group_read_budget_options() {
+ let options = HashMap::new();
+ let core = CoreOptions::new(&options);
+ assert_eq!(core.parquet_row_group_parallelism().unwrap(), 8);
+ assert_eq!(
+ core.parquet_row_group_max_inflight_bytes().unwrap(),
+ 256 * 1024 * 1024
+ );
+
+ let options = HashMap::from([
+ (
+ PARQUET_ROW_GROUP_PARALLELISM_OPTION.to_string(),
+ "3".to_string(),
+ ),
+ (
+ PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION.to_string(),
+ "64 mb".to_string(),
+ ),
+ ]);
+ let core = CoreOptions::new(&options);
+ assert_eq!(core.parquet_row_group_parallelism().unwrap(), 3);
+ assert_eq!(
+ core.parquet_row_group_max_inflight_bytes().unwrap(),
+ 64 * 1024 * 1024
+ );
+
+ for value in ["0", "-1", "invalid"] {
+ let options = HashMap::from([(
+ PARQUET_ROW_GROUP_PARALLELISM_OPTION.to_string(),
+ value.to_string(),
+ )]);
+ assert!(CoreOptions::new(&options)
+ .parquet_row_group_parallelism()
+ .is_err());
+ }
+ for value in ["0", "-1", "invalid", "9223372036854775807 tb"] {
+ let options = HashMap::from([(
+ PARQUET_ROW_GROUP_MAX_INFLIGHT_BYTES_OPTION.to_string(),
+ value.to_string(),
+ )]);
+ assert!(CoreOptions::new(&options)
+ .parquet_row_group_max_inflight_bytes()
+ .is_err());
+ }
+ }
+
#[test]
fn test_source_split_defaults() {
let options = HashMap::new();
diff --git a/crates/paimon/src/table/data_evolution_reader.rs
b/crates/paimon/src/table/data_evolution_reader.rs
index 03c76732..b86cf213 100644
--- a/crates/paimon/src/table/data_evolution_reader.rs
+++ b/crates/paimon/src/table/data_evolution_reader.rs
@@ -22,8 +22,8 @@ use super::data_file_reader::{
append_null_row_id_column, attach_row_id, expand_selected_row_ids,
insert_column_at,
DataFileReader,
};
-use crate::arrow::build_target_arrow_schema;
use crate::arrow::format::FilePredicates;
+use crate::arrow::{build_target_arrow_schema, ParquetReadBudget};
use crate::deletion_vector::{DeletionVector, DeletionVectorFactory};
use crate::io::FileIO;
use crate::spec::{
@@ -114,6 +114,7 @@ pub(crate) struct DataEvolutionReader {
blob_view_rest_env: Option<RESTEnv>,
blob_read_limiter: BlobReadLimiter,
batch_size: Option<usize>,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
}
impl DataEvolutionReader {
@@ -176,6 +177,7 @@ impl DataEvolutionReader {
blob_view_rest_env,
blob_read_limiter: BlobReadLimiter::new(),
batch_size: None,
+ parquet_read_budget: None,
})
}
@@ -184,6 +186,14 @@ impl DataEvolutionReader {
self
}
+ pub(crate) fn with_parquet_read_budget(
+ mut self,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
+ ) -> Self {
+ self.parquet_read_budget = parquet_read_budget;
+ self
+ }
+
/// Read data files in data evolution mode.
pub fn read(self, data_splits: &[DataSplit]) ->
crate::Result<ArrowRecordBatchStream> {
let splits: Vec<DataSplit> = data_splits.to_vec();
@@ -211,7 +221,8 @@ impl DataEvolutionReader {
self.wide_file_read_type.clone(),
Vec::new(),
)
- .with_batch_size(self.batch_size);
+ .with_batch_size(self.batch_size)
+ .with_parquet_read_budget(self.parquet_read_budget.clone());
for split in splits {
let row_ranges = split.row_ranges().map(|r| r.to_vec());
@@ -509,7 +520,8 @@ impl DataEvolutionReader {
false,
None,
)?
- .with_batch_size(self.batch_size);
+ .with_batch_size(self.batch_size)
+ .with_parquet_read_budget(self.parquet_read_budget.clone());
let mut stream = prescan.read(splits)?;
let mut view_structs = HashSet::new();
while let Some(batch) = stream.next().await {
@@ -572,6 +584,7 @@ impl DataEvolutionReader {
let blob_descriptor_fields = self.blob_descriptor_fields.clone();
let blob_as_descriptor = self.blob_as_descriptor;
let batch_size = self.batch_size;
+ let parquet_read_budget = self.parquet_read_budget.clone();
let anchor_deletion_vector = anchor_deletion_vector.clone();
// Batch size for column-merge output. Matches the default Parquet
reader batch size.
const MERGE_BATCH_SIZE: usize = 1024;
@@ -629,6 +642,17 @@ impl DataEvolutionReader {
return;
}
+ // Column evolution advances every active source in lockstep. If
+ // several sources shared the row-group budget, the first source
+ // could occupy every permit while this loop waits for the next
+ // source's initial batch. Use the sequential Parquet path for
+ // multi-source merges; single-source evolution can still prefetch
+ // row groups concurrently.
+ let source_parquet_read_budget = if active_source_indices.len() ==
1 {
+ parquet_read_budget.clone()
+ } else {
+ None
+ };
let mut source_streams: Vec<Option<ArrowRecordBatchStream>> =
source_plan
.sources
.iter()
@@ -646,6 +670,7 @@ impl DataEvolutionReader {
table_fields.clone(),
batch_size,
blob_as_descriptor,
+ source_parquet_read_budget.clone(),
anchor_deletion_vector.as_ref(),
)
.map(Some)
@@ -1133,6 +1158,7 @@ fn open_source_stream(
table_fields: Vec<DataField>,
batch_size: Option<usize>,
blob_as_descriptor: bool,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
anchor_deletion_vector: Option<&DeletionVectorContext>,
) -> crate::Result<ArrowRecordBatchStream> {
let mut row_ranges = row_ranges;
@@ -1196,7 +1222,8 @@ fn open_source_stream(
Vec::new(),
)
.with_batch_size(batch_size)
- .with_blob_as_descriptor(blob_as_descriptor);
+ .with_blob_as_descriptor(blob_as_descriptor)
+ .with_parquet_read_budget(parquet_read_budget);
match source {
FieldSource::DataFile {
diff --git a/crates/paimon/src/table/data_file_reader.rs
b/crates/paimon/src/table/data_file_reader.rs
index b94016d1..591c88e9 100644
--- a/crates/paimon/src/table/data_file_reader.rs
+++ b/crates/paimon/src/table/data_file_reader.rs
@@ -16,8 +16,9 @@
// under the License.
use crate::arrow::build_target_arrow_schema;
-use crate::arrow::format::create_format_reader;
+use crate::arrow::format::create_format_reader_with_budget;
use crate::arrow::schema_evolution::{create_index_mapping, NULL_FIELD_INDEX};
+use crate::arrow::ParquetReadBudget;
use crate::deletion_vector::{DeletionVector, DeletionVectorFactory};
use crate::io::FileIO;
use crate::spec::{
@@ -46,6 +47,7 @@ pub(crate) struct DataFileReader {
row_filter_factory: Option<Arc<dyn crate::arrow::RowFilterFactory>>,
blob_as_descriptor: bool,
batch_size: Option<usize>,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
}
impl DataFileReader {
@@ -67,6 +69,7 @@ impl DataFileReader {
row_filter_factory: None,
blob_as_descriptor: false,
batch_size: None,
+ parquet_read_budget: None,
}
}
@@ -80,6 +83,14 @@ impl DataFileReader {
self
}
+ pub(crate) fn with_parquet_read_budget(
+ mut self,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
+ ) -> Self {
+ self.parquet_read_budget = parquet_read_budget;
+ self
+ }
+
pub(crate) fn with_row_filter_factory(
mut self,
factory: Arc<dyn crate::arrow::RowFilterFactory>,
@@ -279,6 +290,7 @@ impl DataFileReader {
let split = split.clone();
let blob_as_descriptor = self.blob_as_descriptor;
let batch_size = self.batch_size;
+ let parquet_read_budget = self.parquet_read_budget.clone();
let target_schema = build_target_arrow_schema(&read_type)?;
let file_fields = data_fields.clone().unwrap_or_else(||
table_fields.clone());
@@ -325,8 +337,12 @@ impl DataFileReader {
Ok(try_stream! {
let path_to_read = split.data_file_path(&file_meta);
- let format_reader =
- create_format_reader(&path_to_read, blob_as_descriptor,
&format_read_fields)?;
+ let format_reader = create_format_reader_with_budget(
+ &path_to_read,
+ blob_as_descriptor,
+ &format_read_fields,
+ parquet_read_budget,
+ )?;
let input_file = file_io.new_input(&path_to_read)?;
let file_reader = input_file.reader().await?;
let local_ranges = row_ranges.as_ref().map(|ranges| {
@@ -485,6 +501,7 @@ impl DataFileReader {
let file_io = self.file_io.clone();
let split = split.clone();
let blob_as_descriptor = self.blob_as_descriptor;
+ let parquet_read_budget = self.parquet_read_budget.clone();
let target_schema = build_target_arrow_schema(&read_type)?;
let file_fields = data_fields.clone().unwrap_or_else(||
table_fields.clone());
@@ -540,8 +557,12 @@ impl DataFileReader {
Ok(try_stream! {
let path_to_read = split.data_file_path(&file_meta);
- let format_reader =
- create_format_reader(&path_to_read, blob_as_descriptor,
&format_read_fields)?;
+ let format_reader = create_format_reader_with_budget(
+ &path_to_read,
+ blob_as_descriptor,
+ &format_read_fields,
+ parquet_read_budget,
+ )?;
let input_file = file_io.new_input(&path_to_read)?;
let file_reader = input_file.reader().await?;
diff --git a/crates/paimon/src/table/format_read_builder.rs
b/crates/paimon/src/table/format_read_builder.rs
index f97805c0..5d6c1b84 100644
--- a/crates/paimon/src/table/format_read_builder.rs
+++ b/crates/paimon/src/table/format_read_builder.rs
@@ -20,10 +20,13 @@
use super::partition_filter::PartitionFilter;
use super::read_builder::split_scan_predicates;
use super::read_builder::{resolve_projected_fields,
validate_projection_possible};
+use super::table_read::configured_parquet_read_budget;
use super::{Table, TableRead, TableScan};
+use crate::arrow::ParquetReadBudget;
use crate::spec::{DataField, Predicate};
use crate::table::source::RowRange;
use crate::Result;
+use std::sync::Arc;
#[derive(Debug, Clone)]
pub(crate) struct FormatReadBuilder<'a> {
@@ -37,6 +40,7 @@ pub(crate) struct FormatReadBuilder<'a> {
data_predicates: Vec<Predicate>,
limit: Option<usize>,
case_sensitive: bool,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
}
impl<'a> FormatReadBuilder<'a> {
@@ -49,6 +53,7 @@ impl<'a> FormatReadBuilder<'a> {
data_predicates: Vec::new(),
limit: None,
case_sensitive: true,
+ parquet_read_budget: None,
}
}
@@ -102,6 +107,11 @@ impl<'a> FormatReadBuilder<'a> {
self
}
+ pub(crate) fn with_parquet_read_budget(&mut self, budget:
Arc<ParquetReadBudget>) -> &mut Self {
+ self.parquet_read_budget = Some(budget);
+ self
+ }
+
pub(crate) fn new_scan(&self) -> TableScan<'a> {
TableScan::new(
self.table,
@@ -120,12 +130,17 @@ impl<'a> FormatReadBuilder<'a> {
None => self.table.schema().fields().to_vec(),
Some(fields) => fields,
};
+ let parquet_read_budget = match &self.parquet_read_budget {
+ Some(budget) => Arc::clone(budget),
+ None => configured_parquet_read_budget(self.table)?,
+ };
Ok(TableRead::new_format(
self.table,
read_type,
self.data_predicates.clone(),
self.limit,
- ))
+ )
+ .with_parquet_read_budget(parquet_read_budget))
}
/// Resolve the effective read type, deferring projection name resolution
to
diff --git a/crates/paimon/src/table/format_table_read.rs
b/crates/paimon/src/table/format_table_read.rs
index b9ea8e5c..9adb2e83 100644
--- a/crates/paimon/src/table/format_table_read.rs
+++ b/crates/paimon/src/table/format_table_read.rs
@@ -19,8 +19,9 @@
use super::data_file_reader::DataFileReader;
use super::read_builder::split_scan_predicates;
+use super::table_read::configured_parquet_read_budget;
use super::{ArrowRecordBatchStream, Table};
-use crate::arrow::{build_target_arrow_schema, paimon_type_to_arrow};
+use crate::arrow::{build_target_arrow_schema, paimon_type_to_arrow,
ParquetReadBudget};
use crate::spec::{extract_datum, BinaryRow, DataField, DataType, Datum,
Predicate};
use crate::{DataSplit, Error};
use arrow_array::{
@@ -39,6 +40,7 @@ pub(crate) struct FormatTableRead<'a> {
read_type: Vec<DataField>,
data_predicates: Vec<Predicate>,
row_filter_factory: Option<Arc<dyn crate::arrow::RowFilterFactory>>,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
limit: Option<usize>,
}
@@ -54,6 +56,7 @@ impl<'a> FormatTableRead<'a> {
read_type,
data_predicates,
row_filter_factory: None,
+ parquet_read_budget: None,
limit,
}
}
@@ -83,6 +86,18 @@ impl<'a> FormatTableRead<'a> {
self
}
+ pub(crate) fn with_parquet_read_budget(mut self, budget:
Arc<ParquetReadBudget>) -> Self {
+ self.parquet_read_budget = Some(budget);
+ self
+ }
+
+ fn parquet_read_budget(&self) -> crate::Result<Arc<ParquetReadBudget>> {
+ match &self.parquet_read_budget {
+ Some(budget) => Ok(Arc::clone(budget)),
+ None => configured_parquet_read_budget(self.table),
+ }
+ }
+
pub(crate) fn to_arrow(
&self,
data_splits: &[DataSplit],
@@ -105,6 +120,7 @@ impl<'a> FormatTableRead<'a> {
let mut remaining = self.limit;
let batch_size = Some(core_options.read_batch_size()?);
let row_filter_factory = self.row_filter_factory.clone();
+ let parquet_read_budget = Some(self.parquet_read_budget()?);
Ok(try_stream! {
for split in splits {
@@ -120,7 +136,8 @@ impl<'a> FormatTableRead<'a> {
data_read_type.clone(),
data_predicates.clone(),
)
- .with_batch_size(batch_size);
+ .with_batch_size(batch_size)
+ .with_parquet_read_budget(parquet_read_budget.clone());
if let Some(factory) = &row_filter_factory {
reader =
reader.with_row_filter_factory(Arc::clone(factory));
}
diff --git a/crates/paimon/src/table/kv_file_reader.rs
b/crates/paimon/src/table/kv_file_reader.rs
index 64ffce73..b1cae88f 100644
--- a/crates/paimon/src/table/kv_file_reader.rs
+++ b/crates/paimon/src/table/kv_file_reader.rs
@@ -30,7 +30,7 @@ use super::sort_merge::{
AggregateMergeFunction, DeduplicateMergeFunction,
PartialUpdateMergeFunction,
SortMergeReaderBuilder,
};
-use crate::arrow::build_target_arrow_schema;
+use crate::arrow::{build_target_arrow_schema, ParquetReadBudget};
use crate::io::FileIO;
use crate::spec::{
BigIntType, DataField, DataType as PaimonDataType, MergeEngine,
PartialUpdateConfig, Predicate,
@@ -45,6 +45,7 @@ use arrow_array::{RecordBatch, RecordBatchOptions};
use async_stream::try_stream;
use futures::StreamExt;
use std::collections::HashMap;
+use std::sync::Arc;
/// Reads primary-key table data files using sort-merge deduplication.
pub(crate) struct KeyValueFileReader {
@@ -77,6 +78,8 @@ pub(crate) struct KeyValueReadConfig {
pub merge_splits: bool,
/// Optional cap on file streams opened by a single sort-merge group.
pub max_merge_file_streams: Option<usize>,
+ /// Scan-shared Parquet concurrency and projected-byte budget.
+ pub parquet_read_budget: Option<Arc<ParquetReadBudget>>,
}
/// Keep only the conjuncts of `predicates` that reference primary-key columns,
@@ -408,6 +411,7 @@ impl KeyValueFileReader {
let sequence_fields = self.config.sequence_fields;
let read_batch_size = self.config.read_batch_size;
let max_merge_file_streams = self.config.max_merge_file_streams;
+ let parquet_read_budget = self.config.parquet_read_budget;
#[cfg(test)]
let input_batch_sizes = self.input_batch_sizes;
@@ -438,6 +442,17 @@ impl KeyValueFileReader {
continue;
}
ensure_merge_fan_in_limit(file_count, max_merge_file_streams)?;
+ // Sort-merge must first obtain one batch from every input
stream.
+ // A concurrent Parquet reader keeps its row-group permits
until
+ // the complete row group has been consumed, so enabling it on
+ // several lockstep inputs can let the first file occupy the
+ // entire scan budget while the merge waits for the second
file.
+ // Keep multi-file merge inputs on the sequential Parquet path.
+ let group_parquet_read_budget = if file_count == 1 {
+ parquet_read_budget.clone()
+ } else {
+ None
+ };
// Create one stream per data file.
let mut file_streams: Vec<ArrowRecordBatchStream> = Vec::new();
@@ -458,7 +473,8 @@ impl KeyValueFileReader {
internal_read_type.clone(),
pushdown_predicates.clone(),
)
- .with_batch_size(Some(read_batch_size));
+ .with_batch_size(Some(read_batch_size))
+
.with_parquet_read_budget(group_parquet_read_budget.clone());
let stream = reader.read_single_file_stream(
split,
@@ -578,9 +594,12 @@ mod tests {
use crate::table::source::DataSplitBuilder;
use crate::table::table_commit::TableCommit;
use crate::table::{Table, TableWrite};
- use arrow_array::{Array, Int32Array, StringArray};
+ use arrow_array::{Array, Int32Array, Int64Array, Int8Array, StringArray};
use arrow_schema::{DataType as ArrowDataType, Field as ArrowField, Schema
as ArrowSchema};
use futures::TryStreamExt;
+ use parquet::arrow::AsyncArrowWriter;
+ use parquet::file::metadata::ParquetMetaDataReader;
+ use parquet::file::properties::WriterProperties;
use std::sync::Arc;
fn test_file_io() -> FileIO {
@@ -749,6 +768,69 @@ mod tests {
}
}
+ async fn write_multi_row_group_kv_file(
+ file_io: &FileIO,
+ table_path: &str,
+ file_name: &str,
+ sequence: i64,
+ value: i32,
+ ) -> DataFileMeta {
+ let schema = crate::arrow::build_target_arrow_schema(&[
+ DataField::new(
+ SEQUENCE_NUMBER_FIELD_ID,
+ SEQUENCE_NUMBER_FIELD_NAME.to_string(),
+ DataType::BigInt(BigIntType::new()),
+ ),
+ DataField::new(
+ VALUE_KIND_FIELD_ID,
+ VALUE_KIND_FIELD_NAME.to_string(),
+ DataType::TinyInt(TinyIntType::new()),
+ ),
+ DataField::new(0, "id".to_string(), DataType::Int(IntType::new())),
+ DataField::new(1, "value".to_string(),
DataType::Int(IntType::new())),
+ ])
+ .unwrap();
+ let batch = RecordBatch::try_new(
+ schema.clone(),
+ vec![
+ Arc::new(Int64Array::from_value(sequence, 128)),
+ Arc::new(Int8Array::from_value(0, 128)),
+ Arc::new(Int32Array::from_iter_values(0..128)),
+ Arc::new(Int32Array::from_value(value, 128)),
+ ],
+ )
+ .unwrap();
+ let props = WriterProperties::builder()
+ .set_max_row_group_row_count(Some(64))
+ .set_dictionary_enabled(false)
+ .build();
+ let mut bytes = Vec::new();
+ {
+ let mut writer = AsyncArrowWriter::try_new(&mut bytes, schema,
Some(props)).unwrap();
+ writer.write(&batch).await.unwrap();
+ writer.close().await.unwrap();
+ }
+ let parquet_bytes = bytes::Bytes::from(bytes);
+ let metadata = ParquetMetaDataReader::new()
+ .parse_and_finish(&parquet_bytes)
+ .unwrap();
+ assert_eq!(metadata.num_row_groups(), 2);
+
+ let bucket_path = format!("{table_path}/bucket-0");
+ file_io.mkdirs(&format!("{bucket_path}/")).await.unwrap();
+ file_io
+ .new_output(&format!("{bucket_path}/{file_name}"))
+ .unwrap()
+ .write(parquet_bytes.clone())
+ .await
+ .unwrap();
+
+ let mut file = dummy_data_file(file_name.to_string());
+ file.file_size = parquet_bytes.len() as i64;
+ file.row_count = 128;
+ file
+ }
+
#[test]
fn retain_primary_key_conjuncts_semantics() {
let fields = vec![
@@ -863,6 +945,7 @@ mod tests {
read_batch_size: core_options.read_batch_size().unwrap(),
merge_splits: true,
max_merge_file_streams: Some(256),
+ parquet_read_budget: None,
},
);
@@ -926,6 +1009,7 @@ mod tests {
read_batch_size: core_options.read_batch_size().unwrap(),
merge_splits: false,
max_merge_file_streams: None,
+ parquet_read_budget: None,
},
)
.with_input_batch_sizes(input_batch_sizes.clone());
@@ -944,6 +1028,66 @@ mod tests {
assert_eq!(int_column(&batches, "value"), vec![11, 21, 31, 41, 51]);
}
+ #[tokio::test]
+ async fn kv_merge_with_shared_budget_does_not_deadlock_between_files() {
+ let file_io = test_file_io();
+ let table_path = "memory:/kv_shared_parquet_budget";
+ setup_dirs(&file_io, table_path).await;
+ let table = pk_table(
+ &file_io,
+ table_path,
+ &[
+ ("read.batch-size", "1"),
+ ("read.parquet.row-group.parallelism", "2"),
+ ],
+ );
+ let first =
+ write_multi_row_group_kv_file(&file_io, table_path,
"first.parquet", 0, 10).await;
+ let second =
+ write_multi_row_group_kv_file(&file_io, table_path,
"second.parquet", 1, 11).await;
+ let split = DataSplitBuilder::new()
+ .with_snapshot(1)
+ .with_partition(BinaryRow::new(0))
+ .with_bucket(0)
+ .with_bucket_path(format!("{table_path}/bucket-0"))
+ .with_total_buckets(1)
+ .with_data_files(vec![first, second])
+ .build()
+ .unwrap();
+ let core_options = table.schema().core_options();
+ let reader = KeyValueFileReader::new(
+ table.file_io().clone(),
+ KeyValueReadConfig {
+ table_name: table.identifier().full_name(),
+ table_options: table.schema().options().clone(),
+ schema_manager: table.schema_manager().clone(),
+ table_schema_id: table.schema().id(),
+ table_fields: table.schema().fields().to_vec(),
+ read_type: table.schema().fields().to_vec(),
+ predicates: Vec::new(),
+ primary_keys: table.schema().trimmed_primary_keys(),
+ merge_engine: core_options.merge_engine().unwrap(),
+ sequence_fields: Vec::new(),
+ read_batch_size: core_options.read_batch_size().unwrap(),
+ merge_splits: false,
+ max_merge_file_streams: None,
+ parquet_read_budget: Some(Arc::new(ParquetReadBudget::new(2,
256 << 20).unwrap())),
+ },
+ );
+ let batches = tokio::time::timeout(
+ std::time::Duration::from_secs(5),
+ reader.read(&[split]).unwrap().try_collect::<Vec<_>>(),
+ )
+ .await
+ .expect("multi-file sort-merge must not wait forever for a shared
Parquet permit")
+ .unwrap();
+
+ assert_eq!(
+ batches.iter().map(RecordBatch::num_rows).sum::<usize>(),
+ 128
+ );
+ }
+
/// Non-PK equality filter on a dedup PK table read through the sort-merge
/// path must return only matching rows. Before the post-merge residual,
/// the non-PK conjunct was silently dropped and all rows came back.
diff --git a/crates/paimon/src/table/read_builder.rs
b/crates/paimon/src/table/read_builder.rs
index de432a5c..94c89c2b 100644
--- a/crates/paimon/src/table/read_builder.rs
+++ b/crates/paimon/src/table/read_builder.rs
@@ -24,12 +24,13 @@ use super::bucket_filter::{extract_predicate_for_keys,
split_partition_and_data_
use super::format_read_builder::FormatReadBuilder;
use super::incremental_scan::{IncrementalScan, IncrementalScanMode};
use super::partition_filter::PartitionFilter;
-use super::table_read::TableRead;
+use super::table_read::{configured_parquet_read_budget, TableRead};
use super::{Table, TableScan};
use crate::spec::{CoreOptions, DataField, Predicate};
use crate::table::source::RowRange;
use crate::{Error, Result};
use std::collections::{HashMap, HashSet};
+use std::sync::Arc;
#[derive(Debug, Clone, Default)]
struct NormalizedFilter {
@@ -235,6 +236,23 @@ impl<'a> ReadBuilder<'a> {
self
}
+ /// Inject a Parquet budget shared with sibling scan partitions.
+ #[doc(hidden)]
+ pub fn with_parquet_read_budget(
+ &mut self,
+ budget: Arc<crate::arrow::ParquetReadBudget>,
+ ) -> &mut Self {
+ match &mut self.0 {
+ ReadBuilderKind::Paimon(builder) => {
+ builder.with_parquet_read_budget(budget);
+ }
+ ReadBuilderKind::Format(builder) => {
+ builder.with_parquet_read_budget(budget);
+ }
+ }
+ self
+ }
+
/// Create a table scan. Call [TableScan::plan] to get splits.
pub fn new_scan(&self) -> TableScan<'a> {
match &self.0 {
@@ -291,6 +309,7 @@ struct PaimonReadBuilder<'a> {
limit: Option<usize>,
row_ranges: Option<Vec<RowRange>>,
case_sensitive: bool,
+ parquet_read_budget: Option<Arc<crate::arrow::ParquetReadBudget>>,
}
impl<'a> PaimonReadBuilder<'a> {
@@ -303,6 +322,7 @@ impl<'a> PaimonReadBuilder<'a> {
limit: None,
row_ranges: None,
case_sensitive: true,
+ parquet_read_budget: None,
}
}
@@ -420,6 +440,14 @@ impl<'a> PaimonReadBuilder<'a> {
self
}
+ fn with_parquet_read_budget(
+ &mut self,
+ budget: Arc<crate::arrow::ParquetReadBudget>,
+ ) -> &mut Self {
+ self.parquet_read_budget = Some(budget);
+ self
+ }
+
/// Create a table scan. Call [TableScan::plan] to get splits.
///
/// Projection names are resolved here on a best-effort basis: the resolved
@@ -464,11 +492,14 @@ impl<'a> PaimonReadBuilder<'a> {
// Pass the FULL data predicate through (including `And`/`Or`/`Not`).
// Pushdown/stats skip compound nodes; the residual pass enforces the
full
// predicate exactly. Pruning here would drop compound predicates.
- Ok(TableRead::new(
- self.table,
- read_type,
- self.filter.data_predicates.clone(),
- ))
+ let parquet_read_budget = match &self.parquet_read_budget {
+ Some(budget) => Arc::clone(budget),
+ None => configured_parquet_read_budget(self.table)?,
+ };
+ Ok(
+ TableRead::new(self.table, read_type,
self.filter.data_predicates.clone())
+ .with_parquet_read_budget(parquet_read_budget),
+ )
}
/// Resolve the effective read type, deferring projection name resolution
to
diff --git a/crates/paimon/src/table/table_read.rs
b/crates/paimon/src/table/table_read.rs
index 0a179b49..7e7d3c3a 100644
--- a/crates/paimon/src/table/table_read.rs
+++ b/crates/paimon/src/table/table_read.rs
@@ -23,6 +23,7 @@ use super::kv_file_reader::{KeyValueFileReader,
KeyValueReadConfig};
use super::read_builder::split_scan_predicates;
use super::{ArrowRecordBatchStream, Table};
use crate::arrow::build_target_arrow_schema;
+use crate::arrow::ParquetReadBudget;
use crate::spec::{
BigIntType, CoreOptions, DataField, DataType, MergeEngine, Predicate,
TinyIntType,
ROW_KIND_FIELD_ID, ROW_KIND_FIELD_NAME, SEQUENCE_NUMBER_FIELD_ID,
SEQUENCE_NUMBER_FIELD_NAME,
@@ -52,6 +53,16 @@ enum TableReadKind<'a> {
Format(FormatTableRead<'a>),
}
+pub(super) fn configured_parquet_read_budget(
+ table: &Table,
+) -> crate::Result<Arc<ParquetReadBudget>> {
+ let options = table.schema().core_options();
+ Ok(Arc::new(ParquetReadBudget::new(
+ options.parquet_row_group_parallelism()?,
+ options.parquet_row_group_max_inflight_bytes()?,
+ )?))
+}
+
impl<'a> TableRead<'a> {
/// Create a new TableRead with a specific read type (projected fields).
pub fn new(
@@ -132,6 +143,19 @@ impl<'a> TableRead<'a> {
}
}
+ /// Override the Parquet resource budget shared by this read.
+ #[doc(hidden)]
+ pub fn with_parquet_read_budget(self, budget: Arc<ParquetReadBudget>) ->
Self {
+ match self.0 {
+ TableReadKind::Paimon(read) => {
+
Self(TableReadKind::Paimon(read.with_parquet_read_budget(budget)))
+ }
+ TableReadKind::Format(read) => {
+
Self(TableReadKind::Format(read.with_parquet_read_budget(budget)))
+ }
+ }
+ }
+
/// Returns an [`ArrowRecordBatchStream`].
pub fn to_arrow(&self, data_splits: &[DataSplit]) ->
crate::Result<ArrowRecordBatchStream> {
match &self.0 {
@@ -189,6 +213,7 @@ struct PaimonTableRead<'a> {
read_type: Vec<DataField>,
data_predicates: Vec<Predicate>,
row_filter_factory: Option<Arc<dyn crate::arrow::RowFilterFactory>>,
+ parquet_read_budget: Option<Arc<ParquetReadBudget>>,
}
impl<'a> PaimonTableRead<'a> {
@@ -203,6 +228,7 @@ impl<'a> PaimonTableRead<'a> {
read_type,
data_predicates,
row_filter_factory: None,
+ parquet_read_budget: None,
}
}
@@ -241,6 +267,18 @@ impl<'a> PaimonTableRead<'a> {
self
}
+ fn with_parquet_read_budget(mut self, budget: Arc<ParquetReadBudget>) ->
Self {
+ self.parquet_read_budget = Some(budget);
+ self
+ }
+
+ fn parquet_read_budget(&self) -> crate::Result<Arc<ParquetReadBudget>> {
+ match &self.parquet_read_budget {
+ Some(budget) => Ok(Arc::clone(budget)),
+ None => configured_parquet_read_budget(self.table),
+ }
+ }
+
/// Returns an [`ArrowRecordBatchStream`] for an incremental scan plan.
pub fn to_incremental_arrow(
&self,
@@ -276,15 +314,17 @@ impl<'a> PaimonTableRead<'a> {
let table = self.table.clone();
let read_type = self.read_type.clone();
let data_predicates = self.data_predicates.clone();
+ let parquet_read_budget = self.parquet_read_budget()?;
Ok(Box::pin(async_stream::try_stream! {
let mut workers = stream::iter(pairs.into_iter().map(|(before,
after)| {
let table = table.clone();
let read_type = read_type.clone();
let data_predicates = data_predicates.clone();
+ let parquet_read_budget = Arc::clone(&parquet_read_budget);
let worker: ArrowRecordBatchStream =
Box::pin(async_stream::try_stream! {
- let pair_read =
- PaimonTableRead::new(&table, read_type,
data_predicates);
+ let pair_read = PaimonTableRead::new(&table, read_type,
data_predicates)
+ .with_parquet_read_budget(parquet_read_budget);
let mut pair_stream =
pair_read.to_diff_after_image_stream(&before, &after)?;
while let Some(batch) = pair_stream.next().await {
yield batch?;
@@ -356,7 +396,8 @@ impl<'a> PaimonTableRead<'a> {
read_type,
self.data_predicates.clone(),
)
-
.with_batch_size(Some(self.table.schema().core_options().read_batch_size()?));
+
.with_batch_size(Some(self.table.schema().core_options().read_batch_size()?))
+ .with_parquet_read_budget(Some(self.parquet_read_budget()?));
let raw_stream = reader.read(&data_splits)?;
Ok(Box::pin(async_stream::try_stream! {
@@ -414,14 +455,17 @@ impl<'a> PaimonTableRead<'a> {
let table = self.table.clone();
let read_type = self.read_type.clone();
let data_predicates = self.data_predicates.clone();
+ let parquet_read_budget = self.parquet_read_budget()?;
Ok(Box::pin(async_stream::try_stream! {
let mut workers = stream::iter(pairs.into_iter().map(|(before,
after)| {
let table = table.clone();
let read_type = read_type.clone();
let data_predicates = data_predicates.clone();
+ let parquet_read_budget = Arc::clone(&parquet_read_budget);
let worker: ArrowRecordBatchStream =
Box::pin(async_stream::try_stream! {
- let pair_read = PaimonTableRead::new(&table, read_type,
data_predicates);
+ let pair_read = PaimonTableRead::new(&table, read_type,
data_predicates)
+ .with_parquet_read_budget(parquet_read_budget);
let mut pair_stream =
pair_read.to_audit_log_arrow_for_diff(&before,
&after)?;
while let Some(batch) = pair_stream.next().await {
@@ -466,10 +510,12 @@ impl<'a> PaimonTableRead<'a> {
let table = self.table.clone();
let read_type_for_output = self.read_type.clone();
let data_predicates = self.data_predicates.clone();
+ let parquet_read_budget = self.parquet_read_budget()?;
Ok(Box::pin(async_stream::try_stream! {
let core_options = CoreOptions::new(table.schema().options());
- let pair_read = PaimonTableRead::new(&table,
diff_read_type.clone(), data_predicates);
+ let pair_read = PaimonTableRead::new(&table,
diff_read_type.clone(), data_predicates)
+ .with_parquet_read_budget(parquet_read_budget);
let before_stream =
pair_read.read_pk_sorted_for_diff_with_type(&before,
&core_options, &diff_read_type)?;
let after_stream =
@@ -550,10 +596,12 @@ impl<'a> PaimonTableRead<'a> {
let data_predicates = self.data_predicates.clone();
let before = before.to_vec();
let after = after.to_vec();
+ let parquet_read_budget = self.parquet_read_budget()?;
Ok(Box::pin(async_stream::try_stream! {
let core_options = CoreOptions::new(table.schema().options());
- let pair_read = PaimonTableRead::new(&table,
diff_read_type.clone(), data_predicates);
+ let pair_read = PaimonTableRead::new(&table,
diff_read_type.clone(), data_predicates)
+ .with_parquet_read_budget(parquet_read_budget);
let before_stream = pair_read.read_pk_sorted_for_diff_with_type(
&before,
&core_options,
@@ -637,6 +685,10 @@ impl<'a> PaimonTableRead<'a> {
read_batch_size: core_options.read_batch_size()?,
merge_splits: true,
max_merge_file_streams: Some(256),
+ // Diff primes the before and after streams in sequence.
Keeping
+ // a row-group permit across yielded batches can otherwise let
+ // the first side block the second side indefinitely.
+ parquet_read_budget: None,
},
);
reader.read(splits)
@@ -773,6 +825,7 @@ impl<'a> PaimonTableRead<'a> {
read_batch_size: core_options.read_batch_size()?,
merge_splits: false,
max_merge_file_streams: None,
+ parquet_read_budget: Some(self.parquet_read_budget()?),
},
);
reader.read(splits)
@@ -797,7 +850,8 @@ impl<'a> PaimonTableRead<'a> {
core_options.blob_view_resolve_enabled(),
self.table.rest_env().cloned(),
)?
- .with_batch_size(Some(core_options.read_batch_size()?));
+ .with_batch_size(Some(core_options.read_batch_size()?))
+ .with_parquet_read_budget(Some(self.parquet_read_budget()?));
reader.read(data_splits)
}
@@ -815,7 +869,8 @@ impl<'a> PaimonTableRead<'a> {
self.read_type().to_vec(),
self.data_predicates.clone(),
)
-
.with_batch_size(Some(self.table.schema().core_options().read_batch_size()?));
+
.with_batch_size(Some(self.table.schema().core_options().read_batch_size()?))
+ .with_parquet_read_budget(Some(self.parquet_read_budget()?));
// The engine decoder filter is safe only on the plain append/raw path.
// This constructor is also used by raw-convertible primary-key splits,
// where positional merge semantics must remain untouched.
@@ -1401,8 +1456,10 @@ fn pk_split_needs_merge(split: &DataSplit, dv_enabled:
bool) -> bool {
#[cfg(test)]
mod tests {
use super::*;
+ use crate::catalog::Identifier;
+ use crate::io::FileIOBuilder;
use crate::spec::stats::BinaryTableStats;
- use crate::spec::{BinaryRow, DataFileMeta};
+ use crate::spec::{BinaryRow, DataFileMeta, DataType, IntType, Schema,
TableSchema};
use crate::table::query_auth_table;
use crate::table::source::DataSplitBuilder;
@@ -1444,6 +1501,22 @@ mod tests {
.unwrap()
}
+ fn table_with_invalid_parquet_budget(format_table: bool) -> Table {
+ let mut schema = Schema::builder()
+ .column("id", DataType::Int(IntType::new()))
+ .option("read.parquet.row-group.parallelism", "0");
+ if format_table {
+ schema = schema.option("type", "format-table");
+ }
+ Table::new(
+ FileIOBuilder::new("memory").build().unwrap(),
+ Identifier::new("default", "budget_t"),
+ "memory:/budget_t".to_string(),
+ TableSchema::new(0, &schema.build().unwrap()),
+ None,
+ )
+ }
+
#[test]
fn test_pk_split_needs_merge_routing() {
// Raw convertible with known delete counts: raw read.
@@ -1500,6 +1573,23 @@ mod tests {
);
}
+ #[test]
+ fn test_direct_table_read_validates_and_can_override_parquet_budget() {
+ for format_table in [false, true] {
+ let table = table_with_invalid_parquet_budget(format_table);
+ let read = TableRead::new(&table, table.schema.fields().to_vec(),
Vec::new());
+ assert!(matches!(
+ read.to_arrow(&[]),
+ Err(crate::Error::DataInvalid { ref message, .. })
+ if message.contains("row-group.parallelism")
+ ));
+
+ let read = TableRead::new(&table, table.schema.fields().to_vec(),
Vec::new())
+
.with_parquet_read_budget(Arc::new(ParquetReadBudget::default()));
+ assert!(read.to_arrow(&[]).is_ok());
+ }
+ }
+
#[test]
fn test_direct_incremental_read_fails_closed_when_query_auth_enabled() {
let table = query_auth_table();