Skip to main content

lance_file/versions/v2_0/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4//! Lance v2.0 file composition.
5
6use std::{collections::BTreeMap, sync::Arc};
7
8use bytes::Bytes;
9use lance_core::{
10    Result,
11    datatypes::{Field, Schema},
12};
13use lance_encoding::{
14    decoder::PageInfo,
15    encoder::{ArrayFieldEncodingStrategy, EncodedBatch, FieldEncodingStrategy},
16};
17use lance_io::traits::Writer as ObjectWriter;
18
19use crate::{reader::ReadProjection, writer::FileWriterOptions};
20
21mod reader;
22mod writer;
23
24#[cfg(test)]
25pub(crate) use reader::test_projection_length;
26pub(crate) use reader::{
27    decode_column_metadata, finish_metadata, finish_metadata_index, validate_global_buffers,
28};
29pub use reader::{
30    projection_from_column_names, projection_from_field_ids, projection_from_whole_schema,
31};
32pub use writer::Writer;
33
34pub(crate) fn read_projection() -> Arc<dyn ReadProjection> {
35    reader::read_projection()
36}
37
38/// Count physical columns represented by a field in a v2.0 footer.
39pub fn physical_column_count(field: &Field) -> usize {
40    if field.is_blob() || field.is_packed_struct() {
41        1
42    } else {
43        1 + field
44            .children
45            .iter()
46            .map(physical_column_count)
47            .sum::<usize>()
48    }
49}
50
51/// Build persisted field-to-column entries for a v2.0 data file.
52pub fn data_file_columns(schema: &Schema) -> (Vec<i32>, Vec<i32>) {
53    let mut field_ids = Vec::new();
54    let mut column_indices = Vec::new();
55    append_physical_fields(&schema.fields, &mut field_ids, &mut column_indices, &mut 0);
56    (field_ids, column_indices)
57}
58
59pub(super) fn field_id_to_column_index(schema: &Schema) -> BTreeMap<u32, u32> {
60    let (field_ids, column_indices) = data_file_columns(schema);
61    field_ids
62        .into_iter()
63        .zip(column_indices)
64        .map(|(field_id, column_index)| (field_id as u32, column_index as u32))
65        .collect()
66}
67
68fn append_physical_fields(
69    fields: &[Field],
70    field_ids: &mut Vec<i32>,
71    column_indices: &mut Vec<i32>,
72    next_column: &mut i32,
73) {
74    for field in fields {
75        field_ids.push(field.id);
76        column_indices.push(*next_column);
77        *next_column += 1;
78        if !field.is_blob() && !field.is_packed_struct() {
79            append_physical_fields(&field.children, field_ids, column_indices, next_column);
80        }
81    }
82}
83
84fn is_external_metadata_structural_header(
85    fields: &[Field],
86    target_column: usize,
87    next_column: &mut usize,
88) -> Option<bool> {
89    for field in fields {
90        if *next_column == target_column {
91            return Some(field.logical_type.is_struct() && !field.is_packed_struct());
92        }
93        *next_column += 1;
94        if !field.is_blob()
95            && !field.is_packed_struct()
96            && let Some(is_header) =
97                is_external_metadata_structural_header(&field.children, target_column, next_column)
98        {
99            return Some(is_header);
100        }
101    }
102    None
103}
104
105pub(super) fn should_copy_external_metadata_column(
106    schema: &Schema,
107    column_index: usize,
108    has_existing_pages: bool,
109) -> bool {
110    let mut next_column = 0;
111    let is_header =
112        is_external_metadata_structural_header(&schema.fields, column_index, &mut next_column)
113            .unwrap_or(false);
114    !is_header || !has_existing_pages
115}
116
117pub(super) fn finalize_external_metadata_column(
118    schema: &Schema,
119    column_index: usize,
120    pages: &mut Vec<PageInfo>,
121    num_rows: u64,
122) {
123    let mut next_column = 0;
124    let is_header =
125        is_external_metadata_structural_header(&schema.fields, column_index, &mut next_column)
126            .unwrap_or(false);
127    if is_header && !pages.is_empty() {
128        pages[0].num_rows = num_rows;
129        pages[0].priority = 0;
130        pages.truncate(1);
131    }
132}
133
134/// Compose the v2.0 field encoding mechanisms.
135pub fn encoding_strategy() -> Arc<dyn FieldEncodingStrategy> {
136    Arc::new(ArrayFieldEncodingStrategy::new())
137}
138
139/// Create a v2.0 writer with an explicit schema.
140pub fn create_writer(
141    object_writer: Box<dyn ObjectWriter>,
142    schema: Schema,
143    options: FileWriterOptions,
144) -> Result<Writer> {
145    Writer::try_new(object_writer, schema, options)
146}
147
148/// Create a v2.0 writer whose schema is inferred from the first batch.
149pub fn create_lazy_writer(
150    object_writer: Box<dyn ObjectWriter>,
151    options: FileWriterOptions,
152) -> Writer {
153    Writer::new_lazy(object_writer, options)
154}
155
156/// Encode a self-described v2.0 batch.
157pub fn encode_self_described_batch(batch: &EncodedBatch) -> Result<Bytes> {
158    writer::concat_lance_footer(batch, true)
159}
160
161/// Encode a mini-lance v2.0 batch.
162pub fn encode_mini_batch(batch: &EncodedBatch) -> Result<Bytes> {
163    writer::concat_lance_footer(batch, false)
164}