Skip to main content

cobble_data_structure/
structured_scan.rs

1use crate::structured_db::{
2    StructuredColumnFamilySchema, StructuredColumnValue, StructuredScanOptions, combined_resolver,
3    decode_row, load_structured_schema_from_cobble_schema,
4};
5use bytes::Bytes;
6use cobble::{
7    Config, GlobalSnapshotManifest, MergeOperatorResolver, ReadOnlyDb, Result, ScanPlan, ScanSplit,
8    ScanSplitScanner, ShardSnapshotRef,
9};
10use serde::{Deserialize, Serialize};
11use std::sync::Arc;
12
13/// Structured distributed scan plan.
14///
15/// Wraps `cobble::ScanPlan` and produces structured scan splits/scanners.
16pub struct StructuredScanPlan {
17    inner: ScanPlan,
18}
19
20impl StructuredScanPlan {
21    pub fn new(manifest: GlobalSnapshotManifest) -> Self {
22        Self {
23            inner: ScanPlan::new(manifest),
24        }
25    }
26
27    pub fn with_start(mut self, start: Vec<u8>) -> Self {
28        self.inner = self.inner.with_start(start);
29        self
30    }
31
32    pub fn with_end(mut self, end: Vec<u8>) -> Self {
33        self.inner = self.inner.with_end(end);
34        self
35    }
36
37    pub fn splits(&self) -> Vec<StructuredScanSplit> {
38        self.inner.splits().into_iter().map(Into::into).collect()
39    }
40}
41
42/// Structured version of a distributed scan split.
43#[derive(Clone, Debug, Serialize, Deserialize)]
44pub struct StructuredScanSplit {
45    pub shard: ShardSnapshotRef,
46    pub start: Option<Vec<u8>>,
47    pub end: Option<Vec<u8>>,
48    pub start_bucket: Option<u16>,
49    pub start_key_exclusive: Option<Vec<u8>>,
50    pub end_bucket: Option<u16>,
51    pub end_key_inclusive: Option<Vec<u8>>,
52}
53
54pub struct StructuredScanSplitPartition {
55    pub before: StructuredScanSplit,
56    pub after: StructuredScanSplit,
57}
58
59impl From<ScanSplit> for StructuredScanSplit {
60    fn from(value: ScanSplit) -> Self {
61        Self {
62            shard: value.shard,
63            start: value.start,
64            end: value.end,
65            start_bucket: value.start_bucket,
66            start_key_exclusive: value.start_key_exclusive,
67            end_bucket: value.end_bucket,
68            end_key_inclusive: value.end_key_inclusive,
69        }
70    }
71}
72
73impl From<StructuredScanSplit> for ScanSplit {
74    fn from(value: StructuredScanSplit) -> Self {
75        Self {
76            shard: value.shard,
77            start: value.start,
78            end: value.end,
79            start_bucket: value.start_bucket,
80            start_key_exclusive: value.start_key_exclusive,
81            end_bucket: value.end_bucket,
82            end_key_inclusive: value.end_key_inclusive,
83        }
84    }
85}
86
87impl StructuredScanSplit {
88    pub fn split_after(
89        &self,
90        bucket: u16,
91        key_inclusive: Vec<u8>,
92    ) -> Result<StructuredScanSplitPartition> {
93        let partition = ScanSplit::from(self.clone()).split_after(bucket, key_inclusive)?;
94        Ok(StructuredScanSplitPartition {
95            before: partition.before.into(),
96            after: partition.after.into(),
97        })
98    }
99
100    pub fn create_scanner_without_options(
101        &self,
102        config: Config,
103    ) -> Result<StructuredScanSplitScanner> {
104        self.create_scanner_without_options_internal(config, None)
105    }
106
107    pub fn create_scanner(
108        &self,
109        config: Config,
110        options: &StructuredScanOptions,
111    ) -> Result<StructuredScanSplitScanner> {
112        self.create_scanner_internal(config, None, options)
113    }
114
115    pub fn create_scanner_with_resolver_without_options(
116        &self,
117        config: Config,
118        resolver: Arc<dyn MergeOperatorResolver>,
119    ) -> Result<StructuredScanSplitScanner> {
120        self.create_scanner_without_options_internal(config, Some(resolver))
121    }
122
123    pub fn create_scanner_with_resolver(
124        &self,
125        config: Config,
126        resolver: Arc<dyn MergeOperatorResolver>,
127        options: &StructuredScanOptions,
128    ) -> Result<StructuredScanSplitScanner> {
129        self.create_scanner_internal(config, Some(resolver), options)
130    }
131
132    fn create_scanner_internal(
133        &self,
134        config: Config,
135        resolver: Option<Arc<dyn MergeOperatorResolver>>,
136        options: &StructuredScanOptions,
137    ) -> Result<StructuredScanSplitScanner> {
138        let resolver = combined_resolver(resolver);
139        let read_only = ReadOnlyDb::open_with_db_id_and_resolver(
140            config.clone(),
141            self.shard.snapshot_id,
142            self.shard.db_id.clone(),
143            Arc::clone(&resolver),
144        )?;
145        let structured_schema = Arc::new(load_structured_schema_from_cobble_schema(
146            &read_only.current_schema(),
147        )?);
148        let projected_schema = options.resolve_projected_schema_cached(&structured_schema)?;
149        let scanner = ScanSplit::from(self.clone()).create_scanner_with_resolver(
150            config,
151            resolver,
152            options.as_cobble(),
153        )?;
154        Ok(StructuredScanSplitScanner {
155            inner: scanner,
156            structured_schema: projected_schema,
157        })
158    }
159
160    fn create_scanner_without_options_internal(
161        &self,
162        config: Config,
163        resolver: Option<Arc<dyn MergeOperatorResolver>>,
164    ) -> Result<StructuredScanSplitScanner> {
165        let resolver = combined_resolver(resolver);
166        let read_only = ReadOnlyDb::open_with_db_id_and_resolver(
167            config.clone(),
168            self.shard.snapshot_id,
169            self.shard.db_id.clone(),
170            Arc::clone(&resolver),
171        )?;
172        let structured_schema = Arc::new(load_structured_schema_from_cobble_schema(
173            &read_only.current_schema(),
174        )?);
175        let projected_schema = Arc::new(structured_schema.projected(0, None));
176        let scanner = ScanSplit::from(self.clone())
177            .create_scanner_with_resolver_without_options(config, resolver)?;
178        Ok(StructuredScanSplitScanner {
179            inner: scanner,
180            structured_schema: projected_schema,
181        })
182    }
183}
184
185pub struct StructuredScanSplitScanner {
186    inner: ScanSplitScanner,
187    structured_schema: Arc<StructuredColumnFamilySchema>,
188}
189
190impl StructuredScanSplitScanner {
191    pub fn consume_next_row<T, F>(&mut self, mut consumer: F) -> Result<Option<T>>
192    where
193        F: FnMut(&Bytes, &[Option<StructuredColumnValue>]) -> Result<T>,
194    {
195        let structured_schema = Arc::clone(&self.structured_schema);
196        self.inner.consume_next_row(|key, columns| {
197            let decoded = decode_row(&structured_schema, 0, columns.to_vec())?;
198            consumer(key, &decoded)
199        })
200    }
201
202    pub fn consume_next_row_with_bucket<T, F>(&mut self, mut consumer: F) -> Result<Option<T>>
203    where
204        F: FnMut(u16, &Bytes, &[Option<StructuredColumnValue>]) -> Result<T>,
205    {
206        let structured_schema = Arc::clone(&self.structured_schema);
207        self.inner
208            .consume_next_row_with_bucket(|bucket, key, columns| {
209                let decoded = decode_row(&structured_schema, 0, columns.to_vec())?;
210                consumer(bucket, key, &decoded)
211            })
212    }
213}
214
215impl Iterator for StructuredScanSplitScanner {
216    type Item = Result<(u16, Bytes, Vec<Option<StructuredColumnValue>>)>;
217
218    fn next(&mut self) -> Option<Self::Item> {
219        self.inner.next().map(|item| {
220            let (bucket, key, columns) = item?;
221            let decoded = decode_row(&self.structured_schema, 0, columns)?;
222            Ok((bucket, key, decoded))
223        })
224    }
225}
226
227#[cfg(test)]
228#[path = "../tests/unit/structured_scan.rs"]
229mod tests;