genegraph_storage/traits/
lance.rs1use crate::{StorageError, StorageResult};
2use arrow::record_batch::RecordBatch;
3use arrow::record_batch::RecordBatchIterator;
4use log::{debug, info};
5
6use futures::StreamExt;
7use lance::Dataset;
8use lance::dataset::{WriteMode, WriteParams};
9
10pub trait LanceStorage {
11 async fn write_lance_batch_async(&self, uri: String, batch: RecordBatch) -> StorageResult<()> {
13 info!("Writing Lance dataset to {}", uri);
14
15 let schema = batch.schema();
16 let batches = vec![batch];
17 let reader = RecordBatchIterator::new(batches.into_iter().map(Ok), schema);
18
19 let params = WriteParams {
20 mode: WriteMode::Create,
21 ..WriteParams::default()
22 };
23
24 Dataset::write(reader, &uri, Some(params))
25 .await
26 .map_err(|e| StorageError::Lance(e.to_string()))?;
27
28 info!("Successfully wrote Lance dataset to {}", uri);
29 Ok(())
30 }
31
32 async fn read_lance_all_batches_async(&self, uri: String) -> StorageResult<RecordBatch> {
34 info!("Reading Lance dataset from {}", uri);
35
36 let dataset = Dataset::open(&uri)
37 .await
38 .map_err(|e| StorageError::Lance(e.to_string()))?;
39 let scanner = dataset.scan();
40 let mut stream = scanner
41 .try_into_stream()
42 .await
43 .map_err(|e| StorageError::Lance(e.to_string()))?;
44
45 let mut batches = Vec::new();
46 while let Some(batch_result) = stream.next().await {
47 let batch = batch_result.map_err(|e| StorageError::Lance(e.to_string()))?;
48 batches.push(batch);
49 }
50
51 if batches.is_empty() {
52 return Err(StorageError::Invalid("Empty Lance dataset".into()));
53 }
54
55 let schema = batches[0].schema();
56 let combined = arrow::compute::concat_batches(&schema, &batches)
57 .map_err(|e| StorageError::Lance(format!("Failed to concatenate batches: {}", e)))?;
58
59 debug!(
60 "Combined Lance batch for {:?} has {} rows",
61 uri,
62 combined.num_rows()
63 );
64 Ok(combined)
65 }
66}