Skip to main content

lance_file/versions/v2_0/
writer.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4use 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];
42// In 2.1+, we split large pages on read instead of write to avoid empty pages
43// and small pages issues. However, we keep the write-time limit at 32MB to avoid
44// potential regressions in 2.0 format readers.
45//
46// This limit is not applied in the 2.1 writer
47const MAX_PAGE_BYTES: usize = 32 * 1024 * 1024;
48// Total in-memory budget for buffering serialized page metadata before flushing
49// to the spill file. Divided evenly across columns (with a floor of 64 bytes).
50const DEFAULT_SPILL_BUFFER_LIMIT: usize = 256 * 1024;
51
52/// Spills serialized page metadata to a temporary file to bound memory usage.
53///
54/// The spill file is an unstructured sequence of "chunks". Each chunk is a
55/// contiguous run of length-delimited protobuf `Page` messages belonging to a
56/// single column. Chunks from different columns are interleaved in the order
57/// they are flushed (i.e. whenever a column's in-memory buffer exceeds
58/// `per_column_limit`). The `column_chunks` index records the (offset, length)
59/// of every chunk so each column's pages can be read back and reassembled in
60/// order.
61struct PageMetadataSpill {
62    writer: Box<dyn ObjectWriter>,
63    object_store: Arc<ObjectStore>,
64    path: Path,
65    /// Current write position in the spill file.
66    position: u64,
67    /// Per-column buffer of serialized (length-delimited protobuf) page metadata
68    /// that has not yet been flushed to the spill file.
69    column_buffers: Vec<Vec<u8>>,
70    /// Per-column list of chunks that have been flushed to the spill file.
71    /// Each entry is (offset, length) pointing into the spill file.
72    column_chunks: Vec<Vec<(u64, u32)>>,
73    /// Maximum bytes to buffer per column before flushing to the spill file.
74    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
153/// A writer for the Lance v2.0 file grammar.
154pub 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    // The number of rows written for each top-level field (i.e. each entry in
163    // `column_writers`). With `write_batch` every field advances together and
164    // these are all equal, but `write_column` advances one field at a time, so
165    // a single file may end up with columns of differing item counts.
166    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    /// Create a new v2.0 writer with a desired output schema.
185    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    /// Create a new v2.0 writer without a desired output schema.
196    ///
197    /// The output schema will be set based on the first batch of data to arrive.
198    /// If no data arrives and the writer is finished then the write will fail.
199    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    /// Spill page metadata to a sidecar file instead of accumulating in memory.
218    ///
219    /// This can dramatically reduce memory usage when many writers are open
220    /// concurrently (e.g. IVF shuffle with thousands of partition writers).
221    /// The sidecar file is created lazily on the first page write. The caller
222    /// is responsible for cleaning up `path` (e.g. by placing it in a temp
223    /// directory that is removed via RAII).
224    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        // As soon as an encoding task is done we write it.  There is no parallelism
280        // needed here because "writing" is really just submitting the buffer to the
281        // underlying write scheduler (either the OS or object_store's scheduler for
282        // cloud writes).  The only time we might truly await on write_page is if the
283        // scheduler's write queue is full.
284        //
285        // Also, there is no point in trying to make write_page parallel anyways
286        // because we wouldn't want buffers getting mixed up across pages.
287        while let Some(encoding_task) = encoding_tasks.next().await {
288            let encoded_page = encoding_task?;
289            self.write_page(encoded_page).await?;
290        }
291        // Flushing here reaps any upload that has already failed, so the error
292        // is attributed to this batch rather than to whichever later batch or
293        // the shutdown happens to poll the writer next. It does not wait for
294        // in-flight uploads: those are spawned tasks the runtime drives on its
295        // own, and blocking on them would stall the next batch behind them.
296        self.writer.flush().await?;
297        Ok(())
298    }
299
300    /// Schedule batches of data to be written to the file
301    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    // Encode a set of `(field index, array)` pairs, each advancing only its own
400    // column. Each task captures its field's current row offset at encode time,
401    // so `advance_columns` must run after this call (never before); the order of
402    // the returned tasks relative to `write_pages` does not matter.
403    fn encode_columns(
404        &mut self,
405        field_arrays: &[(usize, ArrayRef)],
406        external_buffers: &mut OutOfLineBuffers,
407    ) -> Result<Vec<Vec<EncodeTask>>> {
408        // Snapshot the starting row number of each field before borrowing the
409        // column writers mutably below.
410        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    // Advance the per-field row counters after a set of columns has been
432    // written, keeping `rows_written` (the file's logical length) in sync as the
433    // longest column. Only the written fields move, so their new totals fold into
434    // `rows_written` directly without rescanning every field. (`write_batch`
435    // advances every field uniformly and tracks this inline instead.)
436    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    /// Schedule a batch of data to be written to the file
445    ///
446    /// Note: the future returned by this method may complete before the data has been fully
447    /// flushed to the file (some data may be in the data cache or the I/O cache)
448    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        // First we push each array into its column writer.  This may or may not generate enough
468        // data to trigger an encoding task.  We collect any encoding tasks into a queue.
469        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        // Next, write external buffers
473        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        // `write_batch` advances every field by the same amount, so the longest
483        // column simply grows by `num_rows`. Guard against overflowing the row
484        // counter.
485        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    /// Write a single column, advancing only that column's row counter.
499    ///
500    /// Unlike [`write_batch`](Self::write_batch), which advances every column
501    /// from a single shared row counter, this method advances one column
502    /// independently. Used across calls it produces a single file whose columns
503    /// may have different item counts.
504    ///
505    /// `column_index` refers to a top-level field in the writer's schema (the
506    /// same order as the schema's fields); a nested child cannot be targeted on
507    /// its own. Because each call writes the whole field from a single array, the
508    /// children of a struct field always advance together and stay equal-length;
509    /// only different top-level fields can diverge in length. A column may be
510    /// written across multiple calls; its values are appended. A field that is
511    /// never written ends up as a zero-length column. The writer must have been
512    /// created with an explicit schema (via [`try_new`](Self::try_new)); a lazy
513    /// schema cannot be inferred here because individual calls need not cover
514    /// every field.
515    ///
516    /// ```
517    /// # use arrow_array::{ArrayRef, Int32Array};
518    /// # use std::sync::Arc;
519    /// # use lance_file::writer::FileWriter;
520    /// # async fn example(writer: &mut FileWriter) -> lance_core::Result<()> {
521    /// // Field 0 gets three values, field 1 gets one — a non-rectangular file.
522    /// writer.write_column(0, Arc::new(Int32Array::from(vec![1, 2, 3]))).await?;
523    /// writer.write_column(1, Arc::new(Int32Array::from(vec![10]))).await?;
524    /// # Ok(())
525    /// # }
526    /// ```
527    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        // A never-advanced field simply remains a zero-length column, which the
551        // encoders handle at `finish` time.
552        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        // If spilling, finalize the spill writer and reopen for reading.
588        // The spill file itself is cleaned up by the caller (it lives in a
589        // temp directory managed by the caller's RAII guard).
590        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        // Use descriptor layout for blob v2 fields in the footer to avoid exposing logical child fields.
638        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    /// Add a metadata entry to the schema
655    ///
656    /// This method is useful because sometimes the metadata is not known until after the
657    /// data has been written.  This method allows you to alter the schema metadata.  It
658    /// must be called before `finish` is called.
659    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    /// Prepare the writer when column data and metadata were produced externally.
664    ///
665    /// This is useful for flows that copy already-encoded pages (e.g., binary copy
666    /// during compaction) where the column buffers have been written directly and we
667    /// only need to write the footer and schema metadata. The provided
668    /// `column_metadata` must describe the buffers already persisted by the
669    /// underlying `ObjectWriter`, and `rows_written` should reflect the total number
670    /// of rows in those buffers.
671    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    /// Adds a global buffer to the file
686    ///
687    /// The global buffer can contain any arbitrary bytes.  It will be written to the disk
688    /// immediately.  This method returns the index of the global buffer (this will always
689    /// start at 1 and increment by 1 each time this method is called)
690    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    /// Finishes writing the file
749    ///
750    /// This method will wait until all data has been flushed to the file.  Then it
751    /// will write the file metadata and the footer.  It will not return until all
752    /// data has been flushed and the file has been closed.
753    ///
754    /// Returns a summary of the completed file write.
755    pub async fn finish(&mut self) -> Result<FileWriteSummary> {
756        // 1. flush any remaining data and write out those pages
757        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        // 3. write global buffers (we write the schema here)
778        let global_buffer_offsets = self.write_global_buffers().await?;
779        let num_global_buffers = global_buffer_offsets.len() as u32;
780
781        // 4. write the column metadatas
782        let column_metadata_start = self.writer.tell().await? as u64;
783        let metadata_positions = self.write_column_metadatas().await?;
784
785        // 5. write the column metadata offset table
786        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        // 6. write global buffers offset table
793        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        // 7. write the footer
800        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        // 7. close the writer
810        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        // For multipart uploads, ObjectWriter's Drop impl will abort
820        // the upload when the writer is dropped.
821    }
822
823    pub async fn tell(&mut self) -> Result<u64> {
824        Ok(self.writer.tell().await? as u64)
825    }
826
827    /// Append a buffer whose metadata is supplied by the caller.
828    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
839// Creates a lance footer and appends it to the encoded data
840//
841// The logic here is very similar to logic in the FileWriter except we
842// are using BufMut (put_xyz) instead of AsyncWrite (write_xyz).
843pub fn concat_lance_footer(batch: &EncodedBatch, write_schema: bool) -> Result<Bytes> {
844    // Estimating 1MiB for file footer
845    let mut data = BytesMut::with_capacity(batch.data.len() + 1024 * 1024);
846    data.put(batch.data.clone());
847    // write global buffers (we write the schema here)
848    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    // Write column metadata
864    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    // Write column metadata offsets table
915    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    // Write global buffers offsets table
921    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    // write the footer
929    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}