jerry-024 commented on code in PR #858: URL: https://github.com/apache/paimon-rust/pull/858#discussion_r4056622792
########## crates/paimon/src/table/partition_row_count.rs: ########## @@ -0,0 +1,1075 @@ +// 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. + +//! Real per-partition row counts, computed from manifests alone in bounded memory. +//! +//! [`Table::partition_stats`] and the `$files` system table both materialize every +//! manifest entry together with its per-column statistics, which does not fit in +//! memory for very large tables, and summing `row_count` over files over-counts +//! data-evolution tables, where several column-group files cover the same rows. +//! +//! This module instead: +//! - decodes manifests through [`SlimManifestEntry`], which borrows the handful of +//! fields it needs and skips statistics in place, so nothing per-file is retained; +//! - for data-evolution tables, counts files carrying a `first_row_id` by the +//! *union* of their row-id ranges, collapsing overlapping column-group/blob files; +//! - nets ADD/DELETE entries with a delete set built only from the manifests that +//! contain deletes, the same semantics as the scan's manifest merge. +//! +//! Peak memory is bounded by partitions, live DELETE entries, up to one million +//! retained ADD identities, disjoint row-id ranges, and in-flight manifest buffers +//! rather than all live file metadata and column statistics. Highly fragmented +//! row-id space or large embedded indexes may still increase retained state. + +use std::collections::{BTreeMap, HashMap}; +use std::sync::atomic::{AtomicI64, Ordering}; +use std::sync::Arc; + +use futures::{StreamExt, TryStreamExt}; + +use crate::io::FileIO; +use crate::spec::avro::{ + visit_slim_index_manifest_entries, visit_slim_manifest_entries, SharedSchemaCache, + SlimManifestEntry, +}; +use crate::spec::{ + BinaryRow, CoreOptions, DataField, FileKind, ManifestFileMeta, ManifestList, Predicate, +}; +use crate::table::partition_filter::PartitionFilter; +use crate::table::read_builder::split_scan_predicates; +use crate::table::stats_filter::FileStatsRows; +use crate::table::Table; + +/// Independent I/O and blocking decode limits for data manifests. +const MANIFEST_READ_CONCURRENCY: usize = 32; + +/// ADD identities retained from delete-bearing manifests to avoid fetching them twice. +/// This bounds entry count, not embedded-index payload bytes. +const RETAINED_ADD_BUDGET: i64 = 1_000_000; + +fn manifest_decode_concurrency() -> usize { + std::thread::available_parallelism() + .map_or(2, |parallelism| parallelism.get()) + .clamp(2, MANIFEST_READ_CONCURRENCY) +} + +const DELETION_VECTORS_INDEX_TYPE: &str = "DELETION_VECTORS"; + +/// Real row count of one partition. +#[derive(Debug, Clone, PartialEq, Eq)] +pub struct PartitionRowCount { + /// The partition's typed values, one field per partition key. + pub partition_row: BinaryRow, + /// Rows in the partition: overlapping data-evolution files are counted once + /// and deletion-vector rows are subtracted. For primary-key tables this is + /// the physical count, before merging keys. `None` when it cannot be known + /// exactly (a file without a row count, or a deletion vector without a + /// cardinality) — never a guess. + pub record_count: Option<i64>, +} + +/// Disjoint set of inclusive row-id ranges, coalescing overlapping and adjacent ones. +#[derive(Debug, Default)] +struct RowRangeSet { + ranges: BTreeMap<i64, i64>, +} + +impl RowRangeSet { + fn insert(&mut self, mut start: i64, mut end: i64) { + if start > end { + return; + } + if let Some((&prev_start, &prev_end)) = self.ranges.range(..=start).next_back() { + if prev_end >= end { + return; + } + if prev_end.saturating_add(1) >= start { + start = prev_start; + self.ranges.remove(&prev_start); + } + } + while let Some((&next_start, &next_end)) = + self.ranges.range(start..=end.saturating_add(1)).next() + { + end = end.max(next_end); + self.ranges.remove(&next_start); + } + self.ranges.insert(start, end); + } + + fn merge(&mut self, other: RowRangeSet) { + if self.ranges.is_empty() { + self.ranges = other.ranges; + return; + } + for (start, end) in other.ranges { + self.insert(start, end); + } + } + + fn total(&self) -> i128 { + self.ranges + .iter() + .map(|(start, end)| i128::from(*end) - i128::from(*start) + 1) + .sum() + } +} + +#[derive(Debug, Default)] +struct PartitionAccum { + /// Rows of files without a `first_row_id`. + plain_rows: i128, + /// Row-id ranges of files with a `first_row_id`. + row_ranges: RowRangeSet, + row_count_unknown: bool, +} + +impl PartitionAccum { + fn add_file( + &mut self, + row_count: i64, + first_row_id: Option<i64>, + data_evolution_enabled: bool, + ) { + if row_count < 0 { + self.row_count_unknown = true; + return; + } + if data_evolution_enabled { + if let Some(first) = first_row_id { + if row_count > 0 { + let Some(last) = first.checked_add(row_count - 1) else { + self.row_count_unknown = true; + return; + }; + self.row_ranges.insert(first, last); + } + return; + } + } + self.plain_rows += i128::from(row_count); + } + + fn merge(&mut self, other: PartitionAccum) { + self.plain_rows += other.plain_rows; + self.row_ranges.merge(other.row_ranges); + self.row_count_unknown |= other.row_count_unknown; + } + + fn record_count(&self, deleted: Option<Option<i128>>) -> Option<i64> { + if self.row_count_unknown { + return None; + } + let deleted = deleted.unwrap_or(Some(0))?; + let rows = self.plain_rows + self.row_ranges.total() - deleted; + (rows >= 0).then(|| i64::try_from(rows).ok()).flatten() + } +} + +type PartitionAccums = HashMap<Vec<u8>, PartitionAccum>; + +fn accumulate( + accums: &mut PartitionAccums, + partition: &[u8], + row_count: i64, + first_row_id: Option<i64>, + data_evolution_enabled: bool, +) { + let accum = match accums.get_mut(partition) { + Some(accum) => accum, + None => accums.entry(partition.to_vec()).or_default(), + }; + accum.add_file(row_count, first_row_id, data_evolution_enabled); +} + +/// Identifiers of deleted files, matching the full Paimon `Identifier` semantics. +/// +/// Nested so ADD lookups borrow partition/file-name bytes from the decode buffer +/// and each partition is stored only once. +#[derive(Debug, Default)] +struct DeleteSet { + by_partition: HashMap<Box<[u8]>, DeletedFiles>, +} + +type DeletedFiles = HashMap<Box<str>, Vec<DeletedFile>>; + +#[derive(Debug, PartialEq, Eq)] +struct DeletedFile { + bucket: i32, + level: i32, + extra_files: Vec<Box<str>>, + // Exact identity matching requires the payload. A spill-backed delete set can + // replace this copy if embedded-index memory becomes a measured bottleneck. + embedded_index: Option<Box<[u8]>>, + external_path: Option<Box<str>>, +} + +impl DeletedFile { + fn from_entry(entry: &SlimManifestEntry<'_>) -> Self { + Self { + bucket: entry.bucket, + level: entry.level, + extra_files: entry + .extra_files + .iter() + .map(|value| Box::from(*value)) + .collect(), + embedded_index: entry.embedded_index.map(Box::from), + external_path: entry.external_path.map(Box::from), Review Comment: Thanks. The limit is intentionally entry-based, not a strict byte cap. The PR description documents this trade-off, including the measured RSS increase for retained ADDs and the fact that exact DELETE/embedded-index payloads remain variable-sized. Exact matching is required; byte accounting plus spill would be a separate, substantially larger mechanism. We will keep the current design unless production measurements show this state is the actual memory bottleneck. -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
