1use alopex_core::columnar::encoding::Column;
2use alopex_core::columnar::encoding_v2::Bitmap;
3use alopex_core::columnar::kvs_bridge::key_layout;
4use alopex_core::columnar::segment_v2::{
5 ColumnSegmentV2, InMemorySegmentSource, RecordBatch, SegmentReaderV2,
6};
7use alopex_core::kv::{KVStore, KVTransaction};
8use alopex_core::storage::format::bincode_config;
9use bincode::config::Options;
10
11use crate::ast::expr::BinaryOp;
12use crate::catalog::{ColumnMetadata, RowIdMode, TableMetadata};
13use crate::columnar::statistics::RowGroupStatistics;
14use crate::executor::evaluator::{EvalContext, evaluate};
15use crate::executor::query::iterator::RowIterator;
16use crate::executor::{ExecutorError, Result, Row};
17use crate::planner::typed_expr::{Projection, TypedExpr, TypedExprKind};
18use crate::planner::types::ResolvedType;
19use crate::storage::{SqlTxn, SqlValue};
20use std::collections::BTreeSet;
21
22#[derive(Debug, Clone)]
24pub struct ColumnarScan {
25 pub table_id: u32,
26 pub projected_columns: Vec<usize>,
27 pub pushed_filter: Option<PushdownFilter>,
28 pub residual_filter: Option<TypedExpr>,
29}
30
31#[derive(Debug, Clone, PartialEq)]
33pub enum PushdownFilter {
34 Eq {
35 column_idx: usize,
36 value: SqlValue,
37 },
38 Range {
39 column_idx: usize,
40 min: Option<SqlValue>,
41 max: Option<SqlValue>,
42 },
43 IsNull {
44 column_idx: usize,
45 is_null: bool,
46 },
47 And(Vec<PushdownFilter>),
48 Or(Vec<PushdownFilter>),
49}
50
51impl ColumnarScan {
52 pub fn new(
53 table_id: u32,
54 projected_columns: Vec<usize>,
55 pushed_filter: Option<PushdownFilter>,
56 residual_filter: Option<TypedExpr>,
57 ) -> Self {
58 Self {
59 table_id,
60 projected_columns,
61 pushed_filter,
62 residual_filter,
63 }
64 }
65
66 pub fn should_skip_row_group(&self, stats: &RowGroupStatistics) -> bool {
68 match &self.pushed_filter {
69 None => false,
70 Some(filter) => Self::evaluate_pushdown(filter, stats),
71 }
72 }
73
74 pub fn evaluate_pushdown(filter: &PushdownFilter, stats: &RowGroupStatistics) -> bool {
76 match filter {
77 PushdownFilter::Eq { column_idx, value } => match stats.columns.get(*column_idx) {
78 Some(col_stats) => {
79 if col_stats.total_count == 0 {
80 return true;
81 }
82 if matches!(
83 value.partial_cmp(&col_stats.min),
84 Some(std::cmp::Ordering::Less)
85 ) {
86 return true;
87 }
88 matches!(
89 value.partial_cmp(&col_stats.max),
90 Some(std::cmp::Ordering::Greater)
91 )
92 }
93 None => false,
94 },
95
96 PushdownFilter::Range {
97 column_idx,
98 min,
99 max,
100 } => match stats.columns.get(*column_idx) {
101 Some(col_stats) => {
102 if col_stats.total_count == 0 {
103 return true;
104 }
105 if let Some(filter_min) = min
106 && matches!(
107 col_stats.max.partial_cmp(filter_min),
108 Some(std::cmp::Ordering::Less)
109 )
110 {
111 return true;
112 }
113 if let Some(filter_max) = max
114 && matches!(
115 col_stats.min.partial_cmp(filter_max),
116 Some(std::cmp::Ordering::Greater)
117 )
118 {
119 return true;
120 }
121 false
122 }
123 None => false,
124 },
125
126 PushdownFilter::IsNull {
127 column_idx,
128 is_null,
129 } => match stats.columns.get(*column_idx) {
130 Some(col_stats) => {
131 if *is_null {
132 col_stats.null_count == 0
133 } else {
134 col_stats.null_count == col_stats.total_count
135 }
136 }
137 None => false,
138 },
139
140 PushdownFilter::And(filters) => {
141 if filters.is_empty() {
142 return false;
143 }
144 filters.iter().any(|f| Self::evaluate_pushdown(f, stats))
145 }
146
147 PushdownFilter::Or(filters) => {
148 if filters.is_empty() {
149 return false;
150 }
151 filters.iter().all(|f| Self::evaluate_pushdown(f, stats))
152 }
153 }
154 }
155}
156
157struct LoadedSegment {
163 reader: SegmentReaderV2,
165 row_group_stats: Option<Vec<RowGroupStatistics>>,
167 row_ids: Vec<u64>,
169 row_groups: Vec<alopex_core::columnar::segment_v2::RowGroupMeta>,
171}
172
173pub struct ColumnarScanIterator {
180 segments: Vec<LoadedSegment>,
182 segment_idx: usize,
184 row_group_idx: usize,
186 row_idx: usize,
188 current_batch: Option<RecordBatch>,
190 projected: Vec<usize>,
192 table_meta: TableMetadata,
194 schema: Vec<ColumnMetadata>,
196 scan: ColumnarScan,
198 row_id_col_idx: Option<usize>,
200 next_row_id: u64,
202}
203
204impl ColumnarScanIterator {
205 fn advance(&mut self) -> Option<Result<Row>> {
210 loop {
211 if self.current_batch.is_none() && !self.load_next_batch() {
213 return None; }
215
216 let row_count = match &self.current_batch {
218 Some(batch) => batch.num_rows(),
219 None => continue,
220 };
221
222 if self.row_idx >= row_count {
224 self.current_batch = None;
225 self.row_idx = 0;
226 self.row_group_idx += 1;
227 continue;
228 }
229
230 let row_idx = self.row_idx;
232 match self.convert_current_row(row_idx) {
233 Ok(Some(row)) => {
234 self.row_idx += 1;
235 return Some(Ok(row));
236 }
237 Ok(None) => {
238 self.row_idx += 1;
240 continue;
241 }
242 Err(e) => {
243 self.row_idx += 1;
244 return Some(Err(e));
245 }
246 }
247 }
248 }
249
250 fn load_next_batch(&mut self) -> bool {
254 while self.segment_idx < self.segments.len() {
255 let segment = &self.segments[self.segment_idx];
256 let row_group_count = segment.row_groups.len();
257
258 while self.row_group_idx < row_group_count {
259 let should_skip = match segment.row_group_stats.as_ref() {
261 Some(stats) if stats.len() == row_group_count => {
262 self.scan.should_skip_row_group(&stats[self.row_group_idx])
263 }
264 _ => false,
265 };
266
267 if should_skip {
268 self.row_group_idx += 1;
269 continue;
270 }
271
272 match segment
274 .reader
275 .read_row_group_by_index(&self.projected, self.row_group_idx)
276 {
277 Ok(mut batch) => {
278 if !segment.row_ids.is_empty()
280 && let Some(meta) = segment.row_groups.get(self.row_group_idx)
281 {
282 let start = meta.row_start as usize;
283 let end = start + meta.row_count as usize;
284 if end <= segment.row_ids.len() {
285 batch =
286 batch.with_row_ids(Some(segment.row_ids[start..end].to_vec()));
287 }
288 }
289 self.current_batch = Some(batch);
290 self.row_idx = 0;
291 return true;
292 }
293 Err(_) => {
294 self.row_group_idx += 1;
296 continue;
297 }
298 }
299 }
300
301 self.segment_idx += 1;
303 self.row_group_idx = 0;
304 }
305
306 false
307 }
308
309 fn convert_current_row(&mut self, row_idx: usize) -> Result<Option<Row>> {
313 let batch = self
314 .current_batch
315 .as_ref()
316 .ok_or_else(|| ExecutorError::Columnar("no current batch".into()))?;
317
318 let column_count = self.table_meta.column_count();
319 let mut values = vec![SqlValue::Null; column_count];
320
321 for (pos, &table_col_idx) in self.projected.iter().enumerate() {
322 let column = batch
323 .columns
324 .get(pos)
325 .ok_or_else(|| ExecutorError::Columnar("missing projected column".into()))?;
326 let bitmap = batch.null_bitmaps.get(pos).and_then(|b| b.as_ref());
327 let col_meta = self
328 .table_meta
329 .columns
330 .get(table_col_idx)
331 .ok_or_else(|| ExecutorError::Columnar("column index out of bounds".into()))?;
332 let value = value_from_column(column, bitmap, row_idx, &col_meta.data_type)?;
333 values[table_col_idx] = value;
334 }
335
336 if let Some(predicate) = self.scan.residual_filter.as_ref() {
338 let ctx = EvalContext::new(&values);
339 let keep = matches!(evaluate(predicate, &ctx)?, SqlValue::Boolean(true));
340 if !keep {
341 return Ok(None);
342 }
343 }
344
345 let batch = self
347 .current_batch
348 .as_ref()
349 .ok_or_else(|| ExecutorError::Columnar("no current batch".into()))?;
350
351 let row_id = match self.table_meta.storage_options.row_id_mode {
352 RowIdMode::Direct => {
353 if let Some(row_ids) = batch.row_ids.as_ref() {
354 *row_ids.get(row_idx).ok_or_else(|| {
355 ExecutorError::Columnar(
356 "row_id missing for row in row_id_mode=direct".into(),
357 )
358 })?
359 } else if let Some(idx) = self.row_id_col_idx {
360 let val = values.get(idx).ok_or_else(|| {
361 ExecutorError::Columnar("row_id column missing in projected values".into())
362 })?;
363 match val {
364 SqlValue::Integer(v) if *v >= 0 => *v as u64,
365 SqlValue::BigInt(v) if *v >= 0 => *v as u64,
366 other => {
367 return Err(ExecutorError::Columnar(format!(
368 "row_id column must be non-negative integer, got {}",
369 other.type_name()
370 )));
371 }
372 }
373 } else {
374 let rid = self.next_row_id;
375 self.next_row_id = self.next_row_id.saturating_add(1);
376 rid
377 }
378 }
379 RowIdMode::None => {
380 let rid = self.next_row_id;
381 self.next_row_id = self.next_row_id.saturating_add(1);
382 rid
383 }
384 };
385
386 Ok(Some(Row::new(row_id, values)))
387 }
388}
389
390impl RowIterator for ColumnarScanIterator {
391 fn next_row(&mut self) -> Option<Result<Row>> {
392 self.advance()
393 }
394
395 fn schema(&self) -> &[ColumnMetadata] {
396 &self.schema
397 }
398}
399
400pub fn create_columnar_scan_iterator<'txn, S: KVStore + 'txn>(
415 txn: &mut impl SqlTxn<'txn, S>,
416 table_meta: &TableMetadata,
417 scan: &ColumnarScan,
418) -> Result<ColumnarScanIterator> {
419 debug_assert_eq!(scan.table_id, table_meta.table_id);
420
421 let projected: Vec<usize> = if scan.projected_columns.is_empty() {
422 (0..table_meta.columns.len()).collect()
423 } else {
424 scan.projected_columns.clone()
425 };
426
427 let segment_ids = load_segment_index(txn, table_meta.table_id)?;
428
429 let row_id_col_idx = if table_meta.storage_options.row_id_mode == RowIdMode::Direct {
430 table_meta
431 .columns
432 .iter()
433 .position(|c| c.name.eq_ignore_ascii_case("row_id"))
434 } else {
435 None
436 };
437
438 let mut segments = Vec::with_capacity(segment_ids.len());
440 for segment_id in segment_ids {
441 let segment = load_segment(txn, table_meta.table_id, segment_id)?;
442 let reader =
443 SegmentReaderV2::open(Box::new(InMemorySegmentSource::new(segment.data.clone())))
444 .map_err(|e| ExecutorError::Columnar(e.to_string()))?;
445 let row_group_stats = load_row_group_stats(txn, table_meta.table_id, segment_id);
446
447 segments.push(LoadedSegment {
448 reader,
449 row_group_stats,
450 row_ids: segment.row_ids.clone(),
451 row_groups: segment.meta.row_groups.clone(),
452 });
453 }
454
455 Ok(ColumnarScanIterator {
456 segments,
457 segment_idx: 0,
458 row_group_idx: 0,
459 row_idx: 0,
460 current_batch: None,
461 projected,
462 schema: table_meta.columns.clone(),
463 table_meta: table_meta.clone(),
464 scan: scan.clone(),
465 row_id_col_idx,
466 next_row_id: 0,
467 })
468}
469
470pub fn execute_columnar_scan<'txn, S: KVStore + 'txn>(
472 txn: &mut impl SqlTxn<'txn, S>,
473 table_meta: &TableMetadata,
474 scan: &ColumnarScan,
475) -> Result<Vec<Row>> {
476 debug_assert_eq!(scan.table_id, table_meta.table_id);
477 let projected: Vec<usize> = if scan.projected_columns.is_empty() {
478 (0..table_meta.columns.len()).collect()
479 } else {
480 scan.projected_columns.clone()
481 };
482
483 let segment_ids = load_segment_index(txn, table_meta.table_id)?;
484 if segment_ids.is_empty() {
485 return Ok(Vec::new());
486 }
487
488 let row_id_col_idx = if table_meta.storage_options.row_id_mode == RowIdMode::Direct {
489 table_meta
490 .columns
491 .iter()
492 .position(|c| c.name.eq_ignore_ascii_case("row_id"))
493 } else {
494 None
495 };
496
497 let mut results = Vec::new();
498 let mut next_row_id = 0u64;
499 for segment_id in segment_ids {
500 let segment = load_segment(txn, table_meta.table_id, segment_id)?;
501 let reader =
502 SegmentReaderV2::open(Box::new(InMemorySegmentSource::new(segment.data.clone())))
503 .map_err(|e| ExecutorError::Columnar(e.to_string()))?;
504
505 let row_group_stats = load_row_group_stats(txn, table_meta.table_id, segment_id);
506 let row_group_count = segment.meta.row_groups.len();
507 for rg_index in 0..row_group_count {
508 let should_skip = match row_group_stats.as_ref() {
509 Some(stats) if stats.len() == row_group_count => {
510 scan.should_skip_row_group(&stats[rg_index])
511 }
512 _ => false,
513 };
514 if should_skip {
515 continue;
516 }
517
518 let batch = reader
519 .read_row_group_by_index(&projected, rg_index)
520 .map_err(|e| ExecutorError::Columnar(e.to_string()))?;
521 let batch = if !segment.row_ids.is_empty() {
522 if let Some(meta) = segment.meta.row_groups.get(rg_index) {
523 let start = meta.row_start as usize;
524 let end = start + meta.row_count as usize;
525 if end <= segment.row_ids.len() {
526 batch.with_row_ids(Some(segment.row_ids[start..end].to_vec()))
527 } else {
528 batch
529 }
530 } else {
531 batch
532 }
533 } else {
534 batch
535 };
536 append_rows_from_batch(
537 &mut results,
538 &batch,
539 table_meta,
540 &projected,
541 scan.residual_filter.as_ref(),
542 table_meta.storage_options.row_id_mode,
543 row_id_col_idx,
544 &mut next_row_id,
545 )?;
546 }
547 }
548
549 Ok(results)
550}
551
552pub fn execute_columnar_row_ids<'txn, S: KVStore + 'txn>(
557 txn: &mut impl SqlTxn<'txn, S>,
558 table_meta: &TableMetadata,
559 scan: &ColumnarScan,
560) -> Result<Vec<u64>> {
561 if table_meta.storage_options.storage_type != crate::catalog::StorageType::Columnar {
562 return Err(ExecutorError::Columnar(
563 "execute_columnar_row_ids requires columnar storage".into(),
564 ));
565 }
566
567 let mut needed: BTreeSet<usize> = scan.projected_columns.iter().copied().collect();
568 if let Some(pred) = &scan.residual_filter {
569 collect_column_indices(pred, &mut needed);
570 }
571 let projected: Vec<usize> = needed.into_iter().collect();
572
573 let segment_ids = load_segment_index(txn, table_meta.table_id)?;
574 if segment_ids.is_empty() {
575 return Ok(Vec::new());
576 }
577
578 let row_id_col_idx = if table_meta.storage_options.row_id_mode == RowIdMode::Direct {
579 table_meta
580 .columns
581 .iter()
582 .position(|c| c.name.eq_ignore_ascii_case("row_id"))
583 } else {
584 None
585 };
586
587 let mut results = Vec::new();
588 let mut next_row_id = 0u64;
589 for segment_id in segment_ids {
590 let segment = load_segment(txn, table_meta.table_id, segment_id)?;
591 let reader =
592 SegmentReaderV2::open(Box::new(InMemorySegmentSource::new(segment.data.clone())))
593 .map_err(|e| ExecutorError::Columnar(e.to_string()))?;
594
595 let row_group_stats = load_row_group_stats(txn, table_meta.table_id, segment_id);
596 let row_group_count = segment.meta.row_groups.len();
597 for rg_index in 0..row_group_count {
598 let should_skip = match row_group_stats.as_ref() {
599 Some(stats) if stats.len() == row_group_count => {
600 scan.should_skip_row_group(&stats[rg_index])
601 }
602 _ => false,
603 };
604 if should_skip {
605 continue;
606 }
607
608 let batch = reader
609 .read_row_group_by_index(&projected, rg_index)
610 .map_err(|e| ExecutorError::Columnar(e.to_string()))?;
611 let batch = if !segment.row_ids.is_empty() {
612 if let Some(meta) = segment.meta.row_groups.get(rg_index) {
613 let start = meta.row_start as usize;
614 let end = start + meta.row_count as usize;
615 if end <= segment.row_ids.len() {
616 batch.with_row_ids(Some(segment.row_ids[start..end].to_vec()))
617 } else {
618 batch
619 }
620 } else {
621 batch
622 }
623 } else {
624 batch
625 };
626
627 let row_count = batch.num_rows();
628 for row_idx in 0..row_count {
629 let mut values = vec![SqlValue::Null; table_meta.column_count()];
631 for (pos, &table_col_idx) in projected.iter().enumerate() {
632 let column = batch.columns.get(pos).ok_or_else(|| {
633 ExecutorError::Columnar("missing projected column".into())
634 })?;
635 let bitmap = batch.null_bitmaps.get(pos).and_then(|b| b.as_ref());
636 let value = value_from_column(
637 column,
638 bitmap,
639 row_idx,
640 &table_meta
641 .columns
642 .get(table_col_idx)
643 .ok_or_else(|| {
644 ExecutorError::Columnar("column index out of bounds".into())
645 })?
646 .data_type,
647 )?;
648 values[table_col_idx] = value;
649 }
650
651 if let Some(predicate) = scan.residual_filter.as_ref() {
652 let ctx = EvalContext::new(&values);
653 let keep = matches!(evaluate(predicate, &ctx)?, SqlValue::Boolean(true));
654 if !keep {
655 continue;
656 }
657 }
658
659 let row_id = match table_meta.storage_options.row_id_mode {
660 RowIdMode::Direct => {
661 if let Some(row_ids) = batch.row_ids.as_ref() {
662 *row_ids.get(row_idx).ok_or_else(|| {
663 ExecutorError::Columnar(
664 "row_id missing for row in row_id_mode=direct".into(),
665 )
666 })?
667 } else if let Some(idx) = row_id_col_idx {
668 let val = values.get(idx).ok_or_else(|| {
669 ExecutorError::Columnar(
670 "row_id column missing in projected values".into(),
671 )
672 })?;
673 match val {
674 SqlValue::Integer(v) if *v >= 0 => *v as u64,
675 SqlValue::BigInt(v) if *v >= 0 => *v as u64,
676 other => {
677 return Err(ExecutorError::Columnar(format!(
678 "row_id column must be non-negative integer, got {}",
679 other.type_name()
680 )));
681 }
682 }
683 } else {
684 let rid = next_row_id;
685 next_row_id = next_row_id.saturating_add(1);
686 rid
687 }
688 }
689 RowIdMode::None => {
690 let rid = next_row_id;
691 next_row_id = next_row_id.saturating_add(1);
692 rid
693 }
694 };
695 results.push(row_id);
696 }
697 }
698 }
699
700 Ok(results)
701}
702
703pub fn expr_to_pushdown(expr: &TypedExpr) -> Option<PushdownFilter> {
705 match &expr.kind {
706 TypedExprKind::BinaryOp { left, op, right } => match op {
707 BinaryOp::And => {
708 let l = expr_to_pushdown(left)?;
709 let r = expr_to_pushdown(right)?;
710 Some(PushdownFilter::And(vec![l, r]))
711 }
712 BinaryOp::Or => {
713 let l = expr_to_pushdown(left)?;
714 let r = expr_to_pushdown(right)?;
715 Some(PushdownFilter::Or(vec![l, r]))
716 }
717 BinaryOp::Eq => extract_eq(left, right),
718 BinaryOp::Lt | BinaryOp::LtEq | BinaryOp::Gt | BinaryOp::GtEq => {
719 extract_range(op, left, right)
720 }
721 _ => None,
722 },
723 TypedExprKind::Between {
724 expr,
725 low,
726 high,
727 negated,
728 } => {
729 if *negated {
730 return None;
731 }
732 let (column_idx, value_min, value_max) = match expr.kind {
733 TypedExprKind::ColumnRef { column_index, .. } => {
734 let low_v = literal_value(low)?;
735 let high_v = literal_value(high)?;
736 (column_index, low_v, high_v)
737 }
738 _ => return None,
739 };
740 Some(PushdownFilter::Range {
741 column_idx,
742 min: Some(value_min),
743 max: Some(value_max),
744 })
745 }
746 TypedExprKind::IsNull { expr, negated } => match expr.kind {
747 TypedExprKind::ColumnRef { column_index, .. } => Some(PushdownFilter::IsNull {
748 column_idx: column_index,
749 is_null: !negated,
750 }),
751 _ => None,
752 },
753 _ => None,
754 }
755}
756
757fn extract_eq(left: &TypedExpr, right: &TypedExpr) -> Option<PushdownFilter> {
758 if let Some((col_idx, value)) = extract_column_literal(left, right) {
759 return Some(PushdownFilter::Eq {
760 column_idx: col_idx,
761 value,
762 });
763 }
764 if let Some((col_idx, value)) = extract_column_literal(right, left) {
765 return Some(PushdownFilter::Eq {
766 column_idx: col_idx,
767 value,
768 });
769 }
770 None
771}
772
773fn extract_range(op: &BinaryOp, left: &TypedExpr, right: &TypedExpr) -> Option<PushdownFilter> {
774 match (
775 extract_column_literal(left, right),
776 extract_column_literal(right, left),
777 ) {
778 (Some((col_idx, value)), _) => match op {
779 BinaryOp::Lt | BinaryOp::LtEq => Some(PushdownFilter::Range {
780 column_idx: col_idx,
781 min: None,
782 max: Some(value),
783 }),
784 BinaryOp::Gt | BinaryOp::GtEq => Some(PushdownFilter::Range {
785 column_idx: col_idx,
786 min: Some(value),
787 max: None,
788 }),
789 _ => None,
790 },
791 (_, Some((col_idx, value))) => match op {
792 BinaryOp::Lt | BinaryOp::LtEq => Some(PushdownFilter::Range {
793 column_idx: col_idx,
794 min: Some(value),
795 max: None,
796 }),
797 BinaryOp::Gt | BinaryOp::GtEq => Some(PushdownFilter::Range {
798 column_idx: col_idx,
799 min: None,
800 max: Some(value),
801 }),
802 _ => None,
803 },
804 _ => None,
805 }
806}
807
808fn extract_column_literal(
809 column_expr: &TypedExpr,
810 literal_expr: &TypedExpr,
811) -> Option<(usize, SqlValue)> {
812 match column_expr.kind {
813 TypedExprKind::ColumnRef { column_index, .. } => {
814 let value = literal_value(literal_expr)?;
815 Some((column_index, value))
816 }
817 _ => None,
818 }
819}
820
821fn literal_value(expr: &TypedExpr) -> Option<SqlValue> {
822 match &expr.kind {
823 TypedExprKind::Literal(_) | TypedExprKind::VectorLiteral(_) => {
824 evaluate(expr, &EvalContext::new(&[])).ok()
825 }
826 _ => None,
827 }
828}
829
830pub fn projection_to_columns(projection: &Projection, table_meta: &TableMetadata) -> Vec<usize> {
832 match projection {
833 Projection::All(names) => names
834 .iter()
835 .filter_map(|name| table_meta.columns.iter().position(|c| &c.name == name))
836 .collect(),
837 Projection::Columns(cols) => {
838 let mut indices = BTreeSet::new();
839 for col in cols {
840 collect_column_indices(&col.expr, &mut indices);
841 }
842 if indices.is_empty() {
843 return (0..table_meta.columns.len()).collect();
844 }
845 indices
846 .into_iter()
847 .filter(|idx| *idx < table_meta.columns.len())
848 .collect()
849 }
850 }
851}
852
853pub fn build_columnar_scan_for_filter(
855 table_meta: &TableMetadata,
856 projection: Projection,
857 predicate: &TypedExpr,
858) -> ColumnarScan {
859 let mut projected_columns = projection_to_columns(&projection, table_meta);
860 let mut predicate_indices = BTreeSet::new();
861 collect_column_indices(predicate, &mut predicate_indices);
862 for idx in predicate_indices {
863 if !projected_columns.contains(&idx) {
864 projected_columns.push(idx);
865 }
866 }
867 projected_columns.sort_unstable();
868 let pushed_filter = expr_to_pushdown(predicate);
869 ColumnarScan::new(
870 table_meta.table_id,
871 projected_columns,
872 pushed_filter,
873 Some(predicate.clone()),
874 )
875}
876
877pub fn build_columnar_scan(table_meta: &TableMetadata, projection: &Projection) -> ColumnarScan {
879 let projected_columns = projection_to_columns(projection, table_meta);
880 ColumnarScan::new(table_meta.table_id, projected_columns, None, None)
881}
882
883fn collect_column_indices(expr: &TypedExpr, acc: &mut BTreeSet<usize>) {
885 match &expr.kind {
886 TypedExprKind::ColumnRef { column_index, .. } => {
887 acc.insert(*column_index);
888 }
889 TypedExprKind::BinaryOp { left, right, .. } => {
890 collect_column_indices(left, acc);
891 collect_column_indices(right, acc);
892 }
893 TypedExprKind::UnaryOp { operand, .. } => collect_column_indices(operand, acc),
894 TypedExprKind::Cast { expr, .. } => collect_column_indices(expr, acc),
895 TypedExprKind::Case {
896 operand,
897 branches,
898 else_expr,
899 } => {
900 if let Some(operand) = operand {
901 collect_column_indices(operand, acc);
902 }
903 for branch in branches {
904 collect_column_indices(&branch.when, acc);
905 collect_column_indices(&branch.then, acc);
906 }
907 if let Some(else_expr) = else_expr {
908 collect_column_indices(else_expr, acc);
909 }
910 }
911 TypedExprKind::Between {
912 expr, low, high, ..
913 } => {
914 collect_column_indices(expr, acc);
915 collect_column_indices(low, acc);
916 collect_column_indices(high, acc);
917 }
918 TypedExprKind::Like {
919 expr,
920 pattern,
921 escape,
922 ..
923 } => {
924 collect_column_indices(expr, acc);
925 collect_column_indices(pattern, acc);
926 if let Some(escape) = escape {
927 collect_column_indices(escape, acc);
928 }
929 }
930 TypedExprKind::InList { expr, list, .. } => {
931 collect_column_indices(expr, acc);
932 for item in list {
933 collect_column_indices(item, acc);
934 }
935 }
936 TypedExprKind::IsNull { expr, .. } => collect_column_indices(expr, acc),
937 TypedExprKind::FunctionCall { args, .. } => {
938 for arg in args {
939 collect_column_indices(arg, acc);
940 }
941 }
942 _ => {}
943 }
944}
945
946fn load_segment_index<'txn, S: KVStore + 'txn>(
947 txn: &mut impl SqlTxn<'txn, S>,
948 table_id: u32,
949) -> Result<Vec<u64>> {
950 let key = key_layout::segment_index_key(table_id);
951 let bytes = txn.inner_mut().get(&key)?;
952 if let Some(raw) = bytes {
953 bincode_config()
954 .deserialize(&raw)
955 .map_err(|e| ExecutorError::Columnar(e.to_string()))
956 } else {
957 Ok(Vec::new())
958 }
959}
960
961fn load_segment<'txn, S: KVStore + 'txn>(
962 txn: &mut impl SqlTxn<'txn, S>,
963 table_id: u32,
964 segment_id: u64,
965) -> Result<ColumnSegmentV2> {
966 let key = key_layout::column_segment_key(table_id, segment_id, 0);
967 let bytes = txn
968 .inner_mut()
969 .get(&key)?
970 .ok_or_else(|| ExecutorError::Columnar(format!("segment {segment_id} missing")))?;
971 bincode_config()
972 .deserialize(&bytes)
973 .map_err(|e| ExecutorError::Columnar(e.to_string()))
974}
975
976fn load_row_group_stats<'txn, S: KVStore + 'txn>(
977 txn: &mut impl SqlTxn<'txn, S>,
978 table_id: u32,
979 segment_id: u64,
980) -> Option<Vec<RowGroupStatistics>> {
981 let key = key_layout::row_group_stats_key(table_id, segment_id);
982 match txn.inner_mut().get(&key) {
983 Ok(Some(bytes)) => bincode_config().deserialize(&bytes).ok(),
984 Ok(None) => None,
985 Err(_) => None,
986 }
987}
988
989#[allow(clippy::too_many_arguments)]
990fn append_rows_from_batch(
991 out: &mut Vec<Row>,
992 batch: &alopex_core::columnar::segment_v2::RecordBatch,
993 table_meta: &TableMetadata,
994 projected: &[usize],
995 residual_filter: Option<&TypedExpr>,
996 row_id_mode: RowIdMode,
997 row_id_col_idx: Option<usize>,
998 next_row_id: &mut u64,
999) -> Result<()> {
1000 if batch.columns.len() != projected.len() {
1001 return Err(ExecutorError::Columnar(format!(
1002 "projected column count mismatch: requested {}, got {}",
1003 projected.len(),
1004 batch.columns.len()
1005 )));
1006 }
1007
1008 let row_count = batch.num_rows();
1009 for row_idx in 0..row_count {
1010 let mut values = vec![SqlValue::Null; table_meta.column_count()];
1011 for (pos, &table_col_idx) in projected.iter().enumerate() {
1012 let column = batch
1013 .columns
1014 .get(pos)
1015 .ok_or_else(|| ExecutorError::Columnar("missing projected column".into()))?;
1016 let bitmap = batch.null_bitmaps.get(pos).and_then(|b| b.as_ref());
1017 let value = value_from_column(
1018 column,
1019 bitmap,
1020 row_idx,
1021 &table_meta
1022 .columns
1023 .get(table_col_idx)
1024 .ok_or_else(|| ExecutorError::Columnar("column index out of bounds".into()))?
1025 .data_type,
1026 )?;
1027 values[table_col_idx] = value;
1028 }
1029
1030 if let Some(predicate) = residual_filter {
1031 let ctx = EvalContext::new(&values);
1032 let keep = matches!(evaluate(predicate, &ctx)?, SqlValue::Boolean(true));
1033 if !keep {
1034 continue;
1035 }
1036 }
1037
1038 let row_id = match row_id_mode {
1039 RowIdMode::Direct => {
1040 if let Some(row_ids) = batch.row_ids.as_ref() {
1041 *row_ids.get(row_idx).ok_or_else(|| {
1042 ExecutorError::Columnar(
1043 "row_id missing for row in row_id_mode=direct".into(),
1044 )
1045 })?
1046 } else if let Some(idx) = row_id_col_idx {
1047 let val = values.get(idx).ok_or_else(|| {
1048 ExecutorError::Columnar("row_id column missing in projected values".into())
1049 })?;
1050 match val {
1051 SqlValue::Integer(v) if *v >= 0 => *v as u64,
1052 SqlValue::BigInt(v) if *v >= 0 => *v as u64,
1053 other => {
1054 return Err(ExecutorError::Columnar(format!(
1055 "row_id column must be non-negative integer, got {}",
1056 other.type_name()
1057 )));
1058 }
1059 }
1060 } else {
1061 let rid = *next_row_id;
1062 *next_row_id = next_row_id.saturating_add(1);
1063 rid
1064 }
1065 }
1066 RowIdMode::None => {
1067 let rid = *next_row_id;
1068 *next_row_id = next_row_id.saturating_add(1);
1069 rid
1070 }
1071 };
1072 out.push(Row::new(row_id, values));
1073 }
1074
1075 Ok(())
1076}
1077
1078fn value_from_column(
1079 column: &Column,
1080 bitmap: Option<&Bitmap>,
1081 row_idx: usize,
1082 ty: &ResolvedType,
1083) -> Result<SqlValue> {
1084 if let Some(bm) = bitmap
1085 && !bm.get(row_idx)
1086 {
1087 return Ok(SqlValue::Null);
1088 }
1089
1090 match (ty, column) {
1091 (ResolvedType::Integer, Column::Int64(values)) => {
1092 let v = *values
1093 .get(row_idx)
1094 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1095 Ok(SqlValue::Integer(v as i32))
1096 }
1097 (ResolvedType::BigInt | ResolvedType::Timestamp, Column::Int64(values)) => {
1098 let v = *values
1099 .get(row_idx)
1100 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1101 if matches!(ty, ResolvedType::Timestamp) {
1102 Ok(SqlValue::Timestamp(v))
1103 } else {
1104 Ok(SqlValue::BigInt(v))
1105 }
1106 }
1107 (ResolvedType::Float, Column::Float32(values)) => {
1108 let v = *values
1109 .get(row_idx)
1110 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1111 Ok(SqlValue::Float(v))
1112 }
1113 (ResolvedType::Double, Column::Float64(values)) => {
1114 let v = *values
1115 .get(row_idx)
1116 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1117 Ok(SqlValue::Double(v))
1118 }
1119 (ResolvedType::Boolean, Column::Bool(values)) => {
1120 let v = *values
1121 .get(row_idx)
1122 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1123 Ok(SqlValue::Boolean(v))
1124 }
1125 (ResolvedType::Text, Column::Binary(values)) => {
1126 let raw = values
1127 .get(row_idx)
1128 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1129 String::from_utf8(raw.clone())
1130 .map(SqlValue::Text)
1131 .map_err(|e| ExecutorError::Columnar(e.to_string()))
1132 }
1133 (ResolvedType::Blob, Column::Binary(values)) => {
1134 let raw = values
1135 .get(row_idx)
1136 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1137 Ok(SqlValue::Blob(raw.clone()))
1138 }
1139 (ResolvedType::Vector { .. }, Column::Fixed { values, .. }) => {
1140 let raw = values
1141 .get(row_idx)
1142 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1143 if raw.len() % 4 != 0 {
1144 return Err(ExecutorError::Columnar(
1145 "invalid vector byte length in columnar segment".into(),
1146 ));
1147 }
1148 let floats: Vec<f32> = raw
1149 .chunks_exact(4)
1150 .map(|bytes| f32::from_le_bytes(bytes.try_into().unwrap()))
1151 .collect();
1152 Ok(SqlValue::Vector(floats))
1153 }
1154 (_, Column::Binary(values)) => {
1155 let raw = values
1156 .get(row_idx)
1157 .ok_or_else(|| ExecutorError::Columnar("row index out of bounds".into()))?;
1158 Ok(SqlValue::Blob(raw.clone()))
1159 }
1160 _ => Err(ExecutorError::Columnar(
1161 "unsupported column type for columnar read".into(),
1162 )),
1163 }
1164}
1165#[cfg(test)]
1166mod tests {
1167 use super::*;
1168 use crate::ast::expr::Literal;
1169 use crate::ast::span::Span;
1170 use crate::catalog::{ColumnMetadata, RowIdMode, TableMetadata};
1171 use crate::columnar::statistics::ColumnStatistics;
1172 use crate::planner::TypedCaseWhen;
1173 use crate::planner::typed_expr::TypedExpr;
1174 use crate::planner::typed_expr::TypedExprKind;
1175 use crate::planner::types::ResolvedType;
1176 use crate::storage::TxnBridge;
1177 use alopex_core::kv::memory::MemoryKV;
1178 use bincode::config::Options;
1179 use std::sync::Arc;
1180
1181 #[test]
1182 fn case_promotion_cast_keeps_column_in_projection() {
1183 let span = Span::default();
1184 let column = TypedExpr {
1185 kind: TypedExprKind::ColumnRef {
1186 table: "items".to_string(),
1187 column: "value".to_string(),
1188 column_index: 3,
1189 },
1190 resolved_type: ResolvedType::Integer,
1191 span,
1192 };
1193 let promoted_column = TypedExpr {
1194 kind: TypedExprKind::Cast {
1195 expr: Box::new(column),
1196 target_type: ResolvedType::Double,
1197 },
1198 resolved_type: ResolvedType::Double,
1199 span,
1200 };
1201 let case = TypedExpr {
1202 kind: TypedExprKind::Case {
1203 operand: None,
1204 branches: vec![TypedCaseWhen {
1205 when: TypedExpr {
1206 kind: TypedExprKind::Literal(Literal::Boolean(true)),
1207 resolved_type: ResolvedType::Boolean,
1208 span,
1209 },
1210 then: promoted_column,
1211 }],
1212 else_expr: None,
1213 },
1214 resolved_type: ResolvedType::Double,
1215 span,
1216 };
1217
1218 let mut columns = BTreeSet::new();
1219 collect_column_indices(&case, &mut columns);
1220 assert_eq!(columns, BTreeSet::from([3]));
1221 }
1222
1223 #[test]
1224 fn evaluate_pushdown_eq_prunes_out_of_range() {
1225 let stats = RowGroupStatistics {
1226 row_count: 3,
1227 columns: vec![ColumnStatistics {
1228 min: SqlValue::Integer(1),
1229 max: SqlValue::Integer(3),
1230 null_count: 0,
1231 total_count: 3,
1232 distinct_count: None,
1233 }],
1234 row_id_min: None,
1235 row_id_max: None,
1236 };
1237 let filter = PushdownFilter::Eq {
1238 column_idx: 0,
1239 value: SqlValue::Integer(10),
1240 };
1241 assert!(ColumnarScan::evaluate_pushdown(&filter, &stats));
1242 }
1243
1244 #[test]
1245 fn evaluate_pushdown_range_allows_overlap() {
1246 let stats = RowGroupStatistics {
1247 row_count: 3,
1248 columns: vec![ColumnStatistics {
1249 min: SqlValue::Integer(5),
1250 max: SqlValue::Integer(10),
1251 null_count: 0,
1252 total_count: 3,
1253 distinct_count: None,
1254 }],
1255 row_id_min: None,
1256 row_id_max: None,
1257 };
1258 let filter = PushdownFilter::Range {
1259 column_idx: 0,
1260 min: Some(SqlValue::Integer(8)),
1261 max: Some(SqlValue::Integer(12)),
1262 };
1263 assert!(!ColumnarScan::evaluate_pushdown(&filter, &stats));
1264 }
1265
1266 #[test]
1267 fn evaluate_pushdown_is_null_skips_when_no_nulls() {
1268 let stats = RowGroupStatistics {
1269 row_count: 2,
1270 columns: vec![ColumnStatistics {
1271 min: SqlValue::Integer(1),
1272 max: SqlValue::Integer(2),
1273 null_count: 0,
1274 total_count: 2,
1275 distinct_count: None,
1276 }],
1277 row_id_min: None,
1278 row_id_max: None,
1279 };
1280 let filter = PushdownFilter::IsNull {
1281 column_idx: 0,
1282 is_null: true,
1283 };
1284 assert!(ColumnarScan::evaluate_pushdown(&filter, &stats));
1285 }
1286
1287 #[test]
1288 fn evaluate_pushdown_is_not_null_skips_when_all_null() {
1289 let stats = RowGroupStatistics {
1290 row_count: 2,
1291 columns: vec![ColumnStatistics {
1292 min: SqlValue::Null,
1293 max: SqlValue::Null,
1294 null_count: 2,
1295 total_count: 2,
1296 distinct_count: None,
1297 }],
1298 row_id_min: None,
1299 row_id_max: None,
1300 };
1301 let filter = PushdownFilter::IsNull {
1302 column_idx: 0,
1303 is_null: false,
1304 };
1305 assert!(ColumnarScan::evaluate_pushdown(&filter, &stats));
1306 }
1307
1308 #[test]
1309 fn evaluate_pushdown_and_prunes_if_any_branch_skips() {
1310 let stats = RowGroupStatistics {
1311 row_count: 3,
1312 columns: vec![ColumnStatistics {
1313 min: SqlValue::Integer(1),
1314 max: SqlValue::Integer(3),
1315 null_count: 0,
1316 total_count: 3,
1317 distinct_count: None,
1318 }],
1319 row_id_min: None,
1320 row_id_max: None,
1321 };
1322 let filter = PushdownFilter::And(vec![
1323 PushdownFilter::Eq {
1324 column_idx: 0,
1325 value: SqlValue::Integer(10),
1326 },
1327 PushdownFilter::Eq {
1328 column_idx: 0,
1329 value: SqlValue::Integer(2),
1330 },
1331 ]);
1332 assert!(ColumnarScan::evaluate_pushdown(&filter, &stats));
1333 }
1334
1335 #[test]
1336 fn evaluate_pushdown_or_keeps_if_any_branch_may_match() {
1337 let stats = RowGroupStatistics {
1338 row_count: 3,
1339 columns: vec![ColumnStatistics {
1340 min: SqlValue::Integer(1),
1341 max: SqlValue::Integer(3),
1342 null_count: 0,
1343 total_count: 3,
1344 distinct_count: None,
1345 }],
1346 row_id_min: None,
1347 row_id_max: None,
1348 };
1349 let filter = PushdownFilter::Or(vec![
1350 PushdownFilter::Eq {
1351 column_idx: 0,
1352 value: SqlValue::Integer(10),
1353 },
1354 PushdownFilter::Eq {
1355 column_idx: 0,
1356 value: SqlValue::Integer(2),
1357 },
1358 ]);
1359 assert!(!ColumnarScan::evaluate_pushdown(&filter, &stats));
1360 }
1361
1362 #[test]
1363 fn expr_to_pushdown_converts_eq() {
1364 let expr = TypedExpr {
1365 kind: TypedExprKind::BinaryOp {
1366 left: Box::new(TypedExpr::column_ref(
1367 "t".into(),
1368 "c".into(),
1369 0,
1370 ResolvedType::Integer,
1371 crate::Span::default(),
1372 )),
1373 op: BinaryOp::Eq,
1374 right: Box::new(TypedExpr::literal(
1375 Literal::Number("1".into()),
1376 ResolvedType::Integer,
1377 crate::Span::default(),
1378 )),
1379 },
1380 resolved_type: ResolvedType::Boolean,
1381 span: crate::Span::default(),
1382 };
1383 let filter = expr_to_pushdown(&expr).unwrap();
1384 assert_eq!(
1385 filter,
1386 PushdownFilter::Eq {
1387 column_idx: 0,
1388 value: SqlValue::Integer(1)
1389 }
1390 );
1391 }
1392
1393 #[test]
1394 fn execute_columnar_scan_applies_residual_filter() {
1395 let bridge = TxnBridge::new(Arc::new(MemoryKV::new()));
1396 let mut table = TableMetadata::new(
1397 "users",
1398 vec![
1399 ColumnMetadata::new("id", ResolvedType::Integer),
1400 ColumnMetadata::new("name", ResolvedType::Text),
1401 ],
1402 )
1403 .with_table_id(1);
1404 table.storage_options.storage_type = crate::catalog::StorageType::Columnar;
1405
1406 let schema = alopex_core::columnar::segment_v2::Schema {
1408 columns: vec![
1409 alopex_core::columnar::segment_v2::ColumnSchema {
1410 name: "id".into(),
1411 logical_type: alopex_core::columnar::encoding::LogicalType::Int64,
1412 nullable: false,
1413 fixed_len: None,
1414 },
1415 alopex_core::columnar::segment_v2::ColumnSchema {
1416 name: "name".into(),
1417 logical_type: alopex_core::columnar::encoding::LogicalType::Binary,
1418 nullable: false,
1419 fixed_len: None,
1420 },
1421 ],
1422 };
1423 let batch = alopex_core::columnar::segment_v2::RecordBatch::new(
1424 schema.clone(),
1425 vec![
1426 alopex_core::columnar::encoding::Column::Int64(vec![1]),
1427 alopex_core::columnar::encoding::Column::Binary(vec![b"alice".to_vec()]),
1428 ],
1429 vec![None, None],
1430 );
1431 let mut writer =
1432 alopex_core::columnar::segment_v2::SegmentWriterV2::new(Default::default());
1433 writer.write_batch(batch).unwrap();
1434 let segment = writer.finish().unwrap();
1435
1436 let stats = vec![crate::columnar::statistics::compute_row_group_statistics(
1437 &[vec![SqlValue::Integer(1), SqlValue::Text("alice".into())]],
1438 )];
1439
1440 let mut txn = bridge.begin_write().unwrap();
1441 let segment_bytes = alopex_core::storage::format::bincode_config()
1442 .serialize(&segment)
1443 .unwrap();
1444 let meta_bytes = alopex_core::storage::format::bincode_config()
1445 .serialize(&segment.meta)
1446 .unwrap();
1447 let stats_bytes = alopex_core::storage::format::bincode_config()
1448 .serialize(&stats)
1449 .unwrap();
1450 txn.inner_mut()
1451 .put(
1452 alopex_core::columnar::kvs_bridge::key_layout::column_segment_key(1, 0, 0),
1453 segment_bytes,
1454 )
1455 .unwrap();
1456 txn.inner_mut()
1457 .put(
1458 alopex_core::columnar::kvs_bridge::key_layout::statistics_key(1, 0),
1459 meta_bytes,
1460 )
1461 .unwrap();
1462 txn.inner_mut()
1463 .put(
1464 alopex_core::columnar::kvs_bridge::key_layout::row_group_stats_key(1, 0),
1465 stats_bytes,
1466 )
1467 .unwrap();
1468 let index_bytes = alopex_core::storage::format::bincode_config()
1469 .serialize(&vec![0u64])
1470 .unwrap();
1471 txn.inner_mut()
1472 .put(
1473 alopex_core::columnar::kvs_bridge::key_layout::segment_index_key(1),
1474 index_bytes,
1475 )
1476 .unwrap();
1477 txn.commit().unwrap();
1478
1479 let scan = ColumnarScan::new(
1480 table.table_id,
1481 vec![0, 1],
1482 Some(PushdownFilter::Eq {
1483 column_idx: 0,
1484 value: SqlValue::Integer(1),
1485 }),
1486 Some(TypedExpr {
1487 kind: TypedExprKind::BinaryOp {
1488 left: Box::new(TypedExpr::column_ref(
1489 "users".into(),
1490 "id".into(),
1491 0,
1492 ResolvedType::Integer,
1493 crate::Span::default(),
1494 )),
1495 op: BinaryOp::Eq,
1496 right: Box::new(TypedExpr::literal(
1497 Literal::Number("1".into()),
1498 ResolvedType::Integer,
1499 crate::Span::default(),
1500 )),
1501 },
1502 resolved_type: ResolvedType::Boolean,
1503 span: crate::Span::default(),
1504 }),
1505 );
1506
1507 let mut read_txn = bridge.begin_read().unwrap();
1508 let rows = execute_columnar_scan(&mut read_txn, &table, &scan).unwrap();
1509 assert_eq!(rows.len(), 1);
1510 assert_eq!(rows[0].values[1], SqlValue::Text("alice".into()));
1511 }
1512
1513 #[test]
1514 fn rowid_mode_direct_prefers_rowid_column() {
1515 let bridge = TxnBridge::new(Arc::new(MemoryKV::new()));
1516 let mut table = TableMetadata::new(
1517 "items",
1518 vec![
1519 ColumnMetadata::new("row_id", ResolvedType::BigInt),
1520 ColumnMetadata::new("val", ResolvedType::Integer),
1521 ],
1522 )
1523 .with_table_id(20);
1524 table.storage_options.storage_type = crate::catalog::StorageType::Columnar;
1525 table.storage_options.row_id_mode = RowIdMode::Direct;
1526
1527 let schema = alopex_core::columnar::segment_v2::Schema {
1528 columns: vec![
1529 alopex_core::columnar::segment_v2::ColumnSchema {
1530 name: "row_id".into(),
1531 logical_type: alopex_core::columnar::encoding::LogicalType::Int64,
1532 nullable: false,
1533 fixed_len: None,
1534 },
1535 alopex_core::columnar::segment_v2::ColumnSchema {
1536 name: "val".into(),
1537 logical_type: alopex_core::columnar::encoding::LogicalType::Int64,
1538 nullable: false,
1539 fixed_len: None,
1540 },
1541 ],
1542 };
1543 let batch = alopex_core::columnar::segment_v2::RecordBatch::new(
1544 schema.clone(),
1545 vec![
1546 alopex_core::columnar::encoding::Column::Int64(vec![999]),
1547 alopex_core::columnar::encoding::Column::Int64(vec![7]),
1548 ],
1549 vec![None, None],
1550 );
1551 let mut writer =
1552 alopex_core::columnar::segment_v2::SegmentWriterV2::new(Default::default());
1553 writer.write_batch(batch).unwrap();
1554 let segment = writer.finish().unwrap();
1555 let stats = vec![crate::columnar::statistics::compute_row_group_statistics(
1556 &[vec![SqlValue::BigInt(999), SqlValue::Integer(7)]],
1557 )];
1558
1559 persist_segment_for_test(&bridge, table.table_id, &segment, &stats);
1560
1561 let scan = ColumnarScan::new(table.table_id, vec![0, 1], None, None);
1562 let mut read_txn = bridge.begin_read().unwrap();
1563 let rows = execute_columnar_scan(&mut read_txn, &table, &scan).unwrap();
1564 assert_eq!(rows.len(), 1);
1565 assert_eq!(rows[0].row_id, 999);
1566 assert_eq!(rows[0].values[1], SqlValue::Integer(7));
1567 }
1568
1569 #[test]
1570 fn rowid_mode_none_uses_position() {
1571 let bridge = TxnBridge::new(Arc::new(MemoryKV::new()));
1572 let mut table = TableMetadata::new(
1573 "items",
1574 vec![ColumnMetadata::new("val", ResolvedType::Integer)],
1575 )
1576 .with_table_id(21);
1577 table.storage_options.storage_type = crate::catalog::StorageType::Columnar;
1578 table.storage_options.row_id_mode = RowIdMode::Direct;
1579
1580 let schema = alopex_core::columnar::segment_v2::Schema {
1581 columns: vec![alopex_core::columnar::segment_v2::ColumnSchema {
1582 name: "val".into(),
1583 logical_type: alopex_core::columnar::encoding::LogicalType::Int64,
1584 nullable: false,
1585 fixed_len: None,
1586 }],
1587 };
1588 let batch = alopex_core::columnar::segment_v2::RecordBatch::new(
1589 schema.clone(),
1590 vec![alopex_core::columnar::encoding::Column::Int64(vec![3, 4])],
1591 vec![None],
1592 );
1593 let mut writer =
1594 alopex_core::columnar::segment_v2::SegmentWriterV2::new(Default::default());
1595 writer.write_batch(batch).unwrap();
1596 let segment = writer.finish().unwrap();
1597 let stats = vec![crate::columnar::statistics::compute_row_group_statistics(
1598 &[vec![SqlValue::Integer(3)], vec![SqlValue::Integer(4)]],
1599 )];
1600
1601 persist_segment_for_test(&bridge, table.table_id, &segment, &stats);
1602
1603 let scan = ColumnarScan::new(table.table_id, vec![0], None, None);
1604 let mut read_txn = bridge.begin_read().unwrap();
1605 let rows = execute_columnar_scan(&mut read_txn, &table, &scan).unwrap();
1606 assert_eq!(rows.len(), 2);
1607 assert_eq!(rows[0].row_id, 0);
1608 assert_eq!(rows[1].row_id, 1);
1609 }
1610
1611 fn persist_segment_for_test(
1612 bridge: &TxnBridge<MemoryKV>,
1613 table_id: u32,
1614 segment: &alopex_core::columnar::segment_v2::ColumnSegmentV2,
1615 row_group_stats: &[crate::columnar::statistics::RowGroupStatistics],
1616 ) {
1617 let mut txn = bridge.begin_write().unwrap();
1618 let segment_bytes = alopex_core::storage::format::bincode_config()
1619 .serialize(segment)
1620 .unwrap();
1621 let meta_bytes = alopex_core::storage::format::bincode_config()
1622 .serialize(&segment.meta)
1623 .unwrap();
1624 let stats_bytes = alopex_core::storage::format::bincode_config()
1625 .serialize(row_group_stats)
1626 .unwrap();
1627 txn.inner_mut()
1628 .put(
1629 alopex_core::columnar::kvs_bridge::key_layout::column_segment_key(table_id, 0, 0),
1630 segment_bytes,
1631 )
1632 .unwrap();
1633 txn.inner_mut()
1634 .put(
1635 alopex_core::columnar::kvs_bridge::key_layout::statistics_key(table_id, 0),
1636 meta_bytes,
1637 )
1638 .unwrap();
1639 txn.inner_mut()
1640 .put(
1641 alopex_core::columnar::kvs_bridge::key_layout::row_group_stats_key(table_id, 0),
1642 stats_bytes,
1643 )
1644 .unwrap();
1645 let index_bytes = alopex_core::storage::format::bincode_config()
1646 .serialize(&vec![0u64])
1647 .unwrap();
1648 txn.inner_mut()
1649 .put(
1650 alopex_core::columnar::kvs_bridge::key_layout::segment_index_key(table_id),
1651 index_bytes,
1652 )
1653 .unwrap();
1654 txn.commit().unwrap();
1655 }
1656}