uqa_execution/
spill_scan.rs1use crate::batch::{Batch, RowSchema};
14use crate::physical::{BackwardScanSupport, ExecResult, PhysicalOperator};
15use crate::spill::{SharedSpill, SharedSpillReader, SpillBuffer, SpillDrain};
16
17pub struct SpillScan {
19 schema: RowSchema,
20 buffer: Option<SpillBuffer>,
21 reader: Option<SpillDrain>,
22}
23
24pub struct SharedSpillScan {
27 source: SharedSpill,
28 schema: RowSchema,
29 reader: Option<SharedSpillReader>,
30}
31
32impl SharedSpillScan {
33 pub fn new(source: SharedSpill) -> Self {
34 let schema = source.row_schema().clone();
35 Self {
36 source,
37 schema,
38 reader: None,
39 }
40 }
41}
42
43impl PhysicalOperator for SharedSpillScan {
44 fn row_schema(&self) -> &RowSchema {
45 &self.schema
46 }
47
48 fn estimated_cardinality(&self) -> Option<u64> {
49 u64::try_from(self.source.rows()).ok()
50 }
51
52 fn backward_scan_support(&self) -> BackwardScanSupport {
53 BackwardScanSupport::Materialize
54 }
55
56 fn open(&mut self) -> ExecResult<()> {
57 self.reader = Some(self.source.reader()?);
58 Ok(())
59 }
60
61 fn next(&mut self) -> ExecResult<Option<Batch>> {
62 self.reader
63 .as_mut()
64 .map_or(Ok(None), |reader| reader.next().transpose())
65 }
66
67 fn close(&mut self) -> ExecResult<()> {
68 self.reader = None;
69 Ok(())
70 }
71}
72
73impl SpillScan {
74 pub fn new(schema: impl Into<RowSchema>, buffer: SpillBuffer) -> Self {
75 let schema = schema.into();
76 Self {
77 schema,
78 buffer: Some(buffer),
79 reader: None,
80 }
81 }
82}
83
84impl PhysicalOperator for SpillScan {
85 fn row_schema(&self) -> &RowSchema {
86 &self.schema
87 }
88
89 fn open(&mut self) -> ExecResult<()> {
90 let mut buffer = self.buffer.take().ok_or_else(|| {
91 crate::physical::ExecError::Other("spill scan cannot be reopened".into())
92 })?;
93 self.reader = Some(buffer.drain()?);
94 Ok(())
95 }
96
97 fn next(&mut self) -> ExecResult<Option<Batch>> {
98 let Some(reader) = self.reader.as_mut() else {
99 return Ok(None);
100 };
101 let Some(batch) = reader.next().transpose()? else {
102 return Ok(None);
103 };
104 if batch.schema != self.schema {
105 return Err(crate::physical::ExecError::Other(format!(
106 "spill scan schema mismatch: expected {:?}, got {:?}",
107 self.schema.columns(),
108 batch.schema.columns()
109 )));
110 }
111 Ok(Some(batch))
112 }
113
114 fn close(&mut self) -> ExecResult<()> {
115 self.reader = None;
116 Ok(())
117 }
118}
119
120#[cfg(test)]
121mod tests {
122 use std::collections::BTreeMap;
123
124 use uqa_core::Value;
125
126 use super::*;
127 use crate::physical::run_to_rows;
128
129 #[test]
130 fn scan_streams_a_forced_spill_in_input_order() {
131 let schema = RowSchema::new(vec!["x".into()]);
132 let mut spill = SpillBuffer::new(1);
133 for value in 0..300_i64 {
134 spill
135 .push(Batch::new(
136 schema.clone(),
137 vec![BTreeMap::from([("x".into(), Value::Int(value))])],
138 ))
139 .unwrap();
140 }
141 assert!(spill.has_spilled());
142
143 let mut scan = SpillScan::new(schema.columns().to_vec(), spill);
144 let (_, rows) = run_to_rows(&mut scan).unwrap();
145 assert_eq!(rows.len(), 300);
146 for (expected, row) in rows.iter().enumerate() {
147 assert_eq!(row.get("x"), Some(&Value::Int(expected as i64)));
148 }
149 }
150
151 #[test]
152 fn shared_spill_supports_independent_repeatable_scans() {
153 let schema = RowSchema::new(vec!["x".into()]);
154 let mut spill = SpillBuffer::new(1);
155 for value in 0..2_048_i64 {
156 spill
157 .push(Batch::new(
158 schema.clone(),
159 vec![BTreeMap::from([("x".into(), Value::Int(value))])],
160 ))
161 .unwrap();
162 }
163 let shared = spill.into_shared(schema.columns().to_vec()).unwrap();
164 let mut first = SharedSpillScan::new(shared.clone());
165 let mut second = SharedSpillScan::new(shared);
166 let (_, first_rows) = run_to_rows(&mut first).unwrap();
167 let (_, second_rows) = run_to_rows(&mut second).unwrap();
168 assert_eq!(first_rows, second_rows);
169 assert_eq!(first_rows.len(), 2_048);
170 }
171}