use crate::Dataset;
use crate::Result;
use crate::dataset::DATA_DIR;
use crate::dataset::WriteParams;
use crate::dataset::fragment::write::generate_random_filename;
use crate::datatypes::Schema;
use lance_core::Error;
use lance_encoding::decoder::{ColumnInfo, PageInfo as DecPageInfo};
use lance_file::reader::FileReader as LFReader;
use lance_file::version::ConcreteFileVersion;
use lance_file::versions as file_versions;
use lance_file::writer::{FileWriter, FileWriterOptions};
use lance_io::scheduler::{ScanScheduler, SchedulerConfig};
use lance_table::format::{DataFile, Fragment};
use prost::Message;
use prost_types::Any;
use std::ops::Range;
use std::sync::Arc;
async fn init_writer_if_necessary(
dataset: &Dataset,
version: ConcreteFileVersion,
current_writer: &mut Option<FileWriter>,
current_filename: &mut Option<String>,
) -> Result<bool> {
if current_writer.is_none() {
let filename = format!("{}.lance", generate_random_filename());
let path = dataset.base.clone().join(DATA_DIR).join(filename.as_str());
let object_writer = dataset.object_store.create(&path).await?;
*current_writer = Some(file_versions::create_lazy_writer(
version,
object_writer,
FileWriterOptions::default(),
)?);
*current_filename = Some(filename);
return Ok(true);
}
Ok(false)
}
#[allow(clippy::too_many_arguments)]
async fn finalize_current_output_file(
schema: &Schema,
version: ConcreteFileVersion,
current_writer: &mut Option<FileWriter>,
current_filename: &mut Option<String>,
current_page_table: &[ColumnInfo],
col_pages: &mut [Vec<DecPageInfo>],
col_buffers: &mut [Vec<(u64, u64)>],
total_rows_in_current: u64,
) -> Result<Fragment> {
let mut final_cols: Vec<Arc<ColumnInfo>> = Vec::with_capacity(current_page_table.len());
for (i, column_info) in current_page_table.iter().enumerate() {
let mut pages_vec = std::mem::take(&mut col_pages[i]);
file_versions::finalize_external_metadata_column(
version,
schema,
i,
&mut pages_vec,
total_rows_in_current,
)?;
let pages_arc = Arc::from(pages_vec.into_boxed_slice());
let buffers_vec = std::mem::take(&mut col_buffers[i]);
final_cols.push(Arc::new(ColumnInfo::new(
column_info.index,
pages_arc,
buffers_vec,
column_info.encoding.clone(),
)));
}
let mut writer = current_writer
.take()
.ok_or_else(|| Error::internal("binary copy output writer was not initialized"))?;
flush_footer(&mut writer, schema, &final_cols, total_rows_in_current).await?;
let mut fragment = Fragment::new(0);
let (field_ids, field_column_indices) = file_versions::data_file_columns(version, schema);
let filename = current_filename
.take()
.ok_or_else(|| Error::internal("binary copy output filename was not initialized"))?;
let mut data_file = DataFile::new_unstarted(filename, version);
data_file.fields = field_ids.into();
data_file.column_indices = field_column_indices.into();
fragment.files.push(data_file);
fragment.physical_rows = Some(total_rows_in_current as usize);
Ok(fragment)
}
pub async fn rewrite_files_binary_copy(
version: ConcreteFileVersion,
dataset: &Dataset,
fragments: &[Fragment],
params: &WriteParams,
read_batch_bytes_opt: Option<usize>,
) -> Result<Vec<Fragment>> {
if fragments.is_empty() || fragments.iter().any(|fragment| fragment.files.is_empty()) {
return Err(Error::invalid_input(
"binary copy requires at least one data file",
));
}
let schema = dataset.schema().clone();
let column_count = schema
.fields
.iter()
.map(|field| file_versions::physical_column_count(version, field))
.sum();
let mut out: Vec<Fragment> = Vec::new();
let mut current_writer: Option<FileWriter> = None;
let mut current_filename: Option<String> = None;
let mut current_page_table: Vec<ColumnInfo> = Vec::new();
let mut baseline_col_encoding_bytes: Vec<Vec<u8>> = Vec::new();
let mut col_pages: Vec<Vec<DecPageInfo>> = std::iter::repeat_with(Vec::<DecPageInfo>::new)
.take(column_count)
.collect();
let mut col_buffers: Vec<Vec<(u64, u64)>> = vec![Vec::new(); column_count];
let mut total_rows_in_current: u64 = 0;
let max_rows_per_file = params.max_rows_per_file as u64;
for frag in fragments.iter() {
for df in frag.files.iter() {
let object_store = if let Some(base_id) = df.base_id {
dataset.object_store(Some(base_id)).await?
} else {
dataset.object_store.clone()
};
let full_path = dataset.data_file_dir(df)?.clone().join(df.path.as_str());
let scan_scheduler = ScanScheduler::new(
object_store.clone(),
SchedulerConfig::max_bandwidth(&object_store),
);
let file_scheduler = scan_scheduler
.open_file_with_priority(&full_path, 0, &df.file_size_bytes)
.await?;
let file_meta = LFReader::read_all_metadata(&file_scheduler).await?;
let src_column_infos = file_meta.column_infos.clone();
if current_page_table.is_empty() {
current_page_table = src_column_infos
.iter()
.map(|column_index| ColumnInfo {
index: column_index.index,
buffer_offsets_and_sizes: Arc::from(
Vec::<(u64, u64)>::new().into_boxed_slice(),
),
page_infos: Arc::from(Vec::<DecPageInfo>::new().into_boxed_slice()),
encoding: column_index.encoding.clone(),
})
.collect();
baseline_col_encoding_bytes = src_column_infos
.iter()
.map(|ci| Ok(Any::from_msg(&ci.encoding)?.encode_to_vec()))
.collect::<Result<Vec<_>>>()?;
}
for (col_idx, src_column_info) in src_column_infos.iter().enumerate() {
let has_existing_pages = !col_pages[col_idx].is_empty();
file_versions::copy_external_metadata_column(
version,
&schema,
col_idx,
has_existing_pages,
|| async {
init_writer_if_necessary(
dataset,
version,
&mut current_writer,
&mut current_filename,
)
.await?;
let read_batch_bytes: u64 =
read_batch_bytes_opt.unwrap_or(16 * 1024 * 1024) as u64;
let mut page_index = 0;
while page_index < src_column_info.page_infos.len() {
let mut batch_ranges: Vec<Range<u64>> = Vec::new();
let mut batch_counts: Vec<usize> = Vec::new();
let mut batch_bytes: u64 = 0;
let mut batch_pages: usize = 0;
for current_page in &src_column_info.page_infos[page_index..] {
let page_bytes: u64 = current_page
.buffer_offsets_and_sizes
.iter()
.map(|(_, size)| *size)
.sum();
let would_exceed =
batch_pages > 0 && (batch_bytes + page_bytes > read_batch_bytes);
if would_exceed {
break;
}
batch_counts.push(current_page.buffer_offsets_and_sizes.len());
for (offset, size) in current_page.buffer_offsets_and_sizes.iter() {
if *size > 0 {
batch_ranges.push((*offset)..(*offset + *size));
}
}
batch_bytes += page_bytes;
batch_pages += 1;
page_index += 1;
}
let bytes_vec = if batch_ranges.is_empty() {
Vec::new()
} else {
file_scheduler.submit_request(batch_ranges, 0).await?
};
let mut bytes_iter = bytes_vec.into_iter();
for (local_idx, buffer_count) in batch_counts.iter().enumerate() {
let page_idx = page_index - batch_pages + local_idx;
let page = &src_column_info.page_infos[page_idx];
let mut new_offsets = Vec::with_capacity(*buffer_count);
for (buffer_idx, (_, size)) in
page.buffer_offsets_and_sizes.iter().enumerate()
{
let writer = current_writer.as_mut().ok_or_else(|| {
Error::internal("binary copy output writer was not initialized")
})?;
let bytes = if *size == 0 {
None
} else {
Some(bytes_iter.next().ok_or_else(|| {
Error::execution(format!(
"binary copy: missing page buffer bytes while rewriting data file \
(column {col_idx}, page {page_idx}, buffer {buffer_idx}, expected size {size})",
))
})?)
};
let (start, written) = writer
.write_external_buffer(bytes.as_deref().unwrap_or_default())
.await?;
new_offsets.push((start, written));
}
let new_page_info = DecPageInfo {
num_rows: page.num_rows,
priority: page.priority + total_rows_in_current,
encoding: page.encoding.clone(),
buffer_offsets_and_sizes: Arc::from(new_offsets.into_boxed_slice()),
};
col_pages[col_idx].push(new_page_info);
}
}
if !src_column_info.buffer_offsets_and_sizes.is_empty() {
let src_col_encoding_bytes =
Any::from_msg(&src_column_info.encoding)?.encode_to_vec();
let baseline_bytes = &baseline_col_encoding_bytes[col_idx];
if src_col_encoding_bytes != *baseline_bytes {
return Err(Error::execution(format!(
"binary copy: The ColumnEncoding of column {} is incompatible with the first file, \
making it impossible to safely concatenate buffers",
col_idx
)));
}
let ranges: Vec<Range<u64>> = src_column_info
.buffer_offsets_and_sizes
.iter()
.filter(|(_, size)| *size > 0)
.map(|(offset, size)| (*offset)..(*offset + *size))
.collect();
let bytes_vec = if ranges.is_empty() {
Vec::new()
} else {
file_scheduler.submit_request(ranges, 0).await?
};
let mut bytes_iter = bytes_vec.into_iter();
for (buffer_idx, (_, size)) in
src_column_info.buffer_offsets_and_sizes.iter().enumerate()
{
let writer = current_writer.as_mut().ok_or_else(|| {
Error::internal("binary copy output writer was not initialized")
})?;
let bytes = if *size == 0 {
None
} else {
Some(bytes_iter.next().ok_or_else(|| {
Error::execution(format!(
"binary copy: missing column buffer bytes while rewriting data file \
(column {col_idx}, buffer {buffer_idx}, expected size {size})",
))
})?)
};
let (start, written) = writer
.write_external_buffer(bytes.as_deref().unwrap_or_default())
.await?;
col_buffers[col_idx].push((start, written));
}
}
Ok(())
},
)
.await?;
}
total_rows_in_current += file_meta.num_rows;
if total_rows_in_current >= max_rows_per_file {
let fragment_out = finalize_current_output_file(
&schema,
version,
&mut current_writer,
&mut current_filename,
¤t_page_table,
&mut col_pages,
&mut col_buffers,
total_rows_in_current,
)
.await?;
current_writer = None;
current_page_table.clear();
for v in col_pages.iter_mut() {
v.clear();
}
for v in col_buffers.iter_mut() {
v.clear();
}
out.push(fragment_out);
total_rows_in_current = 0;
}
}
}
if total_rows_in_current > 0 {
init_writer_if_necessary(dataset, version, &mut current_writer, &mut current_filename)
.await?;
let frag = finalize_current_output_file(
&schema,
version,
&mut current_writer,
&mut current_filename,
¤t_page_table,
&mut col_pages,
&mut col_buffers,
total_rows_in_current,
)
.await?;
out.push(frag);
}
Ok(out)
}
async fn flush_footer(
writer: &mut FileWriter,
schema: &Schema,
final_cols: &[Arc<ColumnInfo>],
total_rows_in_current: u64,
) -> Result<()> {
writer.write_external_buffer(&[]).await?;
writer.initialize_with_external_columns(schema.clone(), final_cols, total_rows_in_current)?;
writer.finish().await?;
Ok(())
}