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