Skip to main content

lance_file/versions/v2_1/
mod.rs

1// SPDX-License-Identifier: Apache-2.0
2// SPDX-FileCopyrightText: Copyright The Lance Authors
3
4//! Lance v2.1 file composition.
5
6use std::{collections::BTreeMap, sync::Arc};
7
8use bytes::Bytes;
9use lance_core::{
10    Error, Result,
11    datatypes::{Field, Schema},
12};
13use lance_encoding::{
14    compression_config::CompressionParams,
15    encoder::{
16        ColumnIndexSequence, EncodedBatch, FieldEncoder, FieldEncodingContext,
17        FieldEncodingStrategy,
18        structural::{
19            PrimitiveFieldEncoding, PrimitivePageEncoding, try_create_binary_blob, try_create_list,
20            try_create_struct,
21        },
22    },
23};
24use lance_io::traits::Writer as ObjectWriter;
25
26use crate::{
27    reader::{ReadProjection, structural},
28    writer::FileWriterOptions,
29};
30
31mod compression;
32mod reader;
33mod writer;
34
35pub use reader::{
36    projection_from_column_names, projection_from_field_ids, projection_from_whole_schema,
37};
38pub use writer::Writer;
39
40#[cfg(test)]
41pub(crate) use reader::test_projection_length;
42pub(crate) use reader::{
43    decode_column_metadata, finish_metadata, finish_metadata_index, validate_global_buffers,
44};
45
46pub(crate) fn read_projection() -> Arc<dyn ReadProjection> {
47    structural::read_projection(reader::decode_column)
48}
49
50/// Count physical columns represented by a field in a v2.1 footer.
51pub fn physical_column_count(field: &Field) -> usize {
52    structural::physical_column_count(field)
53}
54
55/// Build persisted field-to-column entries for a v2.1 data file.
56pub fn data_file_columns(schema: &Schema) -> (Vec<i32>, Vec<i32>) {
57    structural::data_file_columns(schema)
58}
59
60pub(super) fn field_id_to_column_index(schema: &Schema) -> BTreeMap<u32, u32> {
61    structural::field_id_to_column_index(schema)
62}
63
64#[derive(Debug)]
65struct FieldStrategy {
66    primitive: PrimitiveFieldEncoding,
67}
68
69impl FieldEncodingStrategy for FieldStrategy {
70    fn create_field_encoder(
71        &self,
72        field: &Field,
73        column_index: &mut ColumnIndexSequence,
74        context: &FieldEncodingContext<'_>,
75    ) -> Result<Box<dyn FieldEncoder>> {
76        if let Some(encoder) =
77            try_create_binary_blob(&self.primitive, field, column_index, context)?
78        {
79            return Ok(encoder);
80        }
81        if field.is_blob() {
82            return Err(Error::invalid_input_source(
83                format!(
84                    "Blob encoding is not available for field '{}' with data type {}",
85                    field.name,
86                    field.data_type()
87                )
88                .into(),
89            ));
90        }
91        if let Some(encoder) = self.primitive.try_create(field, column_index, context)? {
92            return Ok(encoder);
93        }
94        if matches!(
95            field.data_type(),
96            arrow_schema::DataType::FixedSizeList(item, _)
97                if matches!(item.data_type(), arrow_schema::DataType::Struct(_))
98        ) {
99            return Err(Error::not_supported_source(
100                "FixedSizeList<Struct> is not enabled by the selected file format".into(),
101            ));
102        }
103        if matches!(field.data_type(), arrow_schema::DataType::Map(_, _)) {
104            return Err(Error::not_supported_source(
105                "Map data type is not enabled by the selected file format".into(),
106            ));
107        }
108        if let Some(encoder) = try_create_list(field, column_index, context)? {
109            return Ok(encoder);
110        }
111        if let Some(encoder) = try_create_struct(field, column_index, context)? {
112            return Ok(encoder);
113        }
114        Err(Error::not_supported_source(
115            format!(
116                "Lance v2.1 has no field encoding for '{}' with data type {}",
117                field.name,
118                field.data_type()
119            )
120            .into(),
121        ))
122    }
123}
124
125/// Compose the v2.1 field encoding mechanisms.
126pub fn encoding_strategy(params: CompressionParams) -> Arc<dyn FieldEncodingStrategy> {
127    let compression = Arc::new(compression::Strategy::new(params));
128    Arc::new(FieldStrategy {
129        primitive: PrimitiveFieldEncoding::new([
130            PrimitivePageEncoding::reject_sparse(),
131            PrimitivePageEncoding::dense_u16(compression),
132        ]),
133    })
134}
135
136/// Create a v2.1 writer with an explicit schema.
137pub fn create_writer(
138    object_writer: Box<dyn ObjectWriter>,
139    schema: Schema,
140    options: FileWriterOptions,
141) -> Result<Writer> {
142    Writer::try_new(object_writer, schema, options)
143}
144
145/// Create a v2.1 writer with explicit compression tuning.
146pub fn create_writer_with_compression(
147    object_writer: Box<dyn ObjectWriter>,
148    schema: Schema,
149    options: FileWriterOptions,
150    compression: CompressionParams,
151) -> Result<Writer> {
152    Writer::try_new_with_compression(object_writer, schema, options, compression)
153}
154
155/// Create a v2.1 writer whose schema is inferred from the first batch.
156pub fn create_lazy_writer(
157    object_writer: Box<dyn ObjectWriter>,
158    options: FileWriterOptions,
159) -> Writer {
160    Writer::new_lazy(object_writer, options)
161}
162
163/// Create a lazy v2.1 writer with explicit compression tuning.
164pub fn create_lazy_writer_with_compression(
165    object_writer: Box<dyn ObjectWriter>,
166    options: FileWriterOptions,
167    compression: CompressionParams,
168) -> Writer {
169    Writer::new_lazy_with_compression(object_writer, options, compression)
170}
171
172/// Encode a self-described v2.1 batch.
173pub fn encode_self_described_batch(batch: &EncodedBatch) -> Result<Bytes> {
174    writer::concat_lance_footer(batch, true)
175}
176
177/// Encode a mini-lance v2.1 batch.
178pub fn encode_mini_batch(batch: &EncodedBatch) -> Result<Bytes> {
179    writer::concat_lance_footer(batch, false)
180}