lance_file/versions/v2_2/
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_map, try_create_struct, try_create_structural_blob,
21 try_create_structural_fixed_size_list,
22 },
23 },
24};
25use lance_io::traits::Writer as ObjectWriter;
26
27use crate::{
28 reader::{ReadProjection, structural},
29 writer::FileWriterOptions,
30};
31
32mod compression;
33mod reader;
34mod writer;
35
36pub(crate) use reader::{
37 decode_column_metadata, finish_metadata, finish_metadata_index, validate_global_buffers,
38};
39pub use reader::{
40 projection_from_column_names, projection_from_field_ids, projection_from_whole_schema,
41};
42
43pub(crate) fn read_projection() -> Arc<dyn ReadProjection> {
44 structural::read_projection(reader::decode_column)
45}
46pub use writer::Writer;
47
48pub fn physical_column_count(field: &Field) -> usize {
50 structural::physical_column_count(field)
51}
52
53pub fn data_file_columns(schema: &Schema) -> (Vec<i32>, Vec<i32>) {
55 structural::data_file_columns(schema)
56}
57
58pub(super) fn field_id_to_column_index(schema: &Schema) -> BTreeMap<u32, u32> {
59 structural::field_id_to_column_index(schema)
60}
61
62#[derive(Debug)]
63struct FieldStrategy {
64 primitive: PrimitiveFieldEncoding,
65}
66
67impl FieldEncodingStrategy for FieldStrategy {
68 fn create_field_encoder(
69 &self,
70 field: &Field,
71 column_index: &mut ColumnIndexSequence,
72 context: &FieldEncodingContext<'_>,
73 ) -> Result<Box<dyn FieldEncoder>> {
74 if let Some(encoder) =
75 try_create_binary_blob(&self.primitive, field, column_index, context)?
76 {
77 return Ok(encoder);
78 }
79 if let Some(encoder) =
80 try_create_structural_blob(&self.primitive, field, column_index, context)?
81 {
82 return Ok(encoder);
83 }
84 if field.is_blob() {
85 return Err(Error::invalid_input_source(
86 format!(
87 "Blob encoding is not available for field '{}' with data type {}",
88 field.name,
89 field.data_type()
90 )
91 .into(),
92 ));
93 }
94 if let Some(encoder) = try_create_map(field, column_index, context)? {
95 return Ok(encoder);
96 }
97 if let Some(encoder) = try_create_structural_fixed_size_list(field, column_index, context)?
98 {
99 return Ok(encoder);
100 }
101 if let Some(encoder) = self.primitive.try_create(field, column_index, context)? {
102 return Ok(encoder);
103 }
104 if let Some(encoder) = try_create_list(field, column_index, context)? {
105 return Ok(encoder);
106 }
107 if let Some(encoder) = try_create_struct(field, column_index, context)? {
108 return Ok(encoder);
109 }
110 Err(Error::not_supported_source(
111 format!(
112 "Lance v2.2 has no field encoding for '{}' with data type {}",
113 field.name,
114 field.data_type()
115 )
116 .into(),
117 ))
118 }
119}
120
121pub fn encoding_strategy(params: CompressionParams) -> Arc<dyn FieldEncodingStrategy> {
123 let compression = Arc::new(compression::Strategy::new(params));
124 Arc::new(FieldStrategy {
125 primitive: PrimitiveFieldEncoding::new([
126 PrimitivePageEncoding::reject_sparse(),
127 PrimitivePageEncoding::constant(),
128 PrimitivePageEncoding::dense_u32(compression),
129 ]),
130 })
131}
132
133pub fn create_writer(
135 object_writer: Box<dyn ObjectWriter>,
136 schema: Schema,
137 options: FileWriterOptions,
138) -> Result<Writer> {
139 Writer::try_new(object_writer, schema, options)
140}
141
142pub fn create_writer_with_compression(
144 object_writer: Box<dyn ObjectWriter>,
145 schema: Schema,
146 options: FileWriterOptions,
147 compression: CompressionParams,
148) -> Result<Writer> {
149 Writer::try_new_with_compression(object_writer, schema, options, compression)
150}
151
152pub fn create_lazy_writer(
154 object_writer: Box<dyn ObjectWriter>,
155 options: FileWriterOptions,
156) -> Writer {
157 Writer::new_lazy(object_writer, options)
158}
159
160pub fn create_lazy_writer_with_compression(
162 object_writer: Box<dyn ObjectWriter>,
163 options: FileWriterOptions,
164 compression: CompressionParams,
165) -> Writer {
166 Writer::new_lazy_with_compression(object_writer, options, compression)
167}
168
169pub fn encode_self_described_batch(batch: &EncodedBatch) -> Result<Bytes> {
171 writer::concat_lance_footer(batch, true)
172}
173
174pub fn encode_mini_batch(batch: &EncodedBatch) -> Result<Bytes> {
176 writer::concat_lance_footer(batch, false)
177}