Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 15 additions & 3 deletions crates/paimon/src/table/data_evolution_reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ mod blob_fallback;
use super::blob_resolver::{BlobReadLimiter, BLOB_DESCRIPTOR_READ_CONCURRENCY};
use super::data_file_reader::{
append_null_row_id_column, attach_row_id, expand_selected_row_ids, insert_column_at,
DataFileReader,
DataFileReadTiming, DataFileReader,
};
use crate::arrow::format::FilePredicates;
use crate::arrow::{build_target_arrow_schema, ParquetReadBudget};
Expand Down Expand Up @@ -114,6 +114,7 @@ pub(crate) struct DataEvolutionReader {
blob_read_limiter: BlobReadLimiter,
batch_size: Option<usize>,
parquet_read_budget: Option<Arc<ParquetReadBudget>>,
read_timing: Option<Arc<DataFileReadTiming>>,
}

impl DataEvolutionReader {
Expand Down Expand Up @@ -191,6 +192,7 @@ impl DataEvolutionReader {
blob_read_limiter: BlobReadLimiter::new(),
batch_size: None,
parquet_read_budget: None,
read_timing: None,
})
}

Expand All @@ -207,6 +209,11 @@ impl DataEvolutionReader {
self
}

pub(crate) fn with_read_timing(mut self, read_timing: Option<Arc<DataFileReadTiming>>) -> Self {
self.read_timing = read_timing;
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();
Expand Down Expand Up @@ -248,7 +255,8 @@ impl DataEvolutionReader {
},
)
.with_batch_size(self.batch_size)
.with_parquet_read_budget(self.parquet_read_budget.clone());
.with_parquet_read_budget(self.parquet_read_budget.clone())
.with_read_timing(self.read_timing.clone());

for split in splits {
let row_ranges = split.row_ranges().map(|r| r.to_vec());
Expand Down Expand Up @@ -611,6 +619,7 @@ impl DataEvolutionReader {
let blob_as_descriptor = self.blob_as_descriptor;
let batch_size = self.batch_size;
let parquet_read_budget = self.parquet_read_budget.clone();
let read_timing = self.read_timing.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;
Expand Down Expand Up @@ -697,6 +706,7 @@ impl DataEvolutionReader {
batch_size,
blob_as_descriptor,
source_parquet_read_budget.clone(),
read_timing.clone(),
anchor_deletion_vector.as_ref(),
)
.map(Some)
Expand Down Expand Up @@ -1228,6 +1238,7 @@ fn open_source_stream(
batch_size: Option<usize>,
blob_as_descriptor: bool,
parquet_read_budget: Option<Arc<ParquetReadBudget>>,
read_timing: Option<Arc<DataFileReadTiming>>,
anchor_deletion_vector: Option<&DeletionVectorContext>,
) -> crate::Result<ArrowRecordBatchStream> {
let mut row_ranges = row_ranges;
Expand Down Expand Up @@ -1292,7 +1303,8 @@ fn open_source_stream(
)
.with_batch_size(batch_size)
.with_blob_as_descriptor(blob_as_descriptor)
.with_parquet_read_budget(parquet_read_budget);
.with_parquet_read_budget(parquet_read_budget)
.with_read_timing(read_timing);

match source {
FieldSource::DataFile {
Expand Down
86 changes: 83 additions & 3 deletions crates/paimon/src/table/data_file_reader.rs
Original file line number Diff line number Diff line change
Expand Up @@ -20,7 +20,7 @@ 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::io::{FileIO, FileRead};
use crate::spec::{
is_variant_extraction_row_type, DataField, DataFileMeta, DataType, Predicate, ROW_ID_FIELD_NAME,
};
Expand All @@ -33,7 +33,51 @@ use arrow_cast::cast;

use async_stream::try_stream;
use futures::StreamExt;
use std::ops::Range;
use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use std::time::{Duration, Instant};

#[derive(Debug, Default)]
pub(crate) struct DataFileReadTiming {
file_read_nanos: AtomicU64,
parquet_decode_nanos: AtomicU64,
}

impl DataFileReadTiming {
fn add_file_read(&self, duration: Duration) {
self.file_read_nanos
.fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
}

fn add_parquet_decode(&self, duration: Duration) {
self.parquet_decode_nanos
.fetch_add(duration.as_nanos() as u64, Ordering::Relaxed);
}

pub(crate) fn file_read(&self) -> Duration {
Duration::from_nanos(self.file_read_nanos.load(Ordering::Relaxed))
}

pub(crate) fn parquet_decode(&self) -> Duration {
Duration::from_nanos(self.parquet_decode_nanos.load(Ordering::Relaxed))
}
}

struct TimedFileRead {
inner: Box<dyn FileRead>,
timing: Arc<DataFileReadTiming>,
}

#[async_trait::async_trait]
impl FileRead for TimedFileRead {
async fn read(&self, range: Range<u64>) -> crate::Result<bytes::Bytes> {
let start = Instant::now();
let result = self.inner.read(range).await;
self.timing.add_file_read(start.elapsed());
result
}
}

/// Reads data from Parquet files.
#[derive(Clone)]
Expand All @@ -48,6 +92,7 @@ pub(crate) struct DataFileReader {
blob_as_descriptor: bool,
batch_size: Option<usize>,
parquet_read_budget: Option<Arc<ParquetReadBudget>>,
read_timing: Option<Arc<DataFileReadTiming>>,
}

impl DataFileReader {
Expand All @@ -70,6 +115,7 @@ impl DataFileReader {
blob_as_descriptor: false,
batch_size: None,
parquet_read_budget: None,
read_timing: None,
}
}

Expand All @@ -91,6 +137,11 @@ impl DataFileReader {
self
}

pub(crate) fn with_read_timing(mut self, read_timing: Option<Arc<DataFileReadTiming>>) -> Self {
self.read_timing = read_timing;
self
}

pub(crate) fn with_row_filter_factory(
mut self,
factory: Arc<dyn crate::arrow::RowFilterFactory>,
Expand Down Expand Up @@ -291,6 +342,7 @@ impl DataFileReader {
let blob_as_descriptor = self.blob_as_descriptor;
let batch_size = self.batch_size;
let parquet_read_budget = self.parquet_read_budget.clone();
let read_timing = self.read_timing.clone();

let target_schema = build_target_arrow_schema(&read_type)?;
let file_fields = data_fields.clone().unwrap_or_else(|| table_fields.clone());
Expand Down Expand Up @@ -344,7 +396,19 @@ impl DataFileReader {
parquet_read_budget,
)?;
let input_file = file_io.new_input(&path_to_read)?;
let open_start = read_timing.as_ref().map(|_| Instant::now());
let file_reader = input_file.reader().await?;
if let (Some(timing), Some(start)) = (read_timing.as_ref(), open_start) {
timing.add_file_read(start.elapsed());
}
let file_reader: Box<dyn FileRead> = match read_timing.as_ref() {
Some(timing) => Box::new(TimedFileRead {
inner: Box::new(file_reader),
timing: Arc::clone(timing),
}),
None => Box::new(file_reader),
};
let is_parquet = path_to_read.to_ascii_lowercase().ends_with(".parquet");
let local_ranges = row_ranges.as_ref().map(|ranges| {
to_local_row_ranges(ranges, file_meta.first_row_id.unwrap_or(0), file_meta.row_count)
});
Expand All @@ -364,15 +428,31 @@ impl DataFileReader {
let mut row_id_offset = 0usize;

let mut batch_stream = format_reader.read_batch_stream(
Box::new(file_reader),
file_reader,
file_meta.file_size as u64,
&format_read_fields,
file_predicates.as_ref(),
batch_size,
row_selection,
).await?;

while let Some(batch) = batch_stream.next().await {
loop {
let batch = if is_parquet {
if let Some(timing) = read_timing.as_ref() {
std::future::poll_fn(|cx| {
let start = Instant::now();
let batch = batch_stream.as_mut().poll_next(cx);
timing.add_parquet_decode(start.elapsed());
batch
})
.await
} else {
batch_stream.next().await
}
} else {
batch_stream.next().await
};
let Some(batch) = batch else { break };
let batch = batch?;
let num_rows = batch.num_rows();
let batch_schema = batch.schema();
Expand Down
24 changes: 21 additions & 3 deletions crates/paimon/src/table/table_read.rs
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@
// under the License.

use super::data_evolution_reader::DataEvolutionReader;
use super::data_file_reader::DataFileReader;
use super::data_file_reader::{DataFileReadTiming, DataFileReader};
use super::format_table_read::FormatTableRead;
use super::incremental_scan::{IncrementalPlan, IncrementalScanMode, IncrementalSplit};
use super::kv_file_reader::{KeyValueFileReader, KeyValueReadConfig};
Expand Down Expand Up @@ -158,6 +158,15 @@ impl<'a> TableRead<'a> {
}
}

pub(crate) fn with_data_file_read_timing(self, timing: Arc<DataFileReadTiming>) -> Self {
match self.0 {
TableReadKind::Paimon(read) => Self(TableReadKind::Paimon(
read.with_data_file_read_timing(timing),
)),
TableReadKind::Format(read) => Self(TableReadKind::Format(read)),
}
}

/// Returns an [`ArrowRecordBatchStream`].
pub fn to_arrow(&self, data_splits: &[DataSplit]) -> crate::Result<ArrowRecordBatchStream> {
match &self.0 {
Expand Down Expand Up @@ -216,6 +225,7 @@ struct PaimonTableRead<'a> {
data_predicates: Vec<Predicate>,
row_filter_factory: Option<Arc<dyn crate::arrow::RowFilterFactory>>,
parquet_read_budget: Option<Arc<ParquetReadBudget>>,
data_file_read_timing: Option<Arc<DataFileReadTiming>>,
}

impl<'a> PaimonTableRead<'a> {
Expand All @@ -231,6 +241,7 @@ impl<'a> PaimonTableRead<'a> {
data_predicates,
row_filter_factory: None,
parquet_read_budget: None,
data_file_read_timing: None,
}
}

Expand Down Expand Up @@ -274,6 +285,11 @@ impl<'a> PaimonTableRead<'a> {
self
}

fn with_data_file_read_timing(mut self, timing: Arc<DataFileReadTiming>) -> Self {
self.data_file_read_timing = Some(timing);
self
}

fn parquet_read_budget(&self) -> crate::Result<Arc<ParquetReadBudget>> {
match &self.parquet_read_budget {
Some(budget) => Ok(Arc::clone(budget)),
Expand Down Expand Up @@ -856,7 +872,8 @@ impl<'a> PaimonTableRead<'a> {
self.table.rest_env().cloned(),
)?
.with_batch_size(Some(core_options.read_batch_size()?))
.with_parquet_read_budget(Some(self.parquet_read_budget()?));
.with_parquet_read_budget(Some(self.parquet_read_budget()?))
.with_read_timing(self.data_file_read_timing.clone());
reader.read(data_splits)
}

Expand All @@ -875,7 +892,8 @@ impl<'a> PaimonTableRead<'a> {
self.data_predicates.clone(),
)
.with_batch_size(Some(self.table.schema().core_options().read_batch_size()?))
.with_parquet_read_budget(Some(self.parquet_read_budget()?));
.with_parquet_read_budget(Some(self.parquet_read_budget()?))
.with_read_timing(self.data_file_read_timing.clone());
// 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.
Expand Down
Loading
Loading