viirya commented on code in PR #25583: URL: https://github.com/apache/datafusion/pull/25583#discussion_r4116954570
########## datafusion/physical-expr-common/src/metrics/snapshot.rs: ########## @@ -0,0 +1,301 @@ +// 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. + +//! Fixed-membership snapshots of an append-only registry. +//! +//! Registration only appends to a vector. Partition readers share an index that +//! catches up with registration on demand; snapshots remember their original end +//! position even when a later reader has advanced the index past that position. + +use super::Metric; +use parking_lot::Mutex; +use std::collections::HashMap; +use std::fmt; +use std::sync::{Arc, OnceLock}; + +#[derive(Debug, Default)] +pub(super) struct Registry { + pub(super) metrics: Vec<Arc<Metric>>, + // No index allocation or maintenance on the registration path. + index: Option<Box<PartitionIndex>>, +} + +#[derive(Debug, Default)] +struct PartitionIndex { + // Number of registry entries already examined, including global metrics. + indexed: usize, + // Partition ID -> positions in Registry::metrics, in registration order. + positions: HashMap<usize, Vec<usize>>, +} + +impl Registry { + pub(super) fn new(metrics: Vec<Arc<Metric>>) -> Self { + Self { + metrics, + index: None, + } + } + + fn select(&mut self, partition: usize, end: usize) -> Vec<Arc<Metric>> { + let index = self.index.get_or_insert_with(Default::default); + for position in index.indexed..end { + if let Some(partition) = self.metrics[position].partition() { + index.positions.entry(partition).or_default().push(position); + } + } + index.indexed = index.indexed.max(end); + let Some(positions) = index.positions.get(&partition) else { + return Vec::new(); + }; + // Another snapshot may already have indexed registrations after our end. + let len = positions.partition_point(|&position| position < end); + positions[..len] + .iter() + .map(|&i| Arc::clone(&self.metrics[i])) + .collect() + } +} + +#[derive(Clone)] +pub(super) enum Snapshot { + Owned(Vec<Arc<Metric>>), + Deferred(Arc<Deferred>), +} + +pub(super) struct Deferred { + registry: Arc<Mutex<Registry>>, + end: usize, + flat: OnceLock<Vec<Arc<Metric>>>, +} + +impl Default for Snapshot { + fn default() -> Self { + Self::Owned(Vec::new()) + } +} + +impl fmt::Debug for Snapshot { + fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result { + f.debug_list().entries(self.iter()).finish() + } +} + +impl Deferred { + fn materialize(&self) -> Vec<Arc<Metric>> { + self.registry.lock().metrics[..self.end].to_vec() + } +} + +impl Snapshot { + pub(super) fn new(registry: Arc<Mutex<Registry>>, end: usize) -> Self { + Self::Deferred(Arc::new(Deferred { + registry, + end, + flat: OnceLock::new(), + })) + } + + pub(super) fn for_partition(&self, partition: usize) -> Self { + Self::Owned(match self { + Self::Owned(metrics) => metrics + .iter() + .filter(|m| m.partition() == Some(partition)) Review Comment: @alamb This fallback is now explicitly limited to `Standalone` collections, which have no registry index. `All` uses the index, while `Partition` checks the requested ID directly without filtering. ########## datafusion/physical-expr-common/src/metrics/snapshot.rs: ########## @@ -0,0 +1,301 @@ +// 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. + +//! Fixed-membership snapshots of an append-only registry. +//! +//! Registration only appends to a vector. Partition readers share an index that +//! catches up with registration on demand; snapshots remember their original end +//! position even when a later reader has advanced the index past that position. + +use super::Metric; +use parking_lot::Mutex; +use std::collections::HashMap; +use std::fmt; +use std::sync::{Arc, OnceLock}; + +#[derive(Debug, Default)] +pub(super) struct Registry { Review Comment: @alamb The index avoids repeatedly scanning all earlier registrations when reporting successive partitions. Readers incrementally index new registrations, so each entry is indexed once. I've added this rationale beside the index field; registration itself remains a vector append. ########## datafusion/physical-expr-common/src/metrics/snapshot.rs: ########## @@ -0,0 +1,301 @@ +// 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. + +//! Fixed-membership snapshots of an append-only registry. +//! +//! Registration only appends to a vector. Partition readers share an index that +//! catches up with registration on demand; snapshots remember their original end +//! position even when a later reader has advanced the index past that position. + +use super::Metric; +use parking_lot::Mutex; +use std::collections::HashMap; +use std::fmt; +use std::sync::{Arc, OnceLock}; + +#[derive(Debug, Default)] +pub(super) struct Registry { + pub(super) metrics: Vec<Arc<Metric>>, + // No index allocation or maintenance on the registration path. + index: Option<Box<PartitionIndex>>, +} + +#[derive(Debug, Default)] +struct PartitionIndex { + // Number of registry entries already examined, including global metrics. + indexed: usize, + // Partition ID -> positions in Registry::metrics, in registration order. + positions: HashMap<usize, Vec<usize>>, Review Comment: @alamb I considered storing a Vec of metric handles per partition. With the current snapshot boundary, we'd still need registration positions to exclude entries added after an older snapshot. I kept the positions to avoid storing both and adding another Arc reference per indexed metric, and documented that tradeoff. -- 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] --------------------------------------------------------------------- To unsubscribe, e-mail: [email protected] For additional commands, e-mail: [email protected]
