1use std::collections::HashMap;
16use std::fmt::{Display, Formatter};
17
18use common_error::ext::BoxedError;
19use common_recordbatch::OrderOption;
20use datafusion_expr::expr::Expr;
21pub use datatypes::schema::{VectorDistanceMetric, VectorIndexEngineType};
23use datatypes::types::json_type::JsonNativeType;
24use itertools::Itertools;
25use strum::Display;
26
27use crate::storage::{ColumnId, SequenceNumber};
28
29#[derive(Debug, Clone, PartialEq)]
31pub struct VectorSearchRequest {
32 pub column_id: ColumnId,
34 pub query_vector: Vec<f32>,
36 pub k: usize,
38 pub metric: VectorDistanceMetric,
40}
41
42#[derive(Debug, Clone, PartialEq)]
44pub struct VectorSearchMatches {
45 pub keys: Vec<u64>,
47 pub distances: Vec<f32>,
49}
50
51pub trait VectorIndexEngine: Send + Sync {
56 fn add(&mut self, key: u64, vector: &[f32]) -> Result<(), BoxedError>;
58
59 fn search(&self, query: &[f32], k: usize) -> Result<VectorSearchMatches, BoxedError>;
61
62 fn serialized_length(&self) -> usize;
64
65 fn save_to_buffer(&self, buffer: &mut [u8]) -> Result<(), BoxedError>;
67
68 fn reserve(&mut self, capacity: usize) -> Result<(), BoxedError>;
70
71 fn size(&self) -> usize;
73
74 fn capacity(&self) -> usize;
76
77 fn memory_usage(&self) -> usize;
79}
80
81#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Display)]
83pub enum TimeSeriesRowSelector {
84 LastRow,
86}
87
88#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash, Display)]
90pub enum TimeSeriesDistribution {
91 TimeWindowed,
94 PerSeries,
97}
98
99#[derive(Default, Clone, Debug, PartialEq)]
100pub struct ScanRequest {
101 pub projection: Option<Vec<usize>>,
104 pub filters: Vec<Expr>,
106 pub output_ordering: Option<Vec<OrderOption>>,
108 pub limit: Option<usize>,
113 pub series_row_selector: Option<TimeSeriesRowSelector>,
115 pub memtable_max_sequence: Option<SequenceNumber>,
121 pub memtable_min_sequence: Option<SequenceNumber>,
124 pub sst_min_sequence: Option<SequenceNumber>,
127 pub skip_sst_files: bool,
130 pub snapshot_on_scan: bool,
132 pub distribution: Option<TimeSeriesDistribution>,
134 pub vector_search: Option<VectorSearchRequest>,
137 pub json_type_hint: HashMap<String, JsonNativeType>,
139 pub preserve_pk_dictionary_encoding: bool,
141}
142
143impl Display for ScanRequest {
144 fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result {
145 enum Delimiter {
146 None,
147 Init,
148 }
149
150 impl Delimiter {
151 fn as_str(&mut self) -> &str {
152 match self {
153 Delimiter::None => {
154 *self = Delimiter::Init;
155 ""
156 }
157 Delimiter::Init => ", ",
158 }
159 }
160 }
161
162 let mut delimiter = Delimiter::None;
163
164 write!(f, "ScanRequest {{ ")?;
165 if let Some(projection) = &self.projection {
166 write!(f, "{}projection: {:?}", delimiter.as_str(), projection)?;
167 }
168 if !self.filters.is_empty() {
169 write!(
170 f,
171 "{}filters: [{}]",
172 delimiter.as_str(),
173 self.filters
174 .iter()
175 .map(|f| f.to_string())
176 .collect::<Vec<_>>()
177 .join(", ")
178 )?;
179 }
180 if let Some(output_ordering) = &self.output_ordering {
181 write!(
182 f,
183 "{}output_ordering: {:?}",
184 delimiter.as_str(),
185 output_ordering
186 )?;
187 }
188 if let Some(limit) = &self.limit {
189 write!(f, "{}limit: {}", delimiter.as_str(), limit)?;
190 }
191 if let Some(series_row_selector) = &self.series_row_selector {
192 write!(
193 f,
194 "{}series_row_selector: {}",
195 delimiter.as_str(),
196 series_row_selector
197 )?;
198 }
199 if let Some(sequence) = &self.memtable_max_sequence {
200 write!(f, "{}sequence: {}", delimiter.as_str(), sequence)?;
201 }
202 if let Some(sst_min_sequence) = &self.sst_min_sequence {
203 write!(
204 f,
205 "{}sst_min_sequence: {}",
206 delimiter.as_str(),
207 sst_min_sequence
208 )?;
209 }
210 if self.skip_sst_files {
211 write!(
212 f,
213 "{}skip_sst_files: {}",
214 delimiter.as_str(),
215 self.skip_sst_files
216 )?;
217 }
218 if self.snapshot_on_scan {
219 write!(
220 f,
221 "{}snapshot_on_scan: {}",
222 delimiter.as_str(),
223 self.snapshot_on_scan
224 )?;
225 }
226 if self.preserve_pk_dictionary_encoding {
227 write!(
228 f,
229 "{}preserve_pk_dictionary_encoding: true",
230 delimiter.as_str()
231 )?;
232 }
233 if let Some(distribution) = &self.distribution {
234 write!(f, "{}distribution: {}", delimiter.as_str(), distribution)?;
235 }
236 if let Some(vector_search) = &self.vector_search {
237 write!(
238 f,
239 "{}vector_search: column_id={}, k={}, metric={}",
240 delimiter.as_str(),
241 vector_search.column_id,
242 vector_search.k,
243 vector_search.metric
244 )?;
245 }
246 if !self.json_type_hint.is_empty() {
247 write!(
248 f,
249 "{}json_type_hint: {}",
250 delimiter.as_str(),
251 self.json_type_hint
252 .iter()
253 .map(|(column, json_type)| format!("({column}: {json_type})"))
254 .join(", ")
255 )?;
256 }
257 write!(f, " }}")
258 }
259}
260
261#[cfg(test)]
262mod tests {
263 use datafusion_expr::{Operator, binary_expr, col, lit};
264
265 use super::*;
266
267 #[test]
268 fn test_display_scan_request() {
269 let request = ScanRequest {
270 ..Default::default()
271 };
272 assert_eq!(request.to_string(), "ScanRequest { }");
273
274 let projection = Some(vec![1, 2]);
275 let request = ScanRequest {
276 projection,
277 filters: vec![
278 binary_expr(col("i"), Operator::Gt, lit(1)),
279 binary_expr(col("s"), Operator::Eq, lit("x")),
280 ],
281 limit: Some(10),
282 ..Default::default()
283 };
284 assert_eq!(
285 request.to_string(),
286 r#"ScanRequest { projection: [1, 2], filters: [i > Int32(1), s = Utf8("x")], limit: 10 }"#
287 );
288
289 let request = ScanRequest {
290 filters: vec![
291 binary_expr(col("i"), Operator::Gt, lit(1)),
292 binary_expr(col("s"), Operator::Eq, lit("x")),
293 ],
294 limit: Some(10),
295 ..Default::default()
296 };
297 assert_eq!(
298 request.to_string(),
299 r#"ScanRequest { filters: [i > Int32(1), s = Utf8("x")], limit: 10 }"#
300 );
301
302 let projection = Some(vec![1, 2]);
303 let request = ScanRequest {
304 projection,
305 limit: Some(10),
306 ..Default::default()
307 };
308 assert_eq!(
309 request.to_string(),
310 "ScanRequest { projection: [1, 2], limit: 10 }"
311 );
312
313 let request = ScanRequest {
314 snapshot_on_scan: true,
315 ..Default::default()
316 };
317 assert_eq!(
318 request.to_string(),
319 "ScanRequest { snapshot_on_scan: true }"
320 );
321
322 let request = ScanRequest {
323 skip_sst_files: true,
324 ..Default::default()
325 };
326 assert_eq!(request.to_string(), "ScanRequest { skip_sst_files: true }");
327 }
328}