lance_file/versions/v2_0/
mod.rs1use 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
38pub 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
51pub 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
134pub fn encoding_strategy() -> Arc<dyn FieldEncodingStrategy> {
136 Arc::new(ArrayFieldEncodingStrategy::new())
137}
138
139pub 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
148pub fn create_lazy_writer(
150 object_writer: Box<dyn ObjectWriter>,
151 options: FileWriterOptions,
152) -> Writer {
153 Writer::new_lazy(object_writer, options)
154}
155
156pub fn encode_self_described_batch(batch: &EncodedBatch) -> Result<Bytes> {
158 writer::concat_lance_footer(batch, true)
159}
160
161pub fn encode_mini_batch(batch: &EncodedBatch) -> Result<Bytes> {
163 writer::concat_lance_footer(batch, false)
164}