use super::*;
pub struct CollectRequest {
pub lf: LazyFrame,
pub polars_streaming: bool,
pub buffer_start: usize,
pub buffer_end: usize,
pub plan: FillPlan,
}
pub struct FillPlan {
buffer_start: usize,
buffer_end: usize,
num_rows: usize,
count_known: bool,
indexing: bool,
held: Option<(DataFrame, usize)>,
view_start: usize,
view_len: usize,
max_rows: usize,
max_mb: usize,
}
impl FillPlan {
pub fn fit(mut self, df: DataFrame) -> CollectResult {
let returned = df.height();
let bytes_per_row = (returned > 0).then(|| (df.estimated_size() / returned).max(1));
let (df, start, seam) = match self.held.take() {
Some((mut held, held_start)) if held_start + held.height() == self.buffer_start => {
let seam = held.height();
match held.vstack_mut(&df) {
Ok(_) => (held, held_start, Some(seam)),
Err(_) => (df, self.buffer_start, None),
}
}
Some((held, held_start))
if returned > 0 && self.buffer_start + returned == held_start =>
{
match df.vstack(&held) {
Ok(joined) => (joined, self.buffer_start, Some(returned)),
Err(_) => (df, self.buffer_start, None),
}
}
_ => (df, self.buffer_start, None),
};
let (df, start) = self.cut_to_caps(df, start, seam);
CollectResult {
df,
start,
returned,
bytes_per_row,
buffer_start: self.buffer_start,
buffer_end: self.buffer_end,
num_rows: self.num_rows,
count_known: self.count_known,
indexing: self.indexing,
}
}
fn cut_to_caps(&self, df: DataFrame, start: usize, seam: Option<usize>) -> (DataFrame, usize) {
let total = df.height();
if total == 0 {
return (df, start);
}
let mut max_rows = total;
if self.max_rows > 0 {
max_rows = max_rows.min(self.max_rows);
}
if self.max_mb > 0 {
let bytes_per_row = (df.estimated_size() / total).max(1);
max_rows = max_rows.min(self.max_mb * 1024 * 1024 / bytes_per_row);
}
let max_rows = max_rows.max(1);
if max_rows >= total {
return (df, start);
}
let view_off = self.view_start.saturating_sub(start).min(total);
let view_len = self.view_len.max(1).min(total);
let view_center = view_off + view_len / 2;
let mut keep_start = view_center.saturating_sub(max_rows / 2);
if keep_start + max_rows > total {
keep_start = total - max_rows;
}
let kept = max_rows.min(total - keep_start);
(trim_rows(df, keep_start, kept, seam), start + keep_start)
}
}
pub struct CollectResult {
pub(super) df: DataFrame,
pub(super) start: usize,
returned: usize,
bytes_per_row: Option<usize>,
buffer_start: usize,
buffer_end: usize,
num_rows: usize,
count_known: bool,
indexing: bool,
}
impl CollectResult {
pub(crate) fn rows(&self) -> &DataFrame {
&self.df
}
}
pub(super) const STRING_BYTES_GUESS: usize = 40;
pub(super) fn estimate_bytes_per_row(
schema: &Schema,
columns: &[String],
column_bytes: &[(String, usize)],
) -> usize {
let footer_width = |name: &String| {
column_bytes
.iter()
.find(|(n, _)| n == name)
.map(|(_, w)| *w)
};
columns
.iter()
.map(|name| match schema.get(name.as_str()) {
Some(DataType::String) => 16 + footer_width(name).unwrap_or(STRING_BYTES_GUESS - 16),
Some(DataType::Binary) => 16 + binary_stub().len(),
Some(DataType::Boolean) => 1,
Some(DataType::Null) => 0,
Some(dtype) if dtype.is_primitive_numeric() || dtype.is_temporal() => {
match dtype.to_physical() {
DataType::Int8 | DataType::UInt8 => 1,
DataType::Int16 | DataType::UInt16 => 2,
DataType::Int32 | DataType::UInt32 | DataType::Float32 => 4,
DataType::Int128 => 16,
_ => 8,
}
}
Some(DataType::Decimal(..)) => 16,
_ => footer_width(name).unwrap_or(64),
})
.sum::<usize>()
.max(1)
}
pub(super) fn trim_rows(
df: DataFrame,
offset: usize,
len: usize,
seam: Option<usize>,
) -> DataFrame {
if backing_rows(&df, offset, len) > len + len / 4 {
compact_rows(df, offset, len, seam)
} else {
df.slice(offset as i64, len)
}
}
pub(super) fn backing_rows(df: &DataFrame, offset: usize, len: usize) -> usize {
let end = offset + len;
df.columns()
.iter()
.filter_map(Column::as_series)
.map(|s| {
let mut start = 0;
let mut touched = 0;
for chunk in s.chunks() {
let chunk_end = start + chunk.len();
if start < end && offset < chunk_end {
touched += chunk.len();
}
start = chunk_end;
}
touched
})
.max()
.unwrap_or(len)
}
pub(super) fn compact_rows(
df: DataFrame,
offset: usize,
len: usize,
seam: Option<usize>,
) -> DataFrame {
use polars::series::builder::SeriesBuilder;
use polars_arrow::array::builder::ShareStrategy;
#[cfg(test)]
tests::COMPACTIONS.with(|count| count.set(count.get() + 1));
let len = len.min(df.height().saturating_sub(offset));
let pieces = match seam.filter(|&seam| offset < seam && seam < offset + len) {
Some(seam) => vec![(offset, seam - offset), (seam, offset + len - seam)],
None => vec![(offset, len)],
};
let copy = |series: &Series, (offset, len): (usize, usize)| {
let mut builder = SeriesBuilder::new(series.dtype().clone());
builder.reserve(len);
builder.subslice_extend(series, offset, len, ShareStrategy::Never);
builder.freeze(series.name().clone())
};
let columns = df
.into_columns()
.into_iter()
.map(|column| match column {
Column::Scalar(constant) => {
Column::new_scalar(constant.name().clone(), constant.scalar().clone(), len)
}
Column::Series(series) => {
let mut kept = copy(&series, pieces[0]);
for &piece in &pieces[1..] {
if kept.append_owned(copy(&series, piece)).is_err() {
kept = copy(&series, (offset, len));
break;
}
}
kept.into_column()
}
})
.collect();
DataFrame::new(len, columns).unwrap_or_else(|_| DataFrame::empty_with_height(len))
}
pub(super) fn shrink_around_view(
view_start: usize,
view_end: usize,
max_len: usize,
floor: usize,
ceil: usize,
buffer_start: &mut usize,
buffer_end: &mut usize,
) {
if buffer_end.saturating_sub(*buffer_start) <= max_len {
return;
}
let view_len = view_end.saturating_sub(view_start);
if view_len >= max_len {
*buffer_start = view_start;
*buffer_end = (view_start + max_len).min(ceil);
return;
}
let half = (max_len - view_len) / 2;
*buffer_end = (view_end + half).min(ceil);
*buffer_start = buffer_end.saturating_sub(max_len).max(floor);
if *buffer_start > view_start {
*buffer_start = view_start;
}
*buffer_end = (*buffer_start + max_len).min(ceil);
}
pub(super) const MAX_FILES_PER_BUFFER: usize = 16;
pub(super) fn limit_files(
offsets: &[usize],
view_start: usize,
view_end: usize,
start: usize,
end: usize,
max_files: usize,
) -> (usize, usize) {
let (Some((first, last)), Some((view_first, view_last))) = (
files_holding(offsets, start, end.saturating_sub(start)),
files_holding(
offsets,
view_start,
view_end.saturating_sub(view_start).max(1),
),
) else {
return (start, end);
};
let opened = |from: usize, to: usize| (from..=to).filter(|&i| holds_rows(offsets, i)).count();
if opened(first, last) <= max_files {
return (start, end);
}
let (mut lo, mut hi) = (view_first.max(first), view_last.min(last));
let mut files = opened(lo, hi);
while files < max_files && (hi < last || lo > first) {
if hi < last {
hi += 1;
files += usize::from(holds_rows(offsets, hi));
}
if files < max_files && lo > first {
lo -= 1;
files += usize::from(holds_rows(offsets, lo));
}
}
(start.max(offsets[lo]), end.min(offsets[hi + 1]))
}
pub(super) fn holds_rows(offsets: &[usize], i: usize) -> bool {
offsets[i + 1] > offsets[i]
}
pub(super) fn files_with_rows(offsets: &[usize], first: usize, last: usize) -> Vec<usize> {
(first..=last).filter(|&i| holds_rows(offsets, i)).collect()
}
pub(super) fn files_holding(offsets: &[usize], start: usize, len: usize) -> Option<(usize, usize)> {
let files = offsets.len().checked_sub(1)?;
let total = *offsets.last()?;
if files == 0 || len == 0 || start >= total {
return None;
}
let end = (start + len).min(total);
let file_of = |row: usize| offsets.partition_point(|&o| o <= row).saturating_sub(1);
Some((file_of(start), file_of(end - 1).min(files - 1)))
}
pub(super) fn window_of(
lf: &LazyFrame,
files: Option<&RemoteFiles>,
records: Option<&dyn crate::formats::pushdown::Windowed>,
read_as_text: &[PlSmallStr],
start: usize,
len: usize,
all_columns: Vec<Expr>,
) -> PolarsResult<LazyFrame> {
if let Some(records) = records {
return Ok(records.window(start, len)?.select(all_columns));
}
if let Some((files, offsets)) = files.and_then(|f| f.offsets.as_ref().map(|o| (f, o)))
&& let Some((first, last)) = files_holding(offsets, start, len)
{
let urls: Vec<String> = files_with_rows(offsets, first, last)
.into_iter()
.map(|i| files.urls[i].clone())
.collect();
let lf = (files.scan)(&urls, read_as_text)?;
return Ok(lf
.select(all_columns)
.slice((start - offsets[first]) as i64, len as u32));
}
Ok(lf
.clone()
.select(all_columns)
.slice(start as i64, len as u32))
}
#[derive(Clone)]
pub(crate) struct ViewRows {
lf: LazyFrame,
files: Option<RemoteFiles>,
records: Option<Arc<dyn crate::formats::pushdown::Windowed>>,
read_as_text: Vec<PlSmallStr>,
pub(crate) buffer: Option<(DataFrame, usize)>,
pub(crate) num_rows: Option<usize>,
pub(crate) streaming: bool,
pub(crate) whole: bool,
pub(crate) reads_up_to: bool,
}
pub(crate) fn sees_every_row_first(lf: &LazyFrame) -> bool {
use polars::lazy::dsl::DslPlan;
lf.logical_plan.into_iter().any(|node| {
matches!(
node,
DslPlan::Sort { .. } | DslPlan::GroupBy { .. } | DslPlan::Pivot { .. }
)
})
}
pub(crate) fn reads_up_to_a_window(lf: &LazyFrame) -> bool {
use polars::lazy::dsl::{DslPlan, FileScanDsl};
lf.logical_plan.into_iter().any(|node| match node {
DslPlan::Filter { .. } => true,
DslPlan::Scan { scan_type, .. } => !matches!(
**scan_type,
FileScanDsl::Parquet { .. } | FileScanDsl::Ipc { .. }
),
_ => false,
})
}
impl ViewRows {
pub(crate) fn window(
&self,
start: usize,
len: usize,
exprs: Vec<Expr>,
) -> PolarsResult<LazyFrame> {
window_of(
&self.lf,
self.files.as_ref(),
self.records.as_deref(),
&self.read_as_text,
start,
len,
exprs,
)
}
#[cfg(test)]
pub(crate) fn of(lf: LazyFrame, buffer: Option<(DataFrame, usize)>) -> Self {
Self {
whole: sees_every_row_first(&lf),
reads_up_to: reads_up_to_a_window(&lf),
lf,
files: None,
records: None,
read_as_text: Vec::new(),
buffer,
num_rows: None,
streaming: false,
}
}
}
pub(super) fn align_to_row_groups(
offsets: &[usize],
view_start: usize,
view_end: usize,
start: usize,
end: usize,
cap: usize,
) -> (usize, usize) {
let Some(groups) = offsets.len().checked_sub(1).filter(|n| *n > 0) else {
return (start, end);
};
let group_of = |row: usize| {
offsets
.partition_point(|&o| o <= row)
.saturating_sub(1)
.min(groups - 1)
};
let last_row = |s: usize, e: usize| e.saturating_sub(1).max(s);
let (mut lo, mut hi) = (
group_of(view_start),
group_of(last_row(view_start, view_end)),
);
let (want_lo, want_hi) = (group_of(start), group_of(last_row(start, end)));
let fits = |lo: usize, hi: usize| cap == 0 || offsets[hi + 1] - offsets[lo] <= cap;
loop {
if hi < want_hi && fits(lo, hi + 1) {
hi += 1;
} else if lo > want_lo && fits(lo - 1, hi) {
lo -= 1;
} else {
break;
}
}
(offsets[lo], offsets[hi + 1])
}
impl DataTableState {
pub fn scroll_would_trigger_collect(&self, rows: i64) -> bool {
if rows < 0 && self.view.start_row == 0 {
return false;
}
let new_start_row = if self.view.start_row as i64 + rows <= 0 {
0
} else {
if let Some(df) = self.view.df.as_ref()
&& rows > 0
&& df.shape().0 <= self.visible_rows
{
return false;
}
let unclamped = (self.view.start_row as i64 + rows) as usize;
if rows > 0 {
unclamped.min(self.view.num_rows.saturating_sub(self.visible_rows))
} else {
unclamped
}
};
if new_start_row == self.view.start_row {
return false;
}
let view_end = new_start_row
+ self
.visible_rows
.min(self.view.num_rows.saturating_sub(new_start_row));
let within_buffer = new_start_row >= self.view.buffered_start_row
&& view_end <= self.view.buffered_end_row
&& self.view.buffered_end_row > 0;
!within_buffer
}
pub fn slide_table(&mut self, rows: i64) -> bool {
if rows < 0 && self.view.start_row == 0 {
return false;
}
let new_start_row = if self.view.start_row as i64 + rows <= 0 {
0
} else {
if let Some(df) = self.view.df.as_ref()
&& rows > 0
&& df.shape().0 <= self.visible_rows
{
return false;
}
let unclamped = (self.view.start_row as i64 + rows) as usize;
if rows > 0 {
unclamped.min(self.view.num_rows.saturating_sub(self.visible_rows))
} else {
unclamped
}
};
if new_start_row == self.view.start_row {
return false;
}
let view_end = new_start_row
+ self
.visible_rows
.min(self.view.num_rows.saturating_sub(new_start_row));
let within_buffer = new_start_row >= self.view.buffered_start_row
&& view_end <= self.view.buffered_end_row
&& self.view.buffered_end_row > 0;
self.view.start_row = new_start_row;
if within_buffer {
if self.table_state.selected().is_none() {
self.table_state.select(Some(0));
}
false
} else {
true }
}
#[cfg(test)]
pub fn collect(&mut self) {
if self.defer_collect {
return;
}
if !self.view.num_rows_valid {
match collect_lazy(row_count_lf(&self.view.lf), self.polars_streaming) {
Ok(df) => {
self.error = None;
let n = match df.get(0).as_deref().and_then(|row| row.first()) {
Some(AnyValue::UInt64(len)) => *len as usize,
_ => 0,
};
self.set_num_rows(n);
}
Err(e) => {
self.error = Some(e);
self.set_num_rows(0);
}
}
}
let Some(request) = self.prepare_async_collect(None) else {
return;
};
match collect_lazy(request.lf, request.polars_streaming) {
Ok(df) => self.apply_async_collect(request.plan.fit(df)),
Err(e) => self.error = Some(e),
}
}
#[cfg(not(test))]
pub(super) fn collect(&mut self) {
if !self.defer_collect {
self.needs_recollect = true;
}
}
pub(crate) fn binary_stub_exprs(&self) -> Vec<Expr> {
self.view
.column_order
.iter()
.map(|name| {
if matches!(self.view.schema.get(name.as_str()), Some(DataType::Binary)) {
lit(binary_stub()).alias(name.as_str())
} else {
col(name.as_str())
}
})
.collect()
}
pub fn prepare_async_collect(
&mut self,
num_rows_override: Option<usize>,
) -> Option<CollectRequest> {
if self.visible_rows > 0 {
self.proximity_threshold = self.proximity();
}
if let Some(n) = num_rows_override {
self.view.num_rows = n;
self.view.num_rows_valid = true;
}
let count_known = self.view.num_rows_valid;
let bound = self.num_rows_bound();
if count_known {
if self.view.num_rows > 0 {
let max_start = self.view.num_rows.saturating_sub(1);
if self.view.start_row > max_start {
self.view.start_row = max_start;
}
} else {
self.view.start_row = 0;
self.drop_buffer();
self.view.df = None;
self.view.locked_df = None;
return None;
}
}
if self.view.column_order.is_empty() {
self.drop_buffer();
self.view.df = None;
self.view.locked_df = None;
return None;
}
let view_start = self.view.start_row;
let view_end = self.view.start_row + self.visible_rows.min(bound - self.view.start_row);
let within_buffer = view_start >= self.view.buffered_start_row
&& view_end <= self.view.buffered_end_row
&& self.view.buffered_end_row > 0;
let (new_buffer_start, new_buffer_end) = if within_buffer {
let dist_to_start = view_start.saturating_sub(self.view.buffered_start_row);
let dist_to_end = self.view.buffered_end_row.saturating_sub(view_end);
let needs_expansion_back =
dist_to_start <= self.proximity_threshold && self.view.buffered_start_row > 0;
let needs_expansion_forward =
dist_to_end <= self.proximity_threshold && self.view.buffered_end_row < bound;
if !needs_expansion_back && !needs_expansion_forward {
(self.view.buffered_start_row, self.view.buffered_end_row)
} else {
let mut s = if needs_expansion_back {
view_start.saturating_sub(self.reach_rows(self.pages_lookback))
} else {
self.view.buffered_start_row
};
let mut e = if needs_expansion_forward {
(view_end + self.reach_rows(self.pages_lookahead)).min(bound)
} else {
self.view.buffered_end_row
};
self.fit_window(view_start, view_end, &mut s, &mut e);
(s, e)
}
} else {
let had_buffer = self.view.buffered_end_row > 0;
let scrolled_past_end = had_buffer && view_start >= self.view.buffered_end_row;
let scrolled_past_start = had_buffer && view_end <= self.view.buffered_start_row;
let extend_forward_ok = scrolled_past_end
&& (view_start - self.view.buffered_end_row)
<= self.reach_rows(self.pages_lookahead);
let extend_backward_ok = scrolled_past_start
&& (self.view.buffered_start_row - view_end)
<= self.reach_rows(self.pages_lookback);
let mut s;
let mut e;
if extend_forward_ok {
s = self.view.buffered_start_row;
e = (view_end + self.reach_rows(self.pages_lookahead)).min(bound);
} else if extend_backward_ok {
s = view_start.saturating_sub(self.reach_rows(self.pages_lookback));
e = self.view.buffered_end_row;
} else {
s = view_start.saturating_sub(self.reach_rows(self.pages_lookback));
e = (view_end + self.reach_rows(self.pages_lookahead)).min(bound);
let min_initial_len = self.min_buffer_len();
let current_len = e.saturating_sub(s);
if current_len < min_initial_len {
let need = min_initial_len.saturating_sub(current_len);
let can_extend_end = bound.saturating_sub(e);
let can_extend_start = s;
if can_extend_end >= need {
e = (e + need).min(bound);
} else if can_extend_start >= need {
s = s.saturating_sub(need);
} else {
e = (e + can_extend_end).min(bound);
s = s.saturating_sub(need.saturating_sub(can_extend_end));
}
}
}
self.fit_window(view_start, view_end, &mut s, &mut e);
(s, e)
};
let buffer_size = new_buffer_end.saturating_sub(new_buffer_start);
if buffer_size == 0 {
return None;
}
if self.holds_buffer(new_buffer_start, new_buffer_end) {
self.slice_buffer_into_display();
if self.table_state.selected().is_none() {
self.table_state.select(Some(0));
}
return None;
}
let lf = match self.buffer_lf(new_buffer_start, buffer_size) {
Ok(lf) => lf,
Err(e) => {
self.error = Some(e);
return None;
}
};
let num_rows = if count_known {
self.view.num_rows
} else {
new_buffer_end
};
Some(CollectRequest {
lf,
polars_streaming: self.polars_streaming,
buffer_start: new_buffer_start,
buffer_end: new_buffer_end,
plan: self.fill_plan(new_buffer_start, new_buffer_end, num_rows, count_known),
})
}
pub(super) fn fill_plan(
&self,
buffer_start: usize,
buffer_end: usize,
num_rows: usize,
count_known: bool,
) -> FillPlan {
let held = self
.abuts_buffer(buffer_start, buffer_end.saturating_sub(buffer_start))
.then(|| self.view.buffered_df.clone())
.flatten()
.map(|df| (df, self.view.buffered_start_row));
FillPlan {
buffer_start,
buffer_end,
num_rows,
count_known,
indexing: self.indexing().is_some(),
held,
view_start: self.view.start_row,
view_len: self.visible_rows,
max_rows: self.max_buffered_rows,
max_mb: self.max_buffered_mb,
}
}
pub fn apply_async_collect(&mut self, result: CollectResult) {
let CollectResult {
df,
start,
returned: returned_rows,
bytes_per_row,
buffer_start,
buffer_end,
num_rows,
count_known,
indexing,
} = result;
let requested_rows = buffer_end.saturating_sub(buffer_start);
if count_known {
self.view.num_rows = num_rows;
self.view.num_rows_valid = true;
} else if returned_rows < requested_rows
&& (buffer_start == 0 || returned_rows > 0)
&& !indexing
&& self.indexing().is_none()
{
self.view.num_rows = buffer_start + returned_rows;
self.view.num_rows_valid = true;
} else if !self.view.num_rows_valid {
self.view.num_rows = self.view.num_rows.max(buffer_end);
}
self.error = None;
self.remember_pristine_count();
if bytes_per_row.is_some() {
self.view.observed_bytes_per_row = bytes_per_row;
}
let end = start + df.height();
let view_end = self.view.start_row + self.visible_rows.max(1);
let reaches_end = end >= buffer_start + returned_rows;
let shows_view = start <= self.view.start_row
&& (self.view.start_row < end || (returned_rows < requested_rows && reaches_end));
if !shows_view {
self.needs_recollect = true;
return;
}
self.release_display_buffer();
self.view.buffered_start_row = start;
self.view.buffered_end_row = end;
self.view.buffered_df = Some(df);
self.slice_buffer_into_display();
if self.table_state.selected().is_none() {
self.table_state.select(Some(0));
}
if view_end > end && end < self.view.num_rows {
self.needs_recollect = true;
}
}
fn abuts_buffer(&self, start: usize, rows: usize) -> bool {
self.stitches_buffer()
&& (start == self.view.buffered_end_row || start + rows == self.view.buffered_start_row)
}
pub(crate) fn sampled_from(
source: DataTableState,
sample: crate::analysis::sampling::Sample,
schema: &Schema,
rows: Arc<crate::analysis::table_sample::SampleRows>,
through: bool,
path: Option<crate::analysis::table_sample::DrawPath>,
) -> Result<Self> {
let mut view = source.sample_view(DataFrame::empty_with_schema(schema))?;
let frame = scanned_frame(&view.original_lf)
.ok_or_else(|| color_eyre::eyre::eyre!("a sample's frame has no rows to scan"))?;
view.sampled = Some(Box::new(Sampled {
source: Box::new(source),
sample,
rows,
frame,
through,
drawn: None,
path,
}));
Ok(view)
}
pub fn sampled(&self) -> Option<&Sampled> {
self.sampled.as_deref()
}
pub fn unsampled(&self) -> &DataTableState {
self.sampled
.as_ref()
.map_or(self, |sampled| sampled.source.as_ref())
}
pub(crate) fn into_unsampled(mut self) -> DataTableState {
match self.sampled.take() {
Some(sampled) => *sampled.source,
None => self,
}
}
pub(crate) fn sample_grew(&mut self) -> Option<bool> {
let sampled = self.sampled.as_ref()?;
let chunks = sampled.rows.take_new();
if chunks.is_empty() {
return None;
}
let mut frame = (*sampled.frame).clone();
for chunk in &chunks {
frame.vstack_mut(chunk).ok()?;
}
Some(self.rebind_sample(Arc::new(frame), false))
}
pub(crate) fn sample_drawn(&mut self, drawn: crate::analysis::table_sample::Drawn) {
let Some(sampled) = self.sampled.as_mut() else {
return;
};
let ordered = sampled.rows.take_in_source_order().ok().flatten();
sampled.path = drawn.path;
sampled.drawn = Some(drawn);
if let Some(frame) = ordered {
self.rebind_sample(Arc::new(frame), true);
}
}
fn rebind_sample(&mut self, frame: Arc<DataFrame>, reordered: bool) -> bool {
let Some(old) = self.sampled.as_ref().map(|sampled| sampled.frame.clone()) else {
return false;
};
let rows_stand = !reordered
&& self.view.sort_columns.is_empty()
&& self.view.sort_ascending
&& self.scan_is_the_root();
let rows = frame.height();
self.each_frame(|lf| {
crate::analysis::table_sample::rebind(&mut lf.logical_plan, &old, &frame)
});
if let Some(sampled) = self.sampled.as_mut() {
sampled.frame = frame;
}
self.invalidate_num_rows();
if self.is_pristine() {
self.set_num_rows(rows);
} else if self.scan_is_the_root() {
self.pristine_rows = Some(rows);
}
if !rows_stand {
self.drop_buffer();
}
self.needs_recollect = true;
rows_stand
}
pub(crate) fn sample_row_bytes(&self, from_source: bool) -> usize {
let schema = if from_source {
&self.original_schema
} else {
&self.view.schema
};
let columns: Vec<String> = schema
.iter_names()
.filter(|name| name.as_str() != crate::formats::schema_union::DRIFT_COLUMN)
.map(|name| name.to_string())
.collect();
if !from_source && columns.len() == self.view.column_order.len() {
return self.bytes_per_row();
}
estimate_bytes_per_row(schema, &columns, &self.column_bytes)
}
pub fn files_a_page_reads(&self, start: usize, len: usize) -> Option<usize> {
let offsets = self.files_window().and_then(|f| f.offsets.as_ref())?;
let (first, last) = files_holding(offsets, start, len)?;
Some(files_with_rows(offsets, first, last).len())
}
pub(super) fn buffer_lf(&self, start: usize, len: usize) -> PolarsResult<LazyFrame> {
let mut all_columns = self.binary_stub_exprs();
if self.carries_source_rows() {
all_columns.push(col(crate::formats::schema_union::DRIFT_COLUMN));
}
self.window_lf(start, len, all_columns)
}
pub(crate) fn carries_source_rows(&self) -> bool {
self.view.drift_column_present
|| (self.scan_is_the_root() && (self.source_rows_at_open || self.view.view_numbered))
}
pub fn row_numbers_from(&self, start: usize, rows: usize) -> Vec<usize> {
let view = |i: usize| start + i + self.row_start_index;
let places = self
.view
.buffered_df
.as_ref()
.filter(|_| self.carries_source_rows())
.and_then(|df| df.column(crate::formats::schema_union::DRIFT_COLUMN).ok())
.and_then(|column| {
let offset = start.checked_sub(self.view.buffered_start_row)?;
let len = rows.min(column.len().saturating_sub(offset));
let slice = column.slice(offset as i64, len);
let places = slice.u32().ok()?;
let place = |p: usize| {
self.numbering
.as_ref()
.and_then(|lines| lines.line_in_file(p))
.unwrap_or(p)
};
Some(
places
.iter()
.map(|p| p.map(|p| place(p as usize) + self.row_start_index))
.collect::<Vec<_>>(),
)
});
(0..rows)
.map(|i| {
places
.as_ref()
.and_then(|p| p.get(i).copied().flatten())
.unwrap_or_else(|| view(i))
})
.collect()
}
pub(super) fn window_lf(
&self,
start: usize,
len: usize,
all_columns: Vec<Expr>,
) -> PolarsResult<LazyFrame> {
window_of(
&self.view.lf,
self.files_window(),
self.window_now().as_deref(),
&self.read_as_text,
start,
len,
all_columns,
)
}
pub(crate) fn view_rows(&self) -> ViewRows {
ViewRows {
lf: self.view.lf.clone(),
files: self.files_window().cloned(),
records: self.window_now().filter(|_| self.indexing().is_none()),
read_as_text: self.read_as_text.clone(),
buffer: self
.view
.buffered_df
.as_ref()
.filter(|_| self.buffer_on_hand())
.map(|df| (df.clone(), self.view.buffered_start_row)),
num_rows: self.view.num_rows_valid.then_some(self.view.num_rows),
streaming: self.polars_streaming,
whole: sees_every_row_first(&self.view.lf),
reads_up_to: reads_up_to_a_window(&self.view.lf),
}
}
pub(crate) fn go_to_found_row(&mut self, row: usize) -> bool {
if !self.view.num_rows_valid && self.view.num_rows <= row {
self.view.num_rows = row + 1;
}
self.scroll_to_row_centered(row)
}
pub(crate) fn cursor_row(&self) -> usize {
self.view.start_row + self.table_state.selected().unwrap_or(0)
}
pub(super) fn bytes_per_row(&self) -> usize {
self.view.observed_bytes_per_row.unwrap_or_else(|| {
estimate_bytes_per_row(
&self.view.schema,
&self.view.column_order,
&self.column_bytes,
)
})
}
pub fn estimated_row_bytes(&self) -> usize {
self.bytes_per_row()
}
pub fn source_file_count(&self) -> Option<usize> {
self.is_pristine().then(|| self.loaded_file_count())
}
pub(crate) fn loaded_file_count(&self) -> usize {
if !self.drift_files.is_empty() {
self.drift_files.len()
} else if let Some(remote) = &self.remote_files {
remote.urls.len()
} else {
1
}
}
pub(super) fn byte_cap_rows(&self) -> usize {
if self.max_buffered_mb == 0 {
return 0;
}
let max_bytes = self.max_buffered_mb * 1024 * 1024;
(max_bytes / self.bytes_per_row()).max(self.visible_rows.max(1))
}
pub(super) fn remote_window(&self) -> bool {
self.remote_source && self.is_pristine()
}
fn files_window(&self) -> Option<&RemoteFiles> {
self.remote_files.as_ref().filter(|_| self.is_pristine())
}
fn reach_rows(&self, pages: usize) -> usize {
if !self.remote_window() || self.remote_files.is_some() {
return pages * self.visible_rows.max(1);
}
let window = if self.max_buffered_rows > 0 {
self.max_buffered_rows
} else {
DEFAULT_MAX_BUFFERED_ROWS
};
window / 2
}
fn min_buffer_len(&self) -> usize {
self.visible_rows.max(1)
+ self.reach_rows(self.pages_lookahead)
+ self.reach_rows(self.pages_lookback)
}
pub fn at_end(&self) -> bool {
self.view.start_row == self.view.num_rows.saturating_sub(self.visible_rows)
}
fn fit_window(
&self,
view_start: usize,
view_end: usize,
buffer_start: &mut usize,
buffer_end: &mut usize,
) {
let byte_cap = self.byte_cap_rows();
let cap = match (self.max_buffered_rows, byte_cap) {
(0, cap) | (cap, 0) => cap,
(rows, bytes) => rows.min(bytes),
};
if cap > 0 {
shrink_around_view(
view_start,
view_end,
cap,
0,
self.num_rows_bound(),
buffer_start,
buffer_end,
);
}
let Some(offsets) = self
.row_group_offsets
.as_deref()
.filter(|_| self.remote_window())
else {
return;
};
(*buffer_start, *buffer_end) = align_to_row_groups(
offsets,
view_start,
view_end,
*buffer_start,
*buffer_end,
cap,
);
if cap > 0 {
let (floor, ceil) = (*buffer_start, *buffer_end);
shrink_around_view(
view_start,
view_end,
cap,
floor,
ceil,
buffer_start,
buffer_end,
);
}
if let Some(file_offsets) = self.remote_files.as_ref().and_then(|f| f.offsets.as_ref()) {
(*buffer_start, *buffer_end) = limit_files(
file_offsets,
view_start,
view_end,
*buffer_start,
*buffer_end,
MAX_FILES_PER_BUFFER,
);
}
if self.buffer_on_hand() {
let (held_start, held_end) = (self.view.buffered_start_row, self.view.buffered_end_row);
if held_start <= *buffer_start && *buffer_start < held_end && held_end < *buffer_end {
*buffer_start = held_end;
} else if *buffer_start < held_start
&& held_start < *buffer_end
&& *buffer_end <= held_end
{
*buffer_end = held_start;
}
}
}
fn release_display_buffer(&mut self) {
self.widths.rows_arrived();
self.view.buffered_df = None;
self.view.locked_df = None;
self.view.df = None;
}
pub(super) fn slice_buffer_into_display(&mut self) {
let full_df = match self.view.buffered_df.as_ref() {
Some(df) => df,
None => return,
};
if self.view.locked_columns_count > 0 {
let locked_names: Vec<&str> = self
.view
.column_order
.iter()
.take(self.view.locked_columns_count)
.map(|s| s.as_str())
.collect();
if let Ok(locked_df) = full_df.select(locked_names) {
self.view.locked_df = Some(locked_df);
}
} else {
self.view.locked_df = None;
}
let scroll_names: Vec<&str> = self
.view
.column_order
.iter()
.skip(self.frozen_shown() + self.termcol_index)
.map(|s| s.as_str())
.collect();
if scroll_names.is_empty() {
self.view.df = None;
} else {
if let Ok(scroll_df) = full_df.select(scroll_names) {
self.view.df = Some(scroll_df);
}
}
}
pub fn wants_to_load_ahead(&self) -> bool {
if self.visible_rows == 0
|| self.view.buffered_df.is_none()
|| !self.page_on_hand(self.view.start_row)
{
return false;
}
let near = self.proximity();
let view_end = self.view.start_row
+ self
.visible_rows
.min(self.num_rows_bound().saturating_sub(self.view.start_row));
let behind = self.view.start_row - self.view.buffered_start_row <= near
&& self.view.buffered_start_row > 0;
let ahead = self.view.buffered_end_row - view_end <= near
&& self.view.buffered_end_row < self.num_rows_bound();
behind || ahead
}
fn proximity(&self) -> usize {
(self.reach_rows(self.pages_lookahead) / 2).max(self.visible_rows)
}
pub fn buffer_position(&self) -> (u64, usize, usize, usize) {
(
self.len_generation(),
self.view.start_row,
self.view.buffered_start_row,
self.view.buffered_end_row,
)
}
pub(crate) fn page_on_hand(&self, start: usize) -> bool {
let bound = self.num_rows_bound();
let end = start + self.visible_rows.min(bound.saturating_sub(start));
self.view.buffered_df.is_some()
&& self.view.buffered_end_row > 0
&& start >= self.view.buffered_start_row
&& end <= self.view.buffered_end_row
}
pub(crate) fn start_to_draw(&mut self) -> usize {
if self.page_on_hand(self.view.start_row) {
self.view.drawn_start = self.view.start_row;
self.view.start_row
} else if self.page_on_hand(self.view.drawn_start) {
self.view.drawn_start
} else {
self.view.start_row
}
}
}