1use std::collections::HashMap;
16use std::fmt::{Display, Formatter};
17
18use common_recordbatch::OrderOption;
19use datafusion_expr::expr::Expr;
20use datatypes::types::json_type::JsonNativeType;
21use itertools::Itertools;
22use strum::Display;
23
24use crate::storage::SequenceNumber;
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Display)]
28pub enum TimeSeriesRowSelector {
29 #[strum(to_string = "LastRow {{ after_merge: {after_merge} }}")]
31 LastRow {
32 after_merge: bool,
34 },
35}
36
37#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Display)]
39pub enum TimeSeriesDistribution {
40 TimeWindowed,
43 PerSeries,
46}
47
48#[derive(Default, Clone, Debug, PartialEq)]
49pub struct ScanRequest {
50 pub projection: Option<Vec<usize>>,
53 pub filters: Vec<Expr>,
55 pub output_ordering: Option<Vec<OrderOption>>,
57 pub limit: Option<usize>,
62 pub series_row_selector: Option<TimeSeriesRowSelector>,
64 pub memtable_max_sequence: Option<SequenceNumber>,
70 pub memtable_min_sequence: Option<SequenceNumber>,
73 pub sst_min_sequence: Option<SequenceNumber>,
76 pub skip_sst_files: bool,
79 pub snapshot_on_scan: bool,
81 pub exact_sequence_range: bool,
90 pub distribution: Option<TimeSeriesDistribution>,
92 pub json_type_hint: HashMap<String, JsonNativeType>,
94 pub preserve_pk_dictionary_encoding: bool,
96}
97
98impl Display for ScanRequest {
99 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
100 enum Delimiter {
101 None,
102 Init,
103 }
104
105 impl Delimiter {
106 fn as_str(&mut self) -> &str {
107 match self {
108 Delimiter::None => {
109 *self = Delimiter::Init;
110 ""
111 }
112 Delimiter::Init => ", ",
113 }
114 }
115 }
116
117 let mut delimiter = Delimiter::None;
118
119 write!(f, "ScanRequest {{ ")?;
120 if let Some(projection) = &self.projection {
121 write!(f, "{}projection: {:?}", delimiter.as_str(), projection)?;
122 }
123 if !self.filters.is_empty() {
124 write!(
125 f,
126 "{}filters: [{}]",
127 delimiter.as_str(),
128 self.filters
129 .iter()
130 .map(|f| f.to_string())
131 .collect::<Vec<_>>()
132 .join(", ")
133 )?;
134 }
135 if let Some(output_ordering) = &self.output_ordering {
136 write!(
137 f,
138 "{}output_ordering: {:?}",
139 delimiter.as_str(),
140 output_ordering
141 )?;
142 }
143 if let Some(limit) = &self.limit {
144 write!(f, "{}limit: {}", delimiter.as_str(), limit)?;
145 }
146 if let Some(series_row_selector) = &self.series_row_selector {
147 write!(
148 f,
149 "{}series_row_selector: {}",
150 delimiter.as_str(),
151 series_row_selector
152 )?;
153 }
154 if let Some(sequence) = &self.memtable_max_sequence {
155 write!(f, "{}sequence: {}", delimiter.as_str(), sequence)?;
156 }
157 if let Some(sst_min_sequence) = &self.sst_min_sequence {
158 write!(
159 f,
160 "{}sst_min_sequence: {}",
161 delimiter.as_str(),
162 sst_min_sequence
163 )?;
164 }
165 if self.skip_sst_files {
166 write!(
167 f,
168 "{}skip_sst_files: {}",
169 delimiter.as_str(),
170 self.skip_sst_files
171 )?;
172 }
173 if self.snapshot_on_scan {
174 write!(
175 f,
176 "{}snapshot_on_scan: {}",
177 delimiter.as_str(),
178 self.snapshot_on_scan
179 )?;
180 }
181 if self.exact_sequence_range {
182 write!(f, "{}exact_sequence_range: true", delimiter.as_str())?;
183 }
184 if self.preserve_pk_dictionary_encoding {
185 write!(
186 f,
187 "{}preserve_pk_dictionary_encoding: true",
188 delimiter.as_str()
189 )?;
190 }
191 if let Some(distribution) = &self.distribution {
192 write!(f, "{}distribution: {}", delimiter.as_str(), distribution)?;
193 }
194 if !self.json_type_hint.is_empty() {
195 write!(
196 f,
197 "{}json_type_hint: {}",
198 delimiter.as_str(),
199 self.json_type_hint
200 .iter()
201 .map(|(column, json_type)| format!("({column}: {json_type})"))
202 .join(", ")
203 )?;
204 }
205 write!(f, " }}")
206 }
207}
208
209#[cfg(test)]
210mod tests {
211 use datafusion_expr::{Operator, binary_expr, col, lit};
212
213 use super::*;
214
215 #[test]
216 fn test_display_scan_request() {
217 let request = ScanRequest {
218 ..Default::default()
219 };
220 assert_eq!(request.to_string(), "ScanRequest { }");
221
222 let projection = Some(vec![1, 2]);
223 let request = ScanRequest {
224 projection,
225 filters: vec![
226 binary_expr(col("i"), Operator::Gt, lit(1)),
227 binary_expr(col("s"), Operator::Eq, lit("x")),
228 ],
229 limit: Some(10),
230 ..Default::default()
231 };
232 assert_eq!(
233 request.to_string(),
234 r#"ScanRequest { projection: [1, 2], filters: [i > Int32(1), s = Utf8("x")], limit: 10 }"#
235 );
236
237 let request = ScanRequest {
238 filters: vec![
239 binary_expr(col("i"), Operator::Gt, lit(1)),
240 binary_expr(col("s"), Operator::Eq, lit("x")),
241 ],
242 limit: Some(10),
243 ..Default::default()
244 };
245 assert_eq!(
246 request.to_string(),
247 r#"ScanRequest { filters: [i > Int32(1), s = Utf8("x")], limit: 10 }"#
248 );
249
250 let projection = Some(vec![1, 2]);
251 let request = ScanRequest {
252 projection,
253 limit: Some(10),
254 ..Default::default()
255 };
256 assert_eq!(
257 request.to_string(),
258 "ScanRequest { projection: [1, 2], limit: 10 }"
259 );
260
261 let request = ScanRequest {
262 series_row_selector: Some(TimeSeriesRowSelector::LastRow { after_merge: true }),
263 snapshot_on_scan: true,
264 exact_sequence_range: true,
265 ..Default::default()
266 };
267 assert_eq!(
268 request.to_string(),
269 "ScanRequest { series_row_selector: LastRow { after_merge: true }, snapshot_on_scan: true, exact_sequence_range: true }"
270 );
271
272 assert_eq!(
273 TimeSeriesRowSelector::LastRow { after_merge: false }.to_string(),
274 "LastRow { after_merge: false }"
275 );
276
277 let request = ScanRequest {
278 skip_sst_files: true,
279 ..Default::default()
280 };
281 assert_eq!(request.to_string(), "ScanRequest { skip_sst_files: true }");
282 }
283}