akar_processor/physical/
batch_insert.rs1use 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
15pub struct PhysicalBatchInsert {
18 pub table_name: String,
19 pub table_id: u64,
20 pub rows: Vec<Vec<Value>>,
22 pub table_catalog: Arc<TableCatalog>,
23 pub txn_id: Option<u64>,
25 pub undo_sink: Option<Arc<Mutex<Vec<UndoRecord>>>>,
27 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 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}