Skip to main content

akar_processor/physical/
batch_insert.rs

1//! PhysicalBatchInsert — dedicated batch insert operator.
2//!
3//! Wraps `NodeTable::insert_rows_batch()` / `RelTable::insert_rels_batch()`
4//! for use in query plans where multiple CREATE statements can be fused
5//! into a single efficient batch operation.
6
7use crate::physical::types::{OperatorResult, PhysicalOperatorExec};
8use akar_common::types::{PhysicalTypeID, Value};
9use akar_common::vector::{DataChunk, ValueVector};
10use akar_storage::table::TableCatalog;
11use akar_storage::wal::{WalSink, log_insert_record, log_rel_insert_record};
12use akar_transaction::UndoRecord;
13use std::sync::{Arc, Mutex};
14
15/// Physical operator for BATCH INSERT — inserts pre-collected rows/rels
16/// into a table using batch APIs for maximum throughput.
17pub struct PhysicalBatchInsert {
18    pub table_name: String,
19    pub table_id: u64,
20    /// Rows to insert: each row is a Vec<Value> matching column order.
21    pub rows: Vec<Vec<Value>>,
22    pub table_catalog: Arc<TableCatalog>,
23    /// Active transaction id (P52.18).
24    pub txn_id: Option<u64>,
25    /// Undo sink for rollback records (P52.18).
26    pub undo_sink: Option<Arc<Mutex<Vec<UndoRecord>>>>,
27    /// Typed WAL sink so batch rows/edges survive restarts via replay (P60.2).
28    pub wal_sink: Option<WalSink>,
29}
30
31impl PhysicalOperatorExec for PhysicalBatchInsert {
32    fn operator_type(&self) -> &str {
33        "batch_insert"
34    }
35
36    fn execute(&self, _input: Vec<DataChunk>) -> OperatorResult {
37        let num_rows = self.rows.len();
38        if num_rows == 0 {
39            let mut v = ValueVector::new(PhysicalTypeID::Int64, 1);
40            v.resize(1);
41            v.set_i64(0, 0);
42            let arr = akar_common::arrow_vector::ArrowVector::from_legacy(&v).array;
43            return Ok(vec![DataChunk::new(vec![arr], vec![PhysicalTypeID::Int64])]);
44        }
45
46        // Try node table first, then rel table
47        if let Some(mut table) = self.table_catalog.get_node_table_by_name_mut(&self.table_name) {
48            let start = table.num_rows;
49            let count = table
50                .insert_rows_batch_with_txn(&self.rows, self.txn_id)
51                .map_err(|e| format!("BatchInsert node error: {e}"))?;
52            if self.wal_sink.is_some() {
53                for row in self.rows.iter().take(count as usize) {
54                    log_insert_record(&self.wal_sink, self.table_id, row);
55                }
56            }
57            if let Some(sink) = self.undo_sink.as_ref()
58                && let Ok(mut u) = sink.lock()
59            {
60                for row in start..start + count {
61                    u.push(UndoRecord::insert(self.table_id, row));
62                }
63            }
64            tracing::info!(
65                "BATCH INSERT: inserted {count} rows into node table '{}'",
66                self.table_name
67            );
68            let mut v = ValueVector::new(PhysicalTypeID::Int64, 1);
69            v.resize(1);
70            v.set_i64(0, count as i64);
71            let arr = akar_common::arrow_vector::ArrowVector::from_legacy(&v).array;
72            return Ok(vec![DataChunk::new(vec![arr], vec![PhysicalTypeID::Int64])]);
73        }
74
75        if let Some(mut table) = self.table_catalog.get_rel_table_by_name_mut(&self.table_name) {
76            let rels: Vec<(u64, u64, Vec<Value>)> = self
77                .rows
78                .iter()
79                .map(|row| {
80                    let from = match &row[0] {
81                        Value::Int64(v) => *v as u64,
82                        _ => 0,
83                    };
84                    let to = match &row[1] {
85                        Value::Int64(v) => *v as u64,
86                        _ => 0,
87                    };
88                    let props = row[2..].to_vec();
89                    (from, to, props)
90                })
91                .collect();
92            let start = table.edges.len();
93            let count = table
94                .insert_rels_batch(&rels)
95                .map_err(|e| format!("BatchInsert rel error: {e}"))?;
96            for (from, to, props) in &rels {
97                log_rel_insert_record(&self.wal_sink, self.table_id, *from, *to, props);
98            }
99            if let Some(sink) = self.undo_sink.as_ref()
100                && let Ok(mut u) = sink.lock()
101            {
102                for idx in start..start + count as usize {
103                    u.push(UndoRecord::insert(self.table_id, idx as u64));
104                }
105            }
106            tracing::info!(
107                "BATCH INSERT: inserted {count} rels into rel table '{}'",
108                self.table_name
109            );
110            let mut v = ValueVector::new(PhysicalTypeID::Int64, 1);
111            v.resize(1);
112            v.set_i64(0, count as i64);
113            let arr = akar_common::arrow_vector::ArrowVector::from_legacy(&v).array;
114            return Ok(vec![DataChunk::new(vec![arr], vec![PhysicalTypeID::Int64])]);
115        }
116
117        Err(format!(
118            "Table '{}' not found in storage catalog for BatchInsert",
119            self.table_name
120        )
121        .into())
122    }
123}