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 encoding composition.
5
6use std::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::writer::FileWriterOptions;
27
28mod compression;
29mod writer;
30
31pub use writer::Writer;
32
33#[derive(Debug)]
34struct FieldStrategy {
35    primitive: PrimitiveFieldEncoding,
36}
37
38impl FieldEncodingStrategy for FieldStrategy {
39    fn create_field_encoder(
40        &self,
41        field: &Field,
42        column_index: &mut ColumnIndexSequence,
43        context: &FieldEncodingContext<'_>,
44    ) -> Result<Box<dyn FieldEncoder>> {
45        if let Some(encoder) =
46            try_create_binary_blob(&self.primitive, field, column_index, context)?
47        {
48            return Ok(encoder);
49        }
50        if field.is_blob() {
51            return Err(Error::invalid_input_source(
52                format!(
53                    "Blob encoding is not available for field '{}' with data type {}",
54                    field.name,
55                    field.data_type()
56                )
57                .into(),
58            ));
59        }
60        if let Some(encoder) = self.primitive.try_create(field, column_index, context)? {
61            return Ok(encoder);
62        }
63        if matches!(
64            field.data_type(),
65            arrow_schema::DataType::FixedSizeList(item, _)
66                if matches!(item.data_type(), arrow_schema::DataType::Struct(_))
67        ) {
68            return Err(Error::not_supported_source(
69                "FixedSizeList<Struct> is not enabled by the selected file format".into(),
70            ));
71        }
72        if matches!(field.data_type(), arrow_schema::DataType::Map(_, _)) {
73            return Err(Error::not_supported_source(
74                "Map data type is not enabled by the selected file format".into(),
75            ));
76        }
77        if let Some(encoder) = try_create_list(field, column_index, context)? {
78            return Ok(encoder);
79        }
80        if let Some(encoder) = try_create_struct(field, column_index, context)? {
81            return Ok(encoder);
82        }
83        Err(Error::not_supported_source(
84            format!(
85                "Lance v2.1 has no field encoding for '{}' with data type {}",
86                field.name,
87                field.data_type()
88            )
89            .into(),
90        ))
91    }
92}
93
94/// Compose the v2.1 field encoding mechanisms.
95pub fn encoding_strategy(params: CompressionParams) -> Arc<dyn FieldEncodingStrategy> {
96    let compression = Arc::new(compression::Strategy::new(params));
97    Arc::new(FieldStrategy {
98        primitive: PrimitiveFieldEncoding::new([
99            PrimitivePageEncoding::reject_sparse(),
100            PrimitivePageEncoding::dense_u16(compression),
101        ]),
102    })
103}
104
105/// Create a v2.1 writer with an explicit schema.
106pub fn create_writer(
107    object_writer: Box<dyn ObjectWriter>,
108    schema: Schema,
109    options: FileWriterOptions,
110) -> Result<Writer> {
111    Writer::try_new(object_writer, schema, options)
112}
113
114/// Create a v2.1 writer with explicit compression tuning.
115pub fn create_writer_with_compression(
116    object_writer: Box<dyn ObjectWriter>,
117    schema: Schema,
118    options: FileWriterOptions,
119    compression: CompressionParams,
120) -> Result<Writer> {
121    Writer::try_new_with_compression(object_writer, schema, options, compression)
122}
123
124/// Create a v2.1 writer whose schema is inferred from the first batch.
125pub fn create_lazy_writer(
126    object_writer: Box<dyn ObjectWriter>,
127    options: FileWriterOptions,
128) -> Writer {
129    Writer::new_lazy(object_writer, options)
130}
131
132/// Create a lazy v2.1 writer with explicit compression tuning.
133pub fn create_lazy_writer_with_compression(
134    object_writer: Box<dyn ObjectWriter>,
135    options: FileWriterOptions,
136    compression: CompressionParams,
137) -> Writer {
138    Writer::new_lazy_with_compression(object_writer, options, compression)
139}
140
141/// Encode a self-described v2.1 batch.
142pub fn encode_self_described_batch(batch: &EncodedBatch) -> Result<Bytes> {
143    writer::concat_lance_footer(batch, true)
144}
145
146/// Encode a mini-lance v2.1 batch.
147pub fn encode_mini_batch(batch: &EncodedBatch) -> Result<Bytes> {
148    writer::concat_lance_footer(batch, false)
149}