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};
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];
44// In 2.1+, we split large pages on read instead of write to avoid empty pages
45// and small pages issues. However, we keep the write-time limit at 32MB to avoid
46// potential regressions in 2.0 format readers.
47//
48// This limit is not applied in the 2.1 writer
49const MAX_PAGE_BYTES: usize = 32 * 1024 * 1024;
50// Total in-memory budget for buffering serialized page metadata before flushing
51// to the spill file. Divided evenly across columns (with a floor of 64 bytes).
52const DEFAULT_SPILL_BUFFER_LIMIT: usize = 256 * 1024;
53
54/// Spills serialized page metadata to a temporary file to bound memory usage.
55///
56/// The spill file is an unstructured sequence of "chunks". Each chunk is a
57/// contiguous run of length-delimited protobuf `Page` messages belonging to a
58/// single column. Chunks from different columns are interleaved in the order
59/// they are flushed (i.e. whenever a column's in-memory buffer exceeds
60/// `per_column_limit`). The `column_chunks` index records the (offset, length)
61/// of every chunk so each column's pages can be read back and reassembled in
62/// order.
63struct PageMetadataSpill {
64    writer: Box<dyn ObjectWriter>,
65    object_store: Arc<ObjectStore>,
66    path: Path,
67    /// Current write position in the spill file.
68    position: u64,
69    /// Per-column buffer of serialized (length-delimited protobuf) page metadata
70    /// that has not yet been flushed to the spill file.
71    column_buffers: Vec<Vec<u8>>,
72    /// Per-column list of chunks that have been flushed to the spill file.
73    /// Each entry is (offset, length) pointing into the spill file.
74    column_chunks: Vec<Vec<(u64, u32)>>,
75    /// Maximum bytes to buffer per column before flushing to the spill file.
76    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
155/// A writer for the Lance v2.0 file grammar.
156pub 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    // The number of rows written for each top-level field (i.e. each entry in
165    // `column_writers`). With `write_batch` every field advances together and
166    // these are all equal, but `write_column` advances one field at a time, so
167    // a single file may end up with columns of differing item counts.
168    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    /// Create a new v2.0 writer with a desired output schema.
187    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    /// Create a new v2.0 writer without a desired output schema.
198    ///
199    /// The output schema will be set based on the first batch of data to arrive.
200    /// If no data arrives and the writer is finished then the write will fail.
201    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    /// Spill page metadata to a sidecar file instead of accumulating in memory.
220    ///
221    /// This can dramatically reduce memory usage when many writers are open
222    /// concurrently (e.g. IVF shuffle with thousands of partition writers).
223    /// The sidecar file is created lazily on the first page write. The caller
224    /// is responsible for cleaning up `path` (e.g. by placing it in a temp
225    /// directory that is removed via RAII).
226    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        // As soon as an encoding task is done we write it.  There is no parallelism
282        // needed here because "writing" is really just submitting the buffer to the
283        // underlying write scheduler (either the OS or object_store's scheduler for
284        // cloud writes).  The only time we might truly await on write_page is if the
285        // scheduler's write queue is full.
286        //
287        // Also, there is no point in trying to make write_page parallel anyways
288        // because we wouldn't want buffers getting mixed up across pages.
289        while let Some(encoding_task) = encoding_tasks.next().await {
290            let encoded_page = encoding_task?;
291            self.write_page(encoded_page).await?;
292        }
293        // It's important to flush here, we don't know when the next batch will arrive
294        // and the underlying cloud store could have writes in progress that won't advance
295        // until we interact with the writer again.  These in-progress writes will time out
296        // if we don't flush.
297        self.writer.flush().await?;
298        Ok(())
299    }
300
301    /// Schedule batches of data to be written to the file
302    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    // Encode a set of `(field index, array)` pairs, each advancing only its own
421    // column. Each task captures its field's current row offset at encode time,
422    // so `advance_columns` must run after this call (never before); the order of
423    // the returned tasks relative to `write_pages` does not matter.
424    fn encode_columns(
425        &mut self,
426        field_arrays: &[(usize, ArrayRef)],
427        external_buffers: &mut OutOfLineBuffers,
428    ) -> Result<Vec<Vec<EncodeTask>>> {
429        // Snapshot the starting row number of each field before borrowing the
430        // column writers mutably below.
431        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    // Advance the per-field row counters after a set of columns has been
453    // written, keeping `rows_written` (the file's logical length) in sync as the
454    // longest column. Only the written fields move, so their new totals fold into
455    // `rows_written` directly without rescanning every field. (`write_batch`
456    // advances every field uniformly and tracks this inline instead.)
457    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    /// Schedule a batch of data to be written to the file
466    ///
467    /// Note: the future returned by this method may complete before the data has been fully
468    /// flushed to the file (some data may be in the data cache or the I/O cache)
469    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        // First we push each array into its column writer.  This may or may not generate enough
488        // data to trigger an encoding task.  We collect any encoding tasks into a queue.
489        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        // Next, write external buffers
493        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        // `write_batch` advances every field by the same amount, so the longest
503        // column simply grows by `num_rows`. Guard against overflowing the row
504        // counter.
505        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    /// Write a single column, advancing only that column's row counter.
519    ///
520    /// Unlike [`write_batch`](Self::write_batch), which advances every column
521    /// from a single shared row counter, this method advances one column
522    /// independently. Used across calls it produces a single file whose columns
523    /// may have different item counts.
524    ///
525    /// `column_index` refers to a top-level field in the writer's schema (the
526    /// same order as the schema's fields); a nested child cannot be targeted on
527    /// its own. Because each call writes the whole field from a single array, the
528    /// children of a struct field always advance together and stay equal-length;
529    /// only different top-level fields can diverge in length. A column may be
530    /// written across multiple calls; its values are appended. A field that is
531    /// never written ends up as a zero-length column. The writer must have been
532    /// created with an explicit schema (via [`try_new`](Self::try_new)); a lazy
533    /// schema cannot be inferred here because individual calls need not cover
534    /// every field.
535    ///
536    /// ```
537    /// # use arrow_array::{ArrayRef, Int32Array};
538    /// # use std::sync::Arc;
539    /// # use lance_file::writer::FileWriter;
540    /// # async fn example(writer: &mut FileWriter) -> lance_core::Result<()> {
541    /// // Field 0 gets three values, field 1 gets one — a non-rectangular file.
542    /// writer.write_column(0, Arc::new(Int32Array::from(vec![1, 2, 3]))).await?;
543    /// writer.write_column(1, Arc::new(Int32Array::from(vec![10]))).await?;
544    /// # Ok(())
545    /// # }
546    /// ```
547    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        // A never-advanced field simply remains a zero-length column, which the
571        // encoders handle at `finish` time.
572        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        // If spilling, finalize the spill writer and reopen for reading.
608        // The spill file itself is cleaned up by the caller (it lives in a
609        // temp directory managed by the caller's RAII guard).
610        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        // Use descriptor layout for blob v2 fields in the footer to avoid exposing logical child fields.
658        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    /// Add a metadata entry to the schema
675    ///
676    /// This method is useful because sometimes the metadata is not known until after the
677    /// data has been written.  This method allows you to alter the schema metadata.  It
678    /// must be called before `finish` is called.
679    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    /// Prepare the writer when column data and metadata were produced externally.
684    ///
685    /// This is useful for flows that copy already-encoded pages (e.g., binary copy
686    /// during compaction) where the column buffers have been written directly and we
687    /// only need to write the footer and schema metadata. The provided
688    /// `column_metadata` must describe the buffers already persisted by the
689    /// underlying `ObjectWriter`, and `rows_written` should reflect the total number
690    /// of rows in those buffers.
691    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    /// Adds a global buffer to the file
704    ///
705    /// The global buffer can contain any arbitrary bytes.  It will be written to the disk
706    /// immediately.  This method returns the index of the global buffer (this will always
707    /// start at 1 and increment by 1 each time this method is called)
708    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    /// Finishes writing the file
767    ///
768    /// This method will wait until all data has been flushed to the file.  Then it
769    /// will write the file metadata and the footer.  It will not return until all
770    /// data has been flushed and the file has been closed.
771    ///
772    /// Returns a summary of the completed file write.
773    pub async fn finish(&mut self) -> Result<FileWriteSummary> {
774        // 1. flush any remaining data and write out those pages
775        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        // 3. write global buffers (we write the schema here)
796        let global_buffer_offsets = self.write_global_buffers().await?;
797        let num_global_buffers = global_buffer_offsets.len() as u32;
798
799        // 4. write the column metadatas
800        let column_metadata_start = self.writer.tell().await? as u64;
801        let metadata_positions = self.write_column_metadatas().await?;
802
803        // 5. write the column metadata offset table
804        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        // 6. write global buffers offset table
811        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        // 7. write the footer
818        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        // 7. close the writer
828        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        // For multipart uploads, ObjectWriter's Drop impl will abort
838        // the upload when the writer is dropped.
839    }
840
841    pub async fn tell(&mut self) -> Result<u64> {
842        Ok(self.writer.tell().await? as u64)
843    }
844
845    /// Append a buffer whose metadata is supplied by the caller.
846    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
857// Creates a lance footer and appends it to the encoded data
858//
859// The logic here is very similar to logic in the FileWriter except we
860// are using BufMut (put_xyz) instead of AsyncWrite (write_xyz).
861pub fn concat_lance_footer(batch: &EncodedBatch, write_schema: bool) -> Result<Bytes> {
862    // Estimating 1MiB for file footer
863    let mut data = BytesMut::with_capacity(batch.data.len() + 1024 * 1024);
864    data.put(batch.data.clone());
865    // write global buffers (we write the schema here)
866    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    // Write column metadata
882    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    // Write column metadata offsets table
933    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    // Write global buffers offsets table
939    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    // write the footer
947    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}