Skip to content

Commit 72f39c7

Browse files
committed
Reduce V2 duplicate scan requests
Signed-off-by: "Nicholas Gates" <nick@nickgates.com>
1 parent 4109af9 commit 72f39c7

6 files changed

Lines changed: 106 additions & 20 deletions

File tree

vortex-datafusion/src/persistent/opener.rs

Lines changed: 27 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -103,6 +103,8 @@ pub(crate) struct VortexOpener {
103103
pub layout_readers: Arc<DashMap<Path, Weak<dyn LayoutReader>>>,
104104
/// Shared full-file natural split ranges keyed by file path.
105105
pub natural_split_ranges: Arc<DashMap<Path, Arc<[Range<u64>]>>>,
106+
/// Shared V2 file handles keyed by file path.
107+
pub vortex_files: Arc<DashMap<Path, Arc<VortexFile>>>,
106108
/// Whether the query has output ordering specified
107109
pub has_output_ordering: bool,
108110

@@ -139,6 +141,7 @@ impl FileOpener for VortexOpener {
139141
let limit = self.limit;
140142
let layout_readers = Arc::clone(&self.layout_readers);
141143
let natural_split_ranges = Arc::clone(&self.natural_split_ranges);
144+
let vortex_files = Arc::clone(&self.vortex_files);
142145
let has_output_ordering = self.has_output_ordering;
143146
let scan_concurrency = self.scan_concurrency;
144147

@@ -216,10 +219,24 @@ impl FileOpener for VortexOpener {
216219
open_opts = open_opts.with_footer(footer);
217220
}
218221

219-
let vxf = open_opts
220-
.open_read(reader)
221-
.await
222-
.map_err(|e| exec_datafusion_err!("Failed to open Vortex file {e}"))?;
222+
let vxf = if let Some(hit) = vortex_files.get(&file.object_meta.location) {
223+
Arc::clone(hit.value())
224+
} else {
225+
let opened = Arc::new(
226+
open_opts
227+
.open_read(reader)
228+
.await
229+
.map_err(|e| exec_datafusion_err!("Failed to open Vortex file {e}"))?,
230+
);
231+
232+
match vortex_files.entry(file.object_meta.location.clone()) {
233+
Entry::Occupied(entry) => Arc::clone(entry.get()),
234+
Entry::Vacant(entry) => {
235+
entry.insert(Arc::clone(&opened));
236+
opened
237+
}
238+
}
239+
};
223240

224241
// On a miss, cache the parsed footer so other partitions and later executions
225242
// skip the footer fetch and parse. `infer_schema`/`infer_stats` also populate
@@ -915,6 +932,7 @@ mod tests {
915932
metrics_registry: Arc::new(DefaultMetricsRegistry::default()),
916933
layout_readers: Default::default(),
917934
natural_split_ranges: Default::default(),
935+
vortex_files: Default::default(),
918936
has_output_ordering: false,
919937
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
920938
file_metadata_cache: None,
@@ -1111,6 +1129,7 @@ mod tests {
11111129
metrics_registry: Arc::new(DefaultMetricsRegistry::default()),
11121130
layout_readers: Default::default(),
11131131
natural_split_ranges: Default::default(),
1132+
vortex_files: Default::default(),
11141133
has_output_ordering: false,
11151134
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
11161135
file_metadata_cache: None,
@@ -1199,6 +1218,7 @@ mod tests {
11991218
metrics_registry: Arc::new(DefaultMetricsRegistry::default()),
12001219
layout_readers: Default::default(),
12011220
natural_split_ranges: Default::default(),
1221+
vortex_files: Default::default(),
12021222
has_output_ordering: false,
12031223
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
12041224
file_metadata_cache: None,
@@ -1357,6 +1377,7 @@ mod tests {
13571377
metrics_registry: Arc::new(DefaultMetricsRegistry::default()),
13581378
layout_readers: Default::default(),
13591379
natural_split_ranges: Default::default(),
1380+
vortex_files: Default::default(),
13601381
has_output_ordering: false,
13611382
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
13621383
file_metadata_cache: None,
@@ -1418,6 +1439,7 @@ mod tests {
14181439
metrics_registry: Arc::new(DefaultMetricsRegistry::default()),
14191440
layout_readers: Default::default(),
14201441
natural_split_ranges: Default::default(),
1442+
vortex_files: Default::default(),
14211443
has_output_ordering: false,
14221444
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
14231445
file_metadata_cache: None,
@@ -1628,6 +1650,7 @@ mod tests {
16281650
metrics_registry: Arc::new(DefaultMetricsRegistry::default()),
16291651
layout_readers: Default::default(),
16301652
natural_split_ranges: Default::default(),
1653+
vortex_files: Default::default(),
16311654
has_output_ordering: false,
16321655
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
16331656
file_metadata_cache: None,

vortex-datafusion/src/persistent/source.rs

Lines changed: 5 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -32,6 +32,7 @@ use object_store::ObjectStore;
3232
use object_store::path::Path;
3333
use vortex::error::VortexExpect;
3434
use vortex::file::VORTEX_FILE_EXTENSION;
35+
use vortex::file::VortexFile;
3536
use vortex::layout::LayoutReader;
3637
use vortex::layout::scan::v2::scan2_enabled;
3738
use vortex::metrics::DefaultMetricsRegistry;
@@ -200,6 +201,8 @@ pub struct VortexSource {
200201
layout_readers: Arc<DashMap<Path, Weak<dyn LayoutReader>>>,
201202
/// Shared full-file natural split ranges keyed by path.
202203
natural_split_ranges: Arc<DashMap<Path, Arc<[Range<u64>]>>>,
204+
/// Shared V2 file handles keyed by path.
205+
vortex_files: Arc<DashMap<Path, Arc<VortexFile>>>,
203206
expression_convertor: Arc<dyn ExpressionConvertor>,
204207
pub(crate) vortex_reader_factory: Option<Arc<dyn VortexReaderFactory>>,
205208
pub(crate) ordered: bool,
@@ -233,6 +236,7 @@ impl VortexSource {
233236
_unused_df_metrics: Default::default(),
234237
layout_readers: Arc::new(DashMap::default()),
235238
natural_split_ranges: Arc::new(DashMap::default()),
239+
vortex_files: Arc::new(DashMap::default()),
236240
expression_convertor: Arc::new(DefaultExpressionConvertor::default()),
237241
vortex_reader_factory: None,
238242
vx_metrics_registry: Arc::new(DefaultMetricsRegistry::default()),
@@ -367,6 +371,7 @@ impl VortexSource {
367371
metrics_registry: Arc::clone(&self.vx_metrics_registry),
368372
layout_readers: Arc::clone(&self.layout_readers),
369373
natural_split_ranges: Arc::clone(&self.natural_split_ranges),
374+
vortex_files: Arc::clone(&self.vortex_files),
370375
has_output_ordering: !base_config.output_ordering.is_empty() || self.ordered,
371376
expression_convertor: Arc::clone(&self.expression_convertor),
372377
file_metadata_cache: self.file_metadata_cache.clone(),

vortex-file/src/file.rs

Lines changed: 19 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -30,7 +30,10 @@ use vortex_layout::segments::SegmentInfo;
3030
use vortex_layout::segments::SegmentSource;
3131
use vortex_scan::DataSourceRef;
3232
use vortex_scan::ScanRequest;
33+
use vortex_scan::plan::PreparedStateCache;
34+
use vortex_scan::plan::PreparedStateCacheRef;
3335
use vortex_scan::plan::ScanPlanRef;
36+
use vortex_scan::segments::SegmentFutureCache;
3437
use vortex_session::VortexSession;
3538

3639
use crate::FileStatistics;
@@ -58,6 +61,10 @@ pub struct VortexFile {
5861
layout_reader_cache: Option<OnceLock<Arc<dyn LayoutReader>>>,
5962
/// Shared cache for the v2 physical scan plan root.
6063
scan_plan_root_cache: Arc<OnceLock<ScanPlanRef>>,
64+
/// Shared cache for v2 prepared state across row-range scans of this file.
65+
scan_plan_state_cache: PreparedStateCacheRef,
66+
/// Shared cache for v2 in-flight segment futures across row-range scans of this file.
67+
scan_plan_segment_future_cache: Arc<SegmentFutureCache>,
6168
}
6269

6370
fn layout_reader(
@@ -104,6 +111,8 @@ impl VortexFile {
104111
session,
105112
layout_reader_cache: None,
106113
scan_plan_root_cache: Arc::new(OnceLock::new()),
114+
scan_plan_state_cache: Arc::new(PreparedStateCache::default()),
115+
scan_plan_segment_future_cache: Arc::new(SegmentFutureCache::new()),
107116
}
108117
}
109118

@@ -116,6 +125,8 @@ impl VortexFile {
116125
session: self.session,
117126
layout_reader_cache: Some(OnceLock::new()),
118127
scan_plan_root_cache: self.scan_plan_root_cache,
128+
scan_plan_state_cache: self.scan_plan_state_cache,
129+
scan_plan_segment_future_cache: self.scan_plan_segment_future_cache,
119130
}
120131
}
121132

@@ -203,6 +214,14 @@ impl VortexFile {
203214
Ok(root)
204215
}
205216

217+
pub(crate) fn scan_plan_state_cache(&self) -> PreparedStateCacheRef {
218+
Arc::clone(&self.scan_plan_state_cache)
219+
}
220+
221+
pub(crate) fn scan_plan_segment_future_cache(&self) -> Arc<SegmentFutureCache> {
222+
Arc::clone(&self.scan_plan_segment_future_cache)
223+
}
224+
206225
/// Create a [`DataSource`](vortex_scan::DataSource) from this file for scanning.
207226
///
208227
/// Wraps the file's layout reader with [`FileStatsLayoutReader`] (when file-level

vortex-file/src/multi/scan_v2.rs

Lines changed: 6 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -1089,8 +1089,9 @@ impl<T: Send + 'static> Work<T> {
10891089
) -> Self {
10901090
let known_bytes = registered.bytes();
10911091
let future = async move {
1092-
let _registered = registered;
1093-
future.await
1092+
let result = future.await;
1093+
drop(registered);
1094+
result
10941095
}
10951096
.boxed();
10961097
Self {
@@ -1862,7 +1863,7 @@ impl PreparedScanPlanFile {
18621863
},
18631864
);
18641865
let scheduled_segment_source = Arc::clone(&registered_source.source);
1865-
let segment_future_cache = Arc::new(SegmentFutureCache::new());
1866+
let segment_future_cache = file.scan_plan_segment_future_cache();
18661867
let reader = FileReader::new(
18671868
Arc::new(ScheduledSegmentSourceReader::new(
18681869
segment_source_id,
@@ -1872,7 +1873,8 @@ impl PreparedScanPlanFile {
18721873
session.clone(),
18731874
);
18741875

1875-
let mut prepare_ctx = PrepareCtx::new(session.clone());
1876+
let mut prepare_ctx =
1877+
PrepareCtx::with_state_cache(session.clone(), file.scan_plan_state_cache());
18761878
let projection_pushed = push_expr(&root, &projection, file.dtype(), &session)?;
18771879
let mut split_hints = Vec::new();
18781880
extend_split_hints(&projection_pushed, &mut split_hints);

vortex-layout/src/scan/v2/layouts/flat.rs

Lines changed: 35 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -13,6 +13,7 @@ use std::fmt;
1313
use std::ops::Range;
1414
use std::sync::Arc;
1515

16+
use futures::FutureExt;
1617
use futures::future::BoxFuture;
1718
use parking_lot::Mutex;
1819
use vortex_array::ArrayRef;
@@ -21,6 +22,7 @@ use vortex_array::IntoArray;
2122
use vortex_array::arrays::SliceArray;
2223
use vortex_array::expr::Expression;
2324
use vortex_array::serde::SerializedArray;
25+
use vortex_error::VortexError;
2426
use vortex_error::VortexResult;
2527
use vortex_error::vortex_bail;
2628
use vortex_error::vortex_err;
@@ -44,6 +46,7 @@ use vortex_session::VortexSession;
4446
use crate::layout_v2::Flat;
4547
use crate::layout_v2::Layout;
4648
use crate::layout_v2::LayoutRef;
49+
use crate::layouts::SharedArrayFuture;
4750
use crate::segments::SegmentPlanCtx;
4851
use crate::segments::SegmentRequests;
4952

@@ -63,19 +66,38 @@ pub struct FlatScanPlan {
6366
layout: LayoutRef,
6467
}
6568

66-
/// Per-query cache of the parsed (still lazy) array. Concurrent decodes
67-
/// are benign: the segment fetch is deduplicated by the shared segment
68-
/// source, and last-write-wins on the parsed array.
69+
/// Per-query cache of the parsed (still lazy) array.
6970
#[derive(Default)]
7071
pub struct FlatScanState {
71-
array: Mutex<Option<ArrayRef>>,
72+
array: Mutex<Option<SharedArrayFuture>>,
7273
}
7374

7475
struct FlatPreparedRead {
7576
node: Arc<FlatScanPlan>,
7677
state: Arc<FlatScanState>,
7778
}
7879

80+
impl FlatScanPlan {
81+
fn array(&self, io: &FileReader, state: &FlatScanState) -> SharedArrayFuture {
82+
if let Some(hit) = state.array.lock().clone() {
83+
return hit;
84+
}
85+
86+
let mut guard = state.array.lock();
87+
if let Some(hit) = guard.clone() {
88+
return hit;
89+
}
90+
91+
let layout = self.layout.clone();
92+
let io = io.clone();
93+
let future = async move { decode_flat(&layout, &io).await.map_err(Arc::new) }
94+
.boxed()
95+
.shared();
96+
*guard = Some(future.clone());
97+
future
98+
}
99+
}
100+
79101
impl ScanPlan for FlatScanPlan {
80102
fn init_state(&self, _cx: &mut StateCtx<'_>) -> VortexResult<ScanStateRef> {
81103
Ok(Arc::new(FlatScanState::default()))
@@ -90,7 +112,10 @@ impl ScanPlan for FlatScanPlan {
90112
}
91113

92114
fn prepare_read(self: Arc<Self>, cx: &mut PrepareCtx) -> VortexResult<Option<PreparedReadRef>> {
93-
let key = PreparedStateKey::new::<FlatScanState>(Arc::as_ptr(&self) as *const () as usize);
115+
let flat = self.layout.as_opt::<Flat>().ok_or_else(|| {
116+
vortex_err!("expected flat layout, got {}", self.layout.encoding_id())
117+
})?;
118+
let key = PreparedStateKey::new::<FlatScanState>(*flat.data().segment_id() as usize);
94119
let state = cx.shared_state(key, || Ok(FlatScanState::default()))?;
95120
Ok(Some(Arc::new(FlatPreparedRead { node: self, state })))
96121
}
@@ -120,13 +145,11 @@ impl PreparedRead for FlatPreparedRead {
120145
_local: &'a mut ExecutionCtx,
121146
) -> BoxFuture<'a, VortexResult<ArrayRef>> {
122147
Box::pin(async move {
123-
let array = if let Some(hit) = self.state.array.lock().clone() {
124-
hit
125-
} else {
126-
let decoded = decode_flat(&self.node.layout, io).await?;
127-
*self.state.array.lock() = Some(decoded.clone());
128-
decoded
129-
};
148+
let array = self
149+
.node
150+
.array(io, &self.state)
151+
.await
152+
.map_err(VortexError::from)?;
130153
let dense = slice_to_range(array, &range)?;
131154
if rows.selection.len() != dense.len() {
132155
vortex_bail!(

vortex-scan/src/plan/mod.rs

Lines changed: 14 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1242,6 +1242,20 @@ impl PreparedRead for MaskPreparedRead {
12421242
})
12431243
}
12441244

1245+
fn segment_requests(
1246+
&self,
1247+
range: Range<u64>,
1248+
rows: RowScope<'_>,
1249+
cx: &mut SegmentPlanCtx,
1250+
) -> VortexResult<SegmentRequests> {
1251+
let mut requests = self.input.segment_requests(range.clone(), rows, cx)?;
1252+
if requests.is_unknown() {
1253+
return Ok(requests);
1254+
}
1255+
requests.extend(self.validity.segment_requests(range, rows, cx)?);
1256+
Ok(requests)
1257+
}
1258+
12451259
fn release(&self, frontier: u64) -> VortexResult<()> {
12461260
self.input.release(frontier)?;
12471261
self.validity.release(frontier)

0 commit comments

Comments
 (0)