1use std::{collections::HashSet, sync::Arc};
5
6use arrow_array::{RecordBatch, UInt32Array};
7use arrow_ipc::CompressionType;
8use arrow_ipc::reader::FileReader as ArrowFileReader;
9use arrow_ipc::writer::{FileWriter as ArrowFileWriter, IpcWriteOptions};
10use arrow_schema::{ArrowError, DataType, Field, Schema};
11use bytes::Buf;
12use lance_core::error::{CorruptFileSnafu, box_error};
13use lance_core::utils::deletion::DeletionVector;
14use lance_core::utils::tracing::{AUDIT_MODE_CREATE, AUDIT_TYPE_DELETION, TRACE_FILE_AUDIT};
15use lance_core::{Error, Result};
16use lance_io::object_store::ObjectStore;
17use object_store::path::Path;
18use rand::Rng;
19use roaring::bitmap::RoaringBitmap;
20use snafu::ResultExt;
21use tracing::{info, instrument};
22
23use crate::format::{DeletionFile, DeletionFileType};
24
25pub const DELETIONS_DIR: &str = "_deletions";
26
27fn deletion_arrow_schema() -> Arc<Schema> {
29 Arc::new(Schema::new(vec![Field::new(
30 "row_id",
31 DataType::UInt32,
32 false,
33 )]))
34}
35
36pub fn deletion_file_path(base: &Path, fragment_id: u64, deletion_file: &DeletionFile) -> Path {
38 let DeletionFile {
39 read_version,
40 id,
41 file_type,
42 ..
43 } = deletion_file;
44 let suffix = file_type.suffix();
45 base.clone()
46 .join(DELETIONS_DIR)
47 .join(format!("{fragment_id}-{read_version}-{id}.{suffix}"))
48}
49
50pub fn relative_deletion_file_path(fragment_id: u64, deletion_file: &DeletionFile) -> String {
51 let DeletionFile {
52 read_version,
53 id,
54 file_type,
55 ..
56 } = deletion_file;
57 let suffix = file_type.suffix();
58 format!("{DELETIONS_DIR}/{fragment_id}-{read_version}-{id}.{suffix}")
59}
60
61pub async fn write_deletion_file(
66 base: &Path,
67 fragment_id: u64,
68 read_version: u64,
69 removed_rows: &DeletionVector,
70 object_store: &ObjectStore,
71) -> Result<Option<DeletionFile>> {
72 let deletion_file = match removed_rows {
73 DeletionVector::NoDeletions => None,
74 DeletionVector::Set(set) => {
75 let id = rand::rng().random::<u64>();
76 let deletion_file = DeletionFile {
77 read_version,
78 id,
79 file_type: DeletionFileType::Array,
80 num_deleted_rows: Some(set.len()),
81 base_id: None,
82 };
83 let path = deletion_file_path(base, fragment_id, &deletion_file);
84
85 let array = UInt32Array::from_iter(set.iter().copied());
86 let array = Arc::new(array);
87
88 let schema = deletion_arrow_schema();
89 let batch = RecordBatch::try_new(schema.clone(), vec![array])?;
90
91 let mut out: Vec<u8> = Vec::new();
92 let write_options =
93 IpcWriteOptions::default().try_with_compression(Some(CompressionType::ZSTD))?;
94 {
95 let mut writer = ArrowFileWriter::try_new_with_options(
96 &mut out,
97 schema.as_ref(),
98 write_options,
99 )?;
100 writer.write(&batch)?;
101 writer.finish()?;
102 }
104
105 object_store.put(&path, &out).await?;
106
107 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_DELETION, path = path.to_string());
108
109 Some(deletion_file)
110 }
111 DeletionVector::Bitmap(bitmap) => {
112 let id = rand::rng().random::<u64>();
113 let deletion_file = DeletionFile {
114 read_version,
115 id,
116 file_type: DeletionFileType::Bitmap,
117 num_deleted_rows: Some(bitmap.len() as usize),
118 base_id: None,
119 };
120 let path = deletion_file_path(base, fragment_id, &deletion_file);
121
122 let mut out: Vec<u8> = Vec::new();
123 bitmap.serialize_into(&mut out)?;
124
125 object_store.put(&path, &out).await?;
126
127 info!(target: TRACE_FILE_AUDIT, mode=AUDIT_MODE_CREATE, r#type=AUDIT_TYPE_DELETION, path = path.to_string());
128
129 Some(deletion_file)
130 }
131 };
132 Ok(deletion_file)
133}
134
135#[instrument(
136 level = "debug",
137 skip(base, object_store),
138 fields(
139 base = base.as_ref(),
140 bytes_read = tracing::field::Empty
141 )
142)]
143pub async fn read_deletion_file(
144 fragment_id: u64,
145 deletion_file: &DeletionFile,
146 base: &Path,
147 object_store: &ObjectStore,
148) -> Result<DeletionVector> {
149 let span = tracing::Span::current();
150 match deletion_file.file_type {
151 DeletionFileType::Array => {
152 let path = deletion_file_path(base, fragment_id, deletion_file);
153
154 let data = object_store.read_one_all(&path).await?;
155 span.record("bytes_read", data.len());
156 let data = std::io::Cursor::new(data);
157 let mut batches: Vec<RecordBatch> = ArrowFileReader::try_new(data, None)?
158 .collect::<std::result::Result<_, ArrowError>>()
159 .map_err(box_error)
160 .context(CorruptFileSnafu { path: path.clone() })?;
161
162 if batches.len() != 1 {
163 return Err(Error::corrupt_file(
164 path,
165 format!(
166 "Expected exactly one batch in deletion file, got {}",
167 batches.len()
168 ),
169 ));
170 }
171
172 let batch = batches.pop().unwrap();
173 if batch.schema() != deletion_arrow_schema() {
174 return Err(Error::corrupt_file(
175 path,
176 format!(
177 "Expected schema {:?} in deletion file, got {:?}",
178 deletion_arrow_schema(),
179 batch.schema()
180 ),
181 ));
182 }
183
184 let array = batch.columns()[0]
185 .as_any()
186 .downcast_ref::<UInt32Array>()
187 .unwrap();
188
189 let mut set = HashSet::with_capacity(array.len());
190 for val in array.iter() {
191 if let Some(val) = val {
192 set.insert(val);
193 } else {
194 return Err(Error::corrupt_file(
195 path,
196 "Null values are not allowed in deletion files",
197 ));
198 }
199 }
200
201 Ok(DeletionVector::Set(set))
202 }
203 DeletionFileType::Bitmap => {
204 let path = deletion_file_path(base, fragment_id, deletion_file);
205
206 let data = object_store.read_one_all(&path).await?;
207 span.record("bytes_read", data.len());
208 let reader = data.reader();
209 let bitmap = RoaringBitmap::deserialize_from(reader)
210 .map_err(box_error)
211 .context(CorruptFileSnafu { path })?;
212
213 Ok(DeletionVector::Bitmap(bitmap))
214 }
215 }
216}
217
218#[cfg(test)]
219mod test {
220
221 use super::*;
222 use object_store::ObjectStoreExt;
223
224 #[tokio::test]
225 async fn test_write_no_deletions() {
226 let dv = DeletionVector::NoDeletions;
227
228 let (object_store, path) = ObjectStore::from_uri("memory:///no_deletion")
229 .await
230 .unwrap();
231 let file = write_deletion_file(&path, 0, 0, &dv, &object_store)
232 .await
233 .unwrap();
234 assert!(file.is_none());
235 }
236
237 #[tokio::test]
238 async fn test_write_array() {
239 let dv = DeletionVector::Set(HashSet::from_iter(0..100));
240
241 let fragment_id = 21;
242 let read_version = 12;
243
244 let object_store = ObjectStore::memory();
245 let path = Path::from("/write");
246 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
247 .await
248 .unwrap();
249
250 assert!(matches!(
251 file,
252 Some(DeletionFile {
253 file_type: DeletionFileType::Array,
254 ..
255 })
256 ));
257
258 let file = file.unwrap();
259 assert_eq!(file.read_version, read_version);
260 let path = deletion_file_path(&path, fragment_id, &file);
261 assert_eq!(
262 path,
263 Path::from(format!("/write/_deletions/21-12-{}.arrow", file.id))
264 );
265
266 let data = object_store
267 .inner
268 .get(&path)
269 .await
270 .unwrap()
271 .bytes()
272 .await
273 .unwrap();
274 let data = std::io::Cursor::new(data);
275 let mut batches: Vec<RecordBatch> = ArrowFileReader::try_new(data, None)
276 .unwrap()
277 .collect::<std::result::Result<_, ArrowError>>()
278 .unwrap();
279
280 assert_eq!(batches.len(), 1);
281 let batch = batches.pop().unwrap();
282 assert_eq!(batch.schema(), deletion_arrow_schema());
283 let array = batch["row_id"]
284 .as_any()
285 .downcast_ref::<UInt32Array>()
286 .unwrap();
287 let read_dv = DeletionVector::from_iter(array.iter().map(|v| v.unwrap()));
288 assert_eq!(read_dv, dv);
289 }
290
291 #[tokio::test]
292 async fn test_write_bitmap() {
293 let dv = DeletionVector::Bitmap(RoaringBitmap::from_iter(0..100));
294
295 let fragment_id = 21;
296 let read_version = 12;
297
298 let object_store = ObjectStore::memory();
299 let path = Path::from("/bitmap");
300 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
301 .await
302 .unwrap();
303
304 assert!(matches!(
305 file,
306 Some(DeletionFile {
307 file_type: DeletionFileType::Bitmap,
308 ..
309 })
310 ));
311
312 let file = file.unwrap();
313 assert_eq!(file.read_version, read_version);
314 let path = deletion_file_path(&path, fragment_id, &file);
315 assert_eq!(
316 path,
317 Path::from(format!("/bitmap/_deletions/21-12-{}.bin", file.id))
318 );
319
320 let data = object_store
321 .inner
322 .get(&path)
323 .await
324 .unwrap()
325 .bytes()
326 .await
327 .unwrap();
328 let reader = data.reader();
329 let read_bitmap = RoaringBitmap::deserialize_from(reader).unwrap();
330 assert_eq!(read_bitmap, dv.into_iter().collect::<RoaringBitmap>());
331 }
332
333 #[tokio::test]
334 async fn test_roundtrip_array() {
335 let dv = DeletionVector::Set(HashSet::from_iter(0..100));
336
337 let fragment_id = 21;
338 let read_version = 12;
339
340 let object_store = ObjectStore::memory();
341 let path = Path::from("/roundtrip");
342 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
343 .await
344 .unwrap();
345
346 let read_dv = read_deletion_file(fragment_id, &file.unwrap(), &path, &object_store)
347 .await
348 .unwrap();
349 assert_eq!(read_dv, dv);
350 }
351
352 #[tokio::test]
353 async fn test_roundtrip_bitmap() {
354 let dv = DeletionVector::Bitmap(RoaringBitmap::from_iter(0..100));
355
356 let fragment_id = 21;
357 let read_version = 12;
358
359 let object_store = ObjectStore::memory();
360 let path = Path::from("/bitmap");
361 let file = write_deletion_file(&path, fragment_id, read_version, &dv, &object_store)
362 .await
363 .unwrap();
364
365 let read_dv = read_deletion_file(fragment_id, &file.unwrap(), &path, &object_store)
366 .await
367 .unwrap();
368 assert_eq!(read_dv, dv);
369 }
370}