1use core::panic;
5use std::collections::HashMap;
6use std::sync::Arc;
7
8use arrow_array::{ArrayRef, RecordBatch};
9
10use arrow_data::ArrayData;
11use bytes::{Buf, BufMut, Bytes, BytesMut};
12use futures::StreamExt;
13use futures::stream::FuturesOrdered;
14use lance_core::datatypes::{Field, Schema as LanceSchema};
15use lance_core::utils::bit::pad_bytes;
16use lance_core::{Error, Result};
17use lance_encoding::decoder::PageEncoding;
18use lance_encoding::encoder::{
19 ArrayFieldEncodingStrategy, BatchEncoder, EncodeTask, EncodedBatch, EncodedPage,
20 EncodingOptions, FieldEncoder, FieldEncodingStrategy, OutOfLineBuffers,
21};
22use lance_encoding::repdef::RepDefBuilder;
23use lance_io::object_store::ObjectStore;
24use lance_io::traits::Writer as ObjectWriter;
25use log::{debug, warn};
26use object_store::path::Path;
27use prost::Message;
28use prost_types::Any;
29use tokio::io::AsyncWrite;
30use tokio::io::AsyncWriteExt;
31use tracing::instrument;
32
33use crate::datatypes::FieldsWithMeta;
34use crate::format::MAGIC;
35use crate::format::pb;
36use crate::format::pbfile;
37use crate::format::pbfile::DirectEncoding;
38use crate::writer::{
39 ENV_LANCE_FILE_WRITER_MAX_PAGE_BYTES, FileWriteSummary, FileWriterOptions,
40 PAGE_BUFFER_ALIGNMENT,
41};
42
43const PAD_BUFFER: [u8; PAGE_BUFFER_ALIGNMENT] = [72; PAGE_BUFFER_ALIGNMENT];
44const MAX_PAGE_BYTES: usize = 32 * 1024 * 1024;
50const DEFAULT_SPILL_BUFFER_LIMIT: usize = 256 * 1024;
53
54struct PageMetadataSpill {
64 writer: Box<dyn ObjectWriter>,
65 object_store: Arc<ObjectStore>,
66 path: Path,
67 position: u64,
69 column_buffers: Vec<Vec<u8>>,
72 column_chunks: Vec<Vec<(u64, u32)>>,
75 per_column_limit: usize,
77}
78
79impl PageMetadataSpill {
80 async fn new(object_store: Arc<ObjectStore>, path: Path, num_columns: usize) -> Result<Self> {
81 let writer = object_store.create(&path).await?;
82 let per_column_limit = (DEFAULT_SPILL_BUFFER_LIMIT / num_columns.max(1)).max(64);
83 Ok(Self {
84 writer,
85 object_store,
86 path,
87 position: 0,
88 column_buffers: vec![Vec::new(); num_columns],
89 column_chunks: vec![Vec::new(); num_columns],
90 per_column_limit,
91 })
92 }
93
94 async fn append_page(
95 &mut self,
96 column_idx: usize,
97 page: &pbfile::column_metadata::Page,
98 ) -> Result<()> {
99 page.encode_length_delimited(&mut self.column_buffers[column_idx])
100 .map_err(|e| {
101 Error::io_source(Box::new(std::io::Error::new(
102 std::io::ErrorKind::InvalidData,
103 e,
104 )))
105 })?;
106 if self.column_buffers[column_idx].len() >= self.per_column_limit {
107 self.flush_column(column_idx).await?;
108 }
109 Ok(())
110 }
111
112 async fn flush_column(&mut self, column_idx: usize) -> Result<()> {
113 let buf = &self.column_buffers[column_idx];
114 if buf.is_empty() {
115 return Ok(());
116 }
117 let len = buf.len();
118 self.writer.write_all(buf).await?;
119 self.column_chunks[column_idx].push((self.position, len as u32));
120 self.position += len as u64;
121 self.column_buffers[column_idx].clear();
122 Ok(())
123 }
124
125 async fn shutdown_writer(&mut self) -> Result<()> {
126 for col_idx in 0..self.column_buffers.len() {
127 self.flush_column(col_idx).await?;
128 }
129 ObjectWriter::shutdown(self.writer.as_mut()).await?;
130 Ok(())
131 }
132}
133
134fn decode_spilled_chunk(data: &Bytes) -> Result<Vec<pbfile::column_metadata::Page>> {
135 let mut pages = Vec::new();
136 let mut cursor = data.clone();
137 while cursor.has_remaining() {
138 let page =
139 pbfile::column_metadata::Page::decode_length_delimited(&mut cursor).map_err(|e| {
140 Error::io_source(Box::new(std::io::Error::new(
141 std::io::ErrorKind::InvalidData,
142 e,
143 )))
144 })?;
145 pages.push(page);
146 }
147 Ok(pages)
148}
149
150enum PageSpillState {
151 Pending(Arc<ObjectStore>, Path),
152 Active(PageMetadataSpill),
153}
154
155pub struct Writer {
157 writer: Box<dyn ObjectWriter>,
158 schema: Option<LanceSchema>,
159 column_writers: Vec<Box<dyn FieldEncoder>>,
160 column_metadata: Vec<pbfile::ColumnMetadata>,
161 field_id_to_column_indices: Vec<(u32, u32)>,
162 num_columns: u32,
163 rows_written: u64,
164 field_rows_written: Vec<u64>,
169 global_buffers: Vec<(u64, u64)>,
170 schema_metadata: HashMap<String, String>,
171 encoding_strategy: Box<dyn FieldEncodingStrategy>,
172 options: FileWriterOptions,
173 page_spill: Option<PageSpillState>,
174}
175
176fn initial_column_metadata() -> pbfile::ColumnMetadata {
177 pbfile::ColumnMetadata {
178 pages: Vec::new(),
179 buffer_offsets: Vec::new(),
180 buffer_sizes: Vec::new(),
181 encoding: None,
182 }
183}
184
185impl Writer {
186 pub fn try_new(
188 object_writer: Box<dyn ObjectWriter>,
189 schema: LanceSchema,
190 options: FileWriterOptions,
191 ) -> Result<Self> {
192 let mut writer = Self::new_lazy(object_writer, options);
193 writer.initialize(schema)?;
194 Ok(writer)
195 }
196
197 pub fn new_lazy(object_writer: Box<dyn ObjectWriter>, options: FileWriterOptions) -> Self {
202 Self {
203 writer: object_writer,
204 schema: None,
205 column_writers: Vec::new(),
206 column_metadata: Vec::new(),
207 num_columns: 0,
208 rows_written: 0,
209 field_rows_written: Vec::new(),
210 field_id_to_column_indices: Vec::new(),
211 global_buffers: Vec::new(),
212 schema_metadata: HashMap::new(),
213 page_spill: None,
214 encoding_strategy: Box::new(ArrayFieldEncodingStrategy::new()),
215 options,
216 }
217 }
218
219 pub fn with_page_metadata_spill(mut self, object_store: Arc<ObjectStore>, path: Path) -> Self {
227 self.page_spill = Some(PageSpillState::Pending(object_store, path));
228 self
229 }
230
231 async fn do_write_buffer(writer: &mut (impl AsyncWrite + Unpin), buf: &[u8]) -> Result<()> {
232 writer.write_all(buf).await?;
233 let pad_bytes = pad_bytes::<PAGE_BUFFER_ALIGNMENT>(buf.len());
234 writer.write_all(&PAD_BUFFER[..pad_bytes]).await?;
235 Ok(())
236 }
237
238 async fn write_page(&mut self, encoded_page: EncodedPage) -> Result<()> {
239 let buffers = encoded_page.data;
240 let mut buffer_offsets = Vec::with_capacity(buffers.len());
241 let mut buffer_sizes = Vec::with_capacity(buffers.len());
242 for buffer in buffers {
243 buffer_offsets.push(self.writer.tell().await? as u64);
244 buffer_sizes.push(buffer.len() as u64);
245 Self::do_write_buffer(&mut self.writer, &buffer).await?;
246 }
247 let encoded_encoding = match encoded_page.description {
248 PageEncoding::Legacy(array_encoding) => Any::from_msg(&array_encoding)?.encode_to_vec(),
249 PageEncoding::Structural(page_layout) => Any::from_msg(&page_layout)?.encode_to_vec(),
250 };
251 let page = pbfile::column_metadata::Page {
252 buffer_offsets,
253 buffer_sizes,
254 encoding: Some(pbfile::Encoding {
255 location: Some(pbfile::encoding::Location::Direct(DirectEncoding {
256 encoding: encoded_encoding,
257 })),
258 }),
259 length: encoded_page.num_rows,
260 priority: encoded_page.row_number,
261 };
262 let col_idx = encoded_page.column_idx as usize;
263 if matches!(&self.page_spill, Some(PageSpillState::Pending(..))) {
264 let Some(PageSpillState::Pending(store, path)) = self.page_spill.take() else {
265 unreachable!()
266 };
267 self.page_spill = Some(PageSpillState::Active(
268 PageMetadataSpill::new(store, path, self.num_columns as usize).await?,
269 ));
270 }
271 match &mut self.page_spill {
272 Some(PageSpillState::Active(spill)) => spill.append_page(col_idx, &page).await?,
273 None => self.column_metadata[col_idx].pages.push(page),
274 Some(PageSpillState::Pending(..)) => unreachable!(),
275 }
276 Ok(())
277 }
278
279 #[instrument(skip_all, level = "debug")]
280 async fn write_pages(&mut self, mut encoding_tasks: FuturesOrdered<EncodeTask>) -> Result<()> {
281 while let Some(encoding_task) = encoding_tasks.next().await {
290 let encoded_page = encoding_task?;
291 self.write_page(encoded_page).await?;
292 }
293 self.writer.flush().await?;
298 Ok(())
299 }
300
301 pub async fn write_batches(
303 &mut self,
304 batches: impl Iterator<Item = &RecordBatch>,
305 ) -> Result<()> {
306 for batch in batches {
307 self.write_batch(batch).await?;
308 }
309 Ok(())
310 }
311
312 fn verify_field_nullability(arr: &ArrayData, field: &Field) -> Result<()> {
313 if !field.nullable && arr.null_count() > 0 {
314 return Err(Error::invalid_input(format!(
315 "The field `{}` contained null values even though the field is marked non-null in the schema",
316 field.name
317 )));
318 }
319
320 for (child_field, child_arr) in field.children.iter().zip(arr.child_data()) {
321 Self::verify_field_nullability(child_arr, child_field)?;
322 }
323
324 Ok(())
325 }
326
327 fn verify_nullability_constraints(&self, batch: &RecordBatch) -> Result<()> {
328 for (col, field) in batch
329 .columns()
330 .iter()
331 .zip(self.schema.as_ref().unwrap().fields.iter())
332 {
333 Self::verify_field_nullability(&col.to_data(), field)?;
334 }
335 Ok(())
336 }
337
338 fn initialize(&mut self, mut schema: LanceSchema) -> Result<()> {
339 let cache_bytes_per_column = if let Some(data_cache_bytes) = self.options.data_cache_bytes {
340 data_cache_bytes / schema.fields.len() as u64
341 } else {
342 8 * 1024 * 1024
343 };
344
345 let max_page_bytes = self.options.max_page_bytes.unwrap_or_else(|| {
346 std::env::var(ENV_LANCE_FILE_WRITER_MAX_PAGE_BYTES)
347 .map(|s| {
348 s.parse::<u64>().unwrap_or_else(|e| {
349 warn!(
350 "Failed to parse {}: {}, using default",
351 ENV_LANCE_FILE_WRITER_MAX_PAGE_BYTES, e
352 );
353 MAX_PAGE_BYTES as u64
354 })
355 })
356 .unwrap_or(MAX_PAGE_BYTES as u64)
357 });
358
359 schema.validate()?;
360
361 let keep_original_array = self.options.keep_original_array.unwrap_or(false);
362 let encoding_options = EncodingOptions {
363 cache_bytes_per_column,
364 max_page_bytes,
365 keep_original_array,
366 buffer_alignment: PAGE_BUFFER_ALIGNMENT as u64,
367 };
368 let encoder =
369 BatchEncoder::try_new(&schema, self.encoding_strategy.as_ref(), &encoding_options)?;
370 self.num_columns = encoder.num_columns();
371
372 self.field_rows_written = vec![0; encoder.field_encoders.len()];
373 self.column_writers = encoder.field_encoders;
374 self.column_metadata = vec![initial_column_metadata(); self.num_columns as usize];
375 self.field_id_to_column_indices = encoder.field_id_to_column_index;
376 self.schema_metadata
377 .extend(std::mem::take(&mut schema.metadata));
378 self.schema = Some(schema);
379 Ok(())
380 }
381
382 fn ensure_initialized(&mut self, batch: &RecordBatch) -> Result<&LanceSchema> {
383 if self.schema.is_none() {
384 let schema = LanceSchema::try_from(batch.schema().as_ref())?;
385 self.initialize(schema)?;
386 }
387 Ok(self.schema.as_ref().unwrap())
388 }
389
390 #[instrument(skip_all, level = "debug")]
391 fn encode_batch(
392 &mut self,
393 batch: &RecordBatch,
394 external_buffers: &mut OutOfLineBuffers,
395 ) -> Result<Vec<Vec<EncodeTask>>> {
396 let field_arrays = self
397 .schema
398 .as_ref()
399 .unwrap()
400 .fields
401 .iter()
402 .enumerate()
403 .map(|(field_idx, field)| {
404 let array =
405 batch
406 .column_by_name(&field.name)
407 .ok_or(Error::invalid_input_source(
408 format!(
409 "Cannot write batch. The batch was missing the column `{}`",
410 field.name
411 )
412 .into(),
413 ))?;
414 Ok((field_idx, array.clone()))
415 })
416 .collect::<Result<Vec<_>>>()?;
417 self.encode_columns(&field_arrays, external_buffers)
418 }
419
420 fn encode_columns(
425 &mut self,
426 field_arrays: &[(usize, ArrayRef)],
427 external_buffers: &mut OutOfLineBuffers,
428 ) -> Result<Vec<Vec<EncodeTask>>> {
429 let row_numbers = field_arrays
432 .iter()
433 .map(|(field_idx, _)| self.field_rows_written[*field_idx])
434 .collect::<Vec<_>>();
435 field_arrays
436 .iter()
437 .zip(row_numbers)
438 .map(|((field_idx, array), row_number)| {
439 let repdef = RepDefBuilder::default();
440 let num_rows = array.len() as u64;
441 self.column_writers[*field_idx].maybe_encode(
442 array.clone(),
443 external_buffers,
444 repdef,
445 row_number,
446 num_rows,
447 )
448 })
449 .collect::<Result<Vec<_>>>()
450 }
451
452 fn advance_columns(&mut self, field_arrays: &[(usize, ArrayRef)]) {
458 for (field_idx, array) in field_arrays {
459 let new_total = self.field_rows_written[*field_idx] + array.len() as u64;
460 self.field_rows_written[*field_idx] = new_total;
461 self.rows_written = self.rows_written.max(new_total);
462 }
463 }
464
465 pub async fn write_batch(&mut self, batch: &RecordBatch) -> Result<()> {
470 debug!(
471 "write_batch called with {} rows, {} columns, and {} bytes of data",
472 batch.num_rows(),
473 batch.num_columns(),
474 batch.get_array_memory_size()
475 );
476 self.ensure_initialized(batch)?;
477 self.verify_nullability_constraints(batch)?;
478 let num_rows = batch.num_rows() as u64;
479 if num_rows == 0 {
480 return Ok(());
481 }
482 if num_rows > u32::MAX as u64 {
483 return Err(Error::invalid_input_source(
484 "cannot write Lance files with more than 2^32 rows".into(),
485 ));
486 }
487 let mut external_buffers =
490 OutOfLineBuffers::new(self.tell().await?, PAGE_BUFFER_ALIGNMENT as u64);
491 let encoding_tasks = self.encode_batch(batch, &mut external_buffers)?;
492 for external_buffer in external_buffers.take_buffers() {
494 Self::do_write_buffer(&mut self.writer, &external_buffer).await?;
495 }
496
497 let encoding_tasks = encoding_tasks
498 .into_iter()
499 .flatten()
500 .collect::<FuturesOrdered<_>>();
501
502 if self.rows_written.checked_add(num_rows).is_none() {
506 return Err(Error::invalid_input_source(format!("cannot write batch with {} rows because {} rows have already been written and Lance files cannot contain more than 2^64 rows", num_rows, self.rows_written).into()));
507 }
508 for field_rows in self.field_rows_written.iter_mut() {
509 *field_rows += num_rows;
510 }
511 self.rows_written += num_rows;
512
513 self.write_pages(encoding_tasks).await?;
514
515 Ok(())
516 }
517
518 pub async fn write_column(&mut self, column_index: usize, array: ArrayRef) -> Result<()> {
548 let schema = self.schema.as_ref().ok_or_else(|| {
549 Error::invalid_input_source(
550 "write_column requires the writer to be created with an explicit schema".into(),
551 )
552 })?;
553 let field = schema.fields.get(column_index).ok_or_else(|| {
554 Error::invalid_input_source(
555 format!(
556 "write_column: field index {} is out of bounds (schema has {} fields)",
557 column_index,
558 schema.fields.len()
559 )
560 .into(),
561 )
562 })?;
563 if array.len() as u64 > u32::MAX as u64 {
564 return Err(Error::invalid_input_source(
565 "cannot write Lance files with more than 2^32 rows".into(),
566 ));
567 }
568 Self::verify_field_nullability(&array.to_data(), field)?;
569
570 if array.is_empty() {
573 return Ok(());
574 }
575
576 let columns = [(column_index, array)];
577 let mut external_buffers =
578 OutOfLineBuffers::new(self.tell().await?, PAGE_BUFFER_ALIGNMENT as u64);
579 let encoding_tasks = self.encode_columns(&columns, &mut external_buffers)?;
580 for external_buffer in external_buffers.take_buffers() {
581 Self::do_write_buffer(&mut self.writer, &external_buffer).await?;
582 }
583 let encoding_tasks = encoding_tasks
584 .into_iter()
585 .flatten()
586 .collect::<FuturesOrdered<_>>();
587
588 self.advance_columns(&columns);
589 self.write_pages(encoding_tasks).await?;
590 Ok(())
591 }
592
593 async fn write_column_metadata(
594 &mut self,
595 metadata: pbfile::ColumnMetadata,
596 ) -> Result<(u64, u64)> {
597 let metadata_bytes = metadata.encode_to_vec();
598 let position = self.writer.tell().await? as u64;
599 let len = metadata_bytes.len() as u64;
600 self.writer.write_all(&metadata_bytes).await?;
601 Ok((position, len))
602 }
603
604 async fn write_column_metadatas(&mut self) -> Result<Vec<(u64, u64)>> {
605 let metadatas = std::mem::take(&mut self.column_metadata);
606
607 let spill_state = self.page_spill.take();
611 let (spill_chunks, spill_reader) =
612 if let Some(PageSpillState::Active(mut spill)) = spill_state {
613 spill.shutdown_writer().await?;
614 let reader = spill.object_store.open(&spill.path).await?;
615 let chunks = std::mem::take(&mut spill.column_chunks);
616 (chunks, Some(reader))
617 } else {
618 (Vec::new(), None)
619 };
620
621 let mut metadata_positions = Vec::with_capacity(metadatas.len());
622 for (col_idx, mut metadata) in metadatas.into_iter().enumerate() {
623 if let Some(reader) = &spill_reader {
624 let mut pages = Vec::new();
625 for &(offset, len) in &spill_chunks[col_idx] {
626 let data = reader
627 .get_range(offset as usize..(offset as usize + len as usize))
628 .await
629 .map_err(|e| Error::io_source(Box::new(e)))?;
630 pages.extend(decode_spilled_chunk(&data)?);
631 }
632 metadata.pages = pages;
633 }
634 metadata_positions.push(self.write_column_metadata(metadata).await?);
635 }
636
637 Ok(metadata_positions)
638 }
639
640 fn make_file_descriptor(
641 schema: &lance_core::datatypes::Schema,
642 num_rows: u64,
643 ) -> Result<pb::FileDescriptor> {
644 let fields_with_meta = FieldsWithMeta::from(schema);
645 Ok(pb::FileDescriptor {
646 schema: Some(pb::Schema {
647 fields: fields_with_meta.fields.0,
648 metadata: fields_with_meta.metadata,
649 }),
650 length: num_rows,
651 })
652 }
653
654 async fn write_global_buffers(&mut self) -> Result<Vec<(u64, u64)>> {
655 let schema = self.schema.as_mut().ok_or(Error::invalid_input("No schema provided on writer open and no data provided. Schema is unknown and file cannot be created"))?;
656 schema.metadata = std::mem::take(&mut self.schema_metadata);
657 schema
659 .fields
660 .iter_mut()
661 .for_each(|f| f.unload_blobs_recursive());
662
663 let file_descriptor = Self::make_file_descriptor(schema, self.rows_written)?;
664 let file_descriptor_bytes = file_descriptor.encode_to_vec();
665 let file_descriptor_len = file_descriptor_bytes.len() as u64;
666 let file_descriptor_position = self.writer.tell().await? as u64;
667 self.writer.write_all(&file_descriptor_bytes).await?;
668 let mut gbo_table = Vec::with_capacity(1 + self.global_buffers.len());
669 gbo_table.push((file_descriptor_position, file_descriptor_len));
670 gbo_table.append(&mut self.global_buffers);
671 Ok(gbo_table)
672 }
673
674 pub fn add_schema_metadata(&mut self, key: impl Into<String>, value: impl Into<String>) {
680 self.schema_metadata.insert(key.into(), value.into());
681 }
682
683 pub fn initialize_with_external_metadata(
692 &mut self,
693 schema: lance_core::datatypes::Schema,
694 column_metadata: Vec<pbfile::ColumnMetadata>,
695 rows_written: u64,
696 ) {
697 self.schema = Some(schema);
698 self.num_columns = column_metadata.len() as u32;
699 self.column_metadata = column_metadata;
700 self.rows_written = rows_written;
701 }
702
703 pub async fn add_global_buffer(&mut self, buffer: Bytes) -> Result<u32> {
709 let position = self.writer.tell().await? as u64;
710 let len = buffer.len() as u64;
711 Self::do_write_buffer(&mut self.writer, &buffer).await?;
712 self.global_buffers.push((position, len));
713 Ok(self.global_buffers.len() as u32)
714 }
715
716 async fn finish_writers(&mut self) -> Result<()> {
717 let mut col_idx = 0;
718 for mut writer in std::mem::take(&mut self.column_writers) {
719 let mut external_buffers =
720 OutOfLineBuffers::new(self.tell().await?, PAGE_BUFFER_ALIGNMENT as u64);
721 let columns = writer.finish(&mut external_buffers).await?;
722 for buffer in external_buffers.take_buffers() {
723 self.writer.write_all(&buffer).await?;
724 }
725 debug_assert_eq!(
726 columns.len(),
727 writer.num_columns() as usize,
728 "Expected {} columns from column at index {} and got {}",
729 writer.num_columns(),
730 col_idx,
731 columns.len()
732 );
733 for column in columns {
734 for page in column.final_pages {
735 self.write_page(page).await?;
736 }
737 let column_metadata = &mut self.column_metadata[col_idx];
738 let mut buffer_pos = self.writer.tell().await? as u64;
739 for buffer in column.column_buffers {
740 column_metadata.buffer_offsets.push(buffer_pos);
741 let mut size = 0;
742 Self::do_write_buffer(&mut self.writer, &buffer).await?;
743 size += buffer.len() as u64;
744 buffer_pos += size;
745 column_metadata.buffer_sizes.push(size);
746 }
747 let encoded_encoding = Any::from_msg(&column.encoding)?.encode_to_vec();
748 column_metadata.encoding = Some(pbfile::Encoding {
749 location: Some(pbfile::encoding::Location::Direct(pbfile::DirectEncoding {
750 encoding: encoded_encoding,
751 })),
752 });
753 col_idx += 1;
754 }
755 }
756 if col_idx != self.column_metadata.len() {
757 panic!(
758 "Column writers finished with {} columns but we expected {}",
759 col_idx,
760 self.column_metadata.len()
761 );
762 }
763 Ok(())
764 }
765
766 pub async fn finish(&mut self) -> Result<FileWriteSummary> {
774 let mut external_buffers =
776 OutOfLineBuffers::new(self.tell().await?, PAGE_BUFFER_ALIGNMENT as u64);
777 let encoding_tasks = self
778 .column_writers
779 .iter_mut()
780 .map(|writer| writer.flush(&mut external_buffers))
781 .collect::<Result<Vec<_>>>()?;
782 for external_buffer in external_buffers.take_buffers() {
783 Self::do_write_buffer(&mut self.writer, &external_buffer).await?;
784 }
785 let encoding_tasks = encoding_tasks
786 .into_iter()
787 .flatten()
788 .collect::<FuturesOrdered<_>>();
789 self.write_pages(encoding_tasks).await?;
790
791 if !self.column_writers.is_empty() {
792 self.finish_writers().await?;
793 }
794
795 let global_buffer_offsets = self.write_global_buffers().await?;
797 let num_global_buffers = global_buffer_offsets.len() as u32;
798
799 let column_metadata_start = self.writer.tell().await? as u64;
801 let metadata_positions = self.write_column_metadatas().await?;
802
803 let cmo_table_start = self.writer.tell().await? as u64;
805 for (meta_pos, meta_len) in metadata_positions {
806 self.writer.write_u64_le(meta_pos).await?;
807 self.writer.write_u64_le(meta_len).await?;
808 }
809
810 let gbo_table_start = self.writer.tell().await? as u64;
812 for (gbo_pos, gbo_len) in global_buffer_offsets {
813 self.writer.write_u64_le(gbo_pos).await?;
814 self.writer.write_u64_le(gbo_len).await?;
815 }
816
817 self.writer.write_u64_le(column_metadata_start).await?;
819 self.writer.write_u64_le(cmo_table_start).await?;
820 self.writer.write_u64_le(gbo_table_start).await?;
821 self.writer.write_u32_le(num_global_buffers).await?;
822 self.writer.write_u32_le(self.num_columns).await?;
823 self.writer.write_u16_le(0).await?;
824 self.writer.write_u16_le(3).await?;
825 self.writer.write_all(MAGIC).await?;
826
827 let write_result = ObjectWriter::shutdown(self.writer.as_mut()).await?;
829
830 Ok(FileWriteSummary {
831 num_rows: self.rows_written,
832 size_bytes: write_result.size as u64,
833 })
834 }
835
836 pub async fn abort(&mut self) {
837 }
840
841 pub async fn tell(&mut self) -> Result<u64> {
842 Ok(self.writer.tell().await? as u64)
843 }
844
845 pub async fn write_external_buffer(&mut self, bytes: &[u8]) -> Result<(u64, u64)> {
847 let start = self.tell().await?;
848 self.writer.write_all(bytes).await?;
849 Ok((start, bytes.len() as u64))
850 }
851
852 pub fn field_id_to_column_indices(&self) -> &[(u32, u32)] {
853 &self.field_id_to_column_indices
854 }
855}
856
857pub fn concat_lance_footer(batch: &EncodedBatch, write_schema: bool) -> Result<Bytes> {
862 let mut data = BytesMut::with_capacity(batch.data.len() + 1024 * 1024);
864 data.put(batch.data.clone());
865 let global_buffers = if write_schema {
867 let schema_start = data.len() as u64;
868 let lance_schema = lance_core::datatypes::Schema::try_from(batch.schema.as_ref())?;
869 let descriptor = Writer::make_file_descriptor(&lance_schema, batch.num_rows)?;
870 let descriptor_bytes = descriptor.encode_to_vec();
871 let descriptor_len = descriptor_bytes.len() as u64;
872 data.put(descriptor_bytes.as_slice());
873
874 vec![(schema_start, descriptor_len)]
875 } else {
876 vec![]
877 };
878 let col_metadata_start = data.len() as u64;
879
880 let mut col_metadata_positions = Vec::new();
881 for col in &batch.page_table {
883 let position = data.len() as u64;
884 let pages = col
885 .page_infos
886 .iter()
887 .map(|page_info| {
888 let encoded_encoding = match &page_info.encoding {
889 PageEncoding::Legacy(array_encoding) => {
890 Any::from_msg(array_encoding)?.encode_to_vec()
891 }
892 PageEncoding::Structural(page_layout) => {
893 Any::from_msg(page_layout)?.encode_to_vec()
894 }
895 };
896 let (buffer_offsets, buffer_sizes): (Vec<_>, Vec<_>) = page_info
897 .buffer_offsets_and_sizes
898 .as_ref()
899 .iter()
900 .cloned()
901 .unzip();
902 Ok(pbfile::column_metadata::Page {
903 buffer_offsets,
904 buffer_sizes,
905 encoding: Some(pbfile::Encoding {
906 location: Some(pbfile::encoding::Location::Direct(DirectEncoding {
907 encoding: encoded_encoding,
908 })),
909 }),
910 length: page_info.num_rows,
911 priority: page_info.priority,
912 })
913 })
914 .collect::<Result<Vec<_>>>()?;
915 let (buffer_offsets, buffer_sizes): (Vec<_>, Vec<_>) =
916 col.buffer_offsets_and_sizes.iter().cloned().unzip();
917 let encoded_col_encoding = Any::from_msg(&col.encoding)?.encode_to_vec();
918 let column = pbfile::ColumnMetadata {
919 pages,
920 buffer_offsets,
921 buffer_sizes,
922 encoding: Some(pbfile::Encoding {
923 location: Some(pbfile::encoding::Location::Direct(pbfile::DirectEncoding {
924 encoding: encoded_col_encoding,
925 })),
926 }),
927 };
928 let column_bytes = column.encode_to_vec();
929 col_metadata_positions.push((position, column_bytes.len() as u64));
930 data.put(column_bytes.as_slice());
931 }
932 let cmo_table_start = data.len() as u64;
934 for (meta_pos, meta_len) in col_metadata_positions {
935 data.put_u64_le(meta_pos);
936 data.put_u64_le(meta_len);
937 }
938 let gbo_table_start = data.len() as u64;
940 let num_global_buffers = global_buffers.len() as u32;
941 for (gbo_pos, gbo_len) in global_buffers {
942 data.put_u64_le(gbo_pos);
943 data.put_u64_le(gbo_len);
944 }
945
946 data.put_u64_le(col_metadata_start);
948 data.put_u64_le(cmo_table_start);
949 data.put_u64_le(gbo_table_start);
950 data.put_u32_le(num_global_buffers);
951 data.put_u32_le(batch.page_table.len() as u32);
952 data.put_u16_le(2);
953 data.put_u16_le(0);
954 data.put(MAGIC.as_slice());
955
956 Ok(data.freeze())
957}