lance_file/versions/v2_1/
mod.rs1use 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
50pub fn physical_column_count(field: &Field) -> usize {
52 structural::physical_column_count(field)
53}
54
55pub 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
125pub 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
136pub 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
145pub 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
155pub fn create_lazy_writer(
157 object_writer: Box<dyn ObjectWriter>,
158 options: FileWriterOptions,
159) -> Writer {
160 Writer::new_lazy(object_writer, options)
161}
162
163pub 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
172pub fn encode_self_described_batch(batch: &EncodedBatch) -> Result<Bytes> {
174 writer::concat_lance_footer(batch, true)
175}
176
177pub fn encode_mini_batch(batch: &EncodedBatch) -> Result<Bytes> {
179 writer::concat_lance_footer(batch, false)
180}