Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
63 changes: 46 additions & 17 deletions rust/lance-file/src/compatibility_tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -24,11 +24,12 @@ use tokio::io::AsyncWriteExt;
use crate::reader::{FileReader, FileReaderOptions};
use crate::testing::FsFixture;
use crate::version::ConcreteFileVersion;
use crate::versions;
use crate::versions::v1::reader::FileReader as V1Reader;
use crate::versions::v1::writer::{
FileWriter as V1Writer, FileWriterOptions as V1WriterOptions, NotSelfDescribing,
};
use crate::writer::{FileWriter, FileWriterOptions};
use crate::writer::FileWriterOptions;

fn compatibility_fixture_batch() -> RecordBatch {
let row_count = 4097;
Expand Down Expand Up @@ -158,22 +159,50 @@ async fn write_current_fixture(
) -> Vec<u8> {
let fs = FsFixture::default();
let object_writer = fs.object_store.create(&fs.tmp_path).await.unwrap();
let mut writer = FileWriter::try_new(
object_writer,
schema.clone(),
FileWriterOptions {
data_cache_bytes: Some(1),
max_page_bytes: Some(1024),
format_version: Some(version.into()),
..Default::default()
},
)
.unwrap();
for offset in (0..batch.num_rows()).step_by(1024) {
let slice = batch.slice(offset, (batch.num_rows() - offset).min(1024));
writer.write_batch(&slice).await.unwrap();
}
let summary = writer.finish().await.unwrap();
let options = FileWriterOptions {
data_cache_bytes: Some(1),
max_page_bytes: Some(1024),
..Default::default()
};
let summary = match version {
ConcreteFileVersion::V1 => unreachable!("v1 uses its manifest-backed writer"),
ConcreteFileVersion::V2_0 => {
let mut writer =
versions::v2_0::create_writer(object_writer, schema.clone(), options).unwrap();
for offset in (0..batch.num_rows()).step_by(1024) {
let slice = batch.slice(offset, (batch.num_rows() - offset).min(1024));
writer.write_batch(&slice).await.unwrap();
}
writer.finish().await.unwrap()
}
ConcreteFileVersion::V2_1 => {
let mut writer =
versions::v2_1::create_writer(object_writer, schema.clone(), options).unwrap();
for offset in (0..batch.num_rows()).step_by(1024) {
let slice = batch.slice(offset, (batch.num_rows() - offset).min(1024));
writer.write_batch(&slice).await.unwrap();
}
writer.finish().await.unwrap()
}
ConcreteFileVersion::V2_2 => {
let mut writer =
versions::v2_2::create_writer(object_writer, schema.clone(), options).unwrap();
for offset in (0..batch.num_rows()).step_by(1024) {
let slice = batch.slice(offset, (batch.num_rows() - offset).min(1024));
writer.write_batch(&slice).await.unwrap();
}
writer.finish().await.unwrap()
}
ConcreteFileVersion::V2_3 => {
let mut writer =
versions::v2_3::create_writer(object_writer, schema.clone(), options).unwrap();
for offset in (0..batch.num_rows()).step_by(1024) {
let slice = batch.slice(offset, (batch.num_rows() - offset).min(1024));
writer.write_batch(&slice).await.unwrap();
}
writer.finish().await.unwrap()
}
};
fs.object_store
.open(&fs.tmp_path)
.await
Expand Down
41 changes: 40 additions & 1 deletion rust/lance-file/src/versions/v2_0/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -5,9 +5,48 @@

use std::sync::Arc;

use lance_encoding::{array_encoding::ArrayFieldEncodingStrategy, encoder::FieldEncodingStrategy};
use bytes::Bytes;
use lance_core::{Result, datatypes::Schema};
use lance_encoding::{
array_encoding::ArrayFieldEncodingStrategy,
encoder::{EncodedBatch, FieldEncodingStrategy},
};
use lance_io::traits::Writer as ObjectWriter;

use crate::writer::FileWriterOptions;

mod writer;

pub use writer::Writer;

/// Compose the v2.0 field encoding mechanisms.
pub fn encoding_strategy() -> Arc<dyn FieldEncodingStrategy> {
Arc::new(ArrayFieldEncodingStrategy::new())
}

/// Create a v2.0 writer with an explicit schema.
pub fn create_writer(
object_writer: Box<dyn ObjectWriter>,
schema: Schema,
options: FileWriterOptions,
) -> Result<Writer> {
Writer::try_new(object_writer, schema, options)
}

/// Create a v2.0 writer whose schema is inferred from the first batch.
pub fn create_lazy_writer(
object_writer: Box<dyn ObjectWriter>,
options: FileWriterOptions,
) -> Writer {
Writer::new_lazy(object_writer, options)
}

/// Encode a self-described v2.0 batch.
pub fn encode_self_described_batch(batch: &EncodedBatch) -> Result<Bytes> {
writer::concat_lance_footer(batch, true)
}

/// Encode a mini-lance v2.0 batch.
pub fn encode_mini_batch(batch: &EncodedBatch) -> Result<Bytes> {
writer::concat_lance_footer(batch, false)
}
Loading
Loading