use std::fs::File;
use std::io::{BufReader, Read, Seek, SeekFrom, Write};
use std::path::{Path, PathBuf};
use std::sync::Arc;
use crate::batch::{Batch, OwnedPhysicalRow, PhysicalRow, RowSchema};
use crate::physical::ExecResult;
use tempfile::NamedTempFile;
mod format;
use format::{
append_batches, decode_batch, decode_physical_row_record, encode_physical_row_record,
encoded_batch_overhead_size, encoded_batch_size, encoded_physical_row_record_size,
open_spill_reader, read_bounded_spill_record, spill_error, RECORD_PREFIX_BYTES,
};
const SPILL_MAGIC: &[u8] = b"UQA-SPILL\x01\n";
#[derive(Clone, Copy)]
pub(crate) struct EncodedBatchSizer {
physical_width: usize,
bytes: usize,
origin_free_rows: usize,
has_lock_origins: bool,
}
impl EncodedBatchSizer {
pub(crate) fn new(schema: &RowSchema) -> ExecResult<Self> {
Ok(Self {
physical_width: schema.physical_width(),
bytes: encoded_batch_overhead_size(schema)?,
origin_free_rows: 0,
has_lock_origins: false,
})
}
pub(crate) fn append(&mut self, row: &PhysicalRow) -> ExecResult<()> {
let mut additional = encoded_physical_row_record_size(row, self.physical_width)?;
if row.lock_origins().is_empty() {
self.origin_free_rows = self.origin_free_rows.checked_add(1).ok_or_else(|| {
spill_error("incremental spill batch origin-free row count overflow")
})?;
if self.has_lock_origins {
additional = additional.checked_add(8).ok_or_else(|| {
spill_error("incremental spill batch lock-origin size overflow")
})?;
}
} else if !self.has_lock_origins {
let preceding_metadata = self.origin_free_rows.checked_mul(8).ok_or_else(|| {
spill_error("incremental spill batch lock-origin metadata overflow")
})?;
additional = additional
.checked_add(preceding_metadata)
.ok_or_else(|| spill_error("incremental spill batch lock-origin size overflow"))?;
self.has_lock_origins = true;
}
self.bytes = self
.bytes
.checked_add(additional)
.ok_or_else(|| spill_error("incremental spill batch size overflow"))?;
Ok(())
}
pub(crate) fn bytes(self) -> usize {
self.bytes
}
}
pub struct SpillBuffer {
schema: Option<RowSchema>,
batches: Vec<Batch>,
rows: usize,
in_memory_rows: usize,
in_memory_bytes: usize,
max_in_memory_record_bytes: usize,
budget_bytes: usize,
spill_directory: Option<PathBuf>,
spill_file: Option<NamedTempFile>,
spilled_batches: usize,
spilled_rows: usize,
spilled_bytes: usize,
max_spilled_record_bytes: usize,
}
impl SpillBuffer {
pub fn new(budget_bytes: usize) -> Self {
Self {
schema: None,
batches: Vec::new(),
rows: 0,
in_memory_rows: 0,
in_memory_bytes: 0,
max_in_memory_record_bytes: 0,
budget_bytes,
spill_directory: None,
spill_file: None,
spilled_batches: 0,
spilled_rows: 0,
spilled_bytes: 0,
max_spilled_record_bytes: 0,
}
}
pub fn new_in(budget_bytes: usize, directory: impl Into<PathBuf>) -> Self {
let mut buffer = Self::new(budget_bytes);
buffer.spill_directory = Some(directory.into());
buffer
}
pub fn unbounded() -> Self {
Self::new(usize::MAX)
}
pub fn push(&mut self, batch: Batch) -> ExecResult<bool> {
if let Some(schema) = self.schema.as_ref() {
if schema != &batch.schema {
return Err(spill_error(format!(
"spill buffer schema mismatch: expected {:?}, got {:?}",
schema.columns(),
batch.schema.columns()
)));
}
} else {
self.schema = Some(batch.schema.clone());
}
let batch_rows = batch.rows.len();
let next_rows = self
.rows
.checked_add(batch_rows)
.ok_or_else(|| spill_error("spill buffer row count overflow"))?;
let batch_bytes = match Self::encoded_size(&batch) {
Ok(bytes) => bytes,
Err(error) => {
self.retain_batch(batch, usize::MAX);
return Err(error);
}
};
let would_exceed = self
.in_memory_bytes
.checked_add(batch_bytes)
.is_none_or(|bytes| bytes > self.budget_bytes);
let mut spilled = false;
if would_exceed && !self.batches.is_empty() {
if let Err(error) = self.spill_pending() {
self.retain_batch(batch, batch_bytes);
return Err(error);
}
spilled = true;
}
let next_in_memory_rows = self
.in_memory_rows
.checked_add(batch_rows)
.ok_or_else(|| spill_error("spill buffer in-memory row count overflow"))?;
let next_in_memory_bytes = self
.in_memory_bytes
.checked_add(batch_bytes)
.ok_or_else(|| spill_error("spill buffer in-memory byte count overflow"))?;
self.rows = next_rows;
self.in_memory_rows = next_in_memory_rows;
self.in_memory_bytes = next_in_memory_bytes;
self.max_in_memory_record_bytes = self.max_in_memory_record_bytes.max(batch_bytes);
self.batches.push(batch);
if self.in_memory_bytes > self.budget_bytes {
self.spill_pending()?;
spilled = true;
}
Ok(spilled)
}
pub fn encoded_size(batch: &Batch) -> ExecResult<usize> {
encoded_batch_size(batch)
}
pub fn rows(&self) -> usize {
self.rows
}
pub fn in_memory_rows(&self) -> usize {
self.in_memory_rows
}
pub fn in_memory_bytes(&self) -> usize {
self.in_memory_bytes
}
pub fn budget_bytes(&self) -> usize {
self.budget_bytes
}
pub fn over_budget(&self) -> bool {
self.in_memory_bytes > self.budget_bytes
}
pub fn has_spilled(&self) -> bool {
self.spill_file.is_some()
}
pub fn spilled_rows(&self) -> usize {
self.spilled_rows
}
pub fn spilled_batches(&self) -> usize {
self.spilled_batches
}
pub fn spilled_bytes(&self) -> usize {
self.spilled_bytes
}
pub fn spill_path(&self) -> Option<&Path> {
self.spill_file.as_ref().map(NamedTempFile::path)
}
pub fn spill_if_over_budget(&mut self) -> ExecResult<bool> {
if !self.over_budget() || self.batches.is_empty() {
return Ok(false);
}
self.spill_pending()
}
pub fn spill_pending(&mut self) -> ExecResult<bool> {
if self.batches.is_empty() {
return Ok(false);
}
let next_spilled_batches = self
.spilled_batches
.checked_add(self.batches.len())
.ok_or_else(|| spill_error("spill batch count overflow"))?;
let next_spilled_rows = self
.spilled_rows
.checked_add(self.in_memory_rows)
.ok_or_else(|| spill_error("spill row count overflow"))?;
let next_spilled_bytes = self
.spilled_bytes
.checked_add(self.in_memory_bytes)
.ok_or_else(|| spill_error("spill byte count overflow"))?;
let next_max_spilled_record_bytes = self
.max_spilled_record_bytes
.max(self.max_in_memory_record_bytes);
if let Some(file) = self.spill_file.as_mut() {
append_batches(file.as_file_mut(), &self.batches)?;
} else {
let mut file = self.create_spill_file()?;
append_batches(file.as_file_mut(), &self.batches)?;
self.spill_file = Some(file);
}
self.spilled_batches = next_spilled_batches;
self.spilled_rows = next_spilled_rows;
self.spilled_bytes = next_spilled_bytes;
self.max_spilled_record_bytes = next_max_spilled_record_bytes;
self.batches.clear();
self.in_memory_rows = 0;
self.in_memory_bytes = 0;
self.max_in_memory_record_bytes = 0;
Ok(true)
}
pub fn reader(&self) -> ExecResult<SpillReader<'_>> {
let reader = self
.spill_file
.as_ref()
.map(open_spill_reader)
.transpose()?;
let disk_finished = reader.is_none();
Ok(SpillReader {
reader,
memory: self.batches.iter(),
disk_finished,
failed: false,
max_record_bytes: self.max_spilled_record_bytes,
expected_schema: self.schema.clone(),
})
}
pub fn read_rows(&self) -> ExecResult<SpillRows<SpillReader<'_>>> {
self.reader().map(SpillRows::new)
}
pub fn drain(&mut self) -> ExecResult<SpillDrain> {
let reader = self
.spill_file
.as_ref()
.map(open_spill_reader)
.transpose()?;
let spill_file = self.spill_file.take();
let memory = std::mem::take(&mut self.batches).into_iter();
let expected_schema = self.schema.take();
self.rows = 0;
self.in_memory_rows = 0;
self.in_memory_bytes = 0;
self.max_in_memory_record_bytes = 0;
self.spilled_batches = 0;
self.spilled_rows = 0;
self.spilled_bytes = 0;
let max_record_bytes = std::mem::take(&mut self.max_spilled_record_bytes);
let disk_finished = reader.is_none();
Ok(SpillDrain {
reader,
spill_file,
memory,
disk_finished,
failed: false,
max_record_bytes,
expected_schema,
})
}
pub fn drain_all(&mut self) -> ExecResult<Vec<Batch>> {
self.drain()?.collect()
}
pub fn drain_rows(&mut self) -> ExecResult<SpillRows<SpillDrain>> {
self.drain().map(SpillRows::new)
}
pub fn clear(&mut self) {
self.schema = None;
self.batches.clear();
self.spill_file = None;
self.rows = 0;
self.in_memory_rows = 0;
self.in_memory_bytes = 0;
self.max_in_memory_record_bytes = 0;
self.spilled_batches = 0;
self.spilled_rows = 0;
self.spilled_bytes = 0;
self.max_spilled_record_bytes = 0;
}
pub fn into_shared(mut self, schema: impl Into<RowSchema>) -> ExecResult<SharedSpill> {
let schema = schema.into();
if let Some(actual) = self.schema.as_ref() {
if actual != &schema {
return Err(spill_error(format!(
"shared spill schema mismatch: expected {:?}, got {:?}",
schema.columns(),
actual.columns()
)));
}
}
let rows = self.rows;
let storage = if self.spill_file.is_none() {
let batches = std::mem::take(&mut self.batches);
SharedSpillStorage::Memory(batches)
} else {
self.spill_pending()?;
SharedSpillStorage::Disk(
self.spill_file
.take()
.expect("spill file exists after flushing shared materialization"),
)
};
Ok(SharedSpill {
inner: Arc::new(SharedSpillInner {
storage,
schema,
rows,
max_record_bytes: self.max_spilled_record_bytes,
}),
})
}
fn create_spill_file(&self) -> ExecResult<NamedTempFile> {
let mut file = match &self.spill_directory {
Some(directory) => NamedTempFile::new_in(directory).map_err(|error| {
spill_error(format!(
"failed to create spill file in {}: {error}",
directory.display()
))
})?,
None => NamedTempFile::new()
.map_err(|error| spill_error(format!("failed to create spill file: {error}")))?,
};
file.as_file_mut()
.write_all(SPILL_MAGIC)
.map_err(|error| spill_error(format!("failed to initialize spill file: {error}")))?;
file.as_file_mut()
.flush()
.map_err(|error| spill_error(format!("failed to flush spill header: {error}")))?;
Ok(file)
}
fn retain_batch(&mut self, batch: Batch, encoded_bytes: usize) {
self.rows = self.rows.saturating_add(batch.rows.len());
self.in_memory_rows = self.in_memory_rows.saturating_add(batch.rows.len());
self.in_memory_bytes = self.in_memory_bytes.saturating_add(encoded_bytes);
self.max_in_memory_record_bytes = self.max_in_memory_record_bytes.max(encoded_bytes);
self.batches.push(batch);
}
}
enum SharedSpillStorage {
Memory(Vec<Batch>),
Disk(NamedTempFile),
}
struct SharedSpillInner {
storage: SharedSpillStorage,
schema: RowSchema,
rows: usize,
max_record_bytes: usize,
}
#[derive(Clone)]
pub struct SharedSpill {
inner: Arc<SharedSpillInner>,
}
impl SharedSpill {
pub fn schema(&self) -> &[String] {
self.inner.schema.columns()
}
pub fn row_schema(&self) -> &RowSchema {
&self.inner.schema
}
pub fn rows(&self) -> usize {
self.inner.rows
}
pub fn has_spilled(&self) -> bool {
matches!(self.inner.storage, SharedSpillStorage::Disk(_))
}
pub fn reader(&self) -> ExecResult<SharedSpillReader> {
let source = Arc::clone(&self.inner);
Self::reader_from_source(source)
}
pub fn into_reader(self) -> ExecResult<SharedSpillReader> {
match Arc::try_unwrap(self.inner) {
Ok(SharedSpillInner {
storage: SharedSpillStorage::Memory(batches),
schema,
max_record_bytes,
..
}) => Ok(SharedSpillReader {
reader: SharedSpillReaderSource::OwnedMemory(batches.into_iter()),
source: None,
failed: false,
max_record_bytes,
expected_schema: Some(schema),
}),
Ok(inner) => Self::reader_from_source(Arc::new(inner)),
Err(source) => Self::reader_from_source(source),
}
}
fn reader_from_source(source: Arc<SharedSpillInner>) -> ExecResult<SharedSpillReader> {
let reader = match &source.storage {
SharedSpillStorage::Memory(_) => SharedSpillReaderSource::Memory { next_batch: 0 },
SharedSpillStorage::Disk(file) => {
SharedSpillReaderSource::Disk(open_spill_reader(file)?)
}
};
let max_record_bytes = source.max_record_bytes;
let expected_schema = Some(source.schema.clone());
Ok(SharedSpillReader {
reader,
source: Some(source),
failed: false,
max_record_bytes,
expected_schema,
})
}
pub fn read_rows(&self) -> ExecResult<SpillRows<SharedSpillReader>> {
self.reader().map(SpillRows::new)
}
}
enum SharedSpillReaderSource {
Memory { next_batch: usize },
OwnedMemory(std::vec::IntoIter<Batch>),
Disk(BufReader<File>),
}
fn validate_decoded_schema(batch: Batch, expected_schema: Option<&RowSchema>) -> ExecResult<Batch> {
if expected_schema.is_none_or(|expected| expected == &batch.schema) {
return Ok(batch);
}
let expected = expected_schema.expect("schema presence checked above");
Err(spill_error(format!(
"spill batch schema mismatch: expected {:?}, got {:?}",
expected.columns(),
batch.schema.columns()
)))
}
pub struct SharedSpillReader {
reader: SharedSpillReaderSource,
source: Option<Arc<SharedSpillInner>>,
failed: bool,
max_record_bytes: usize,
expected_schema: Option<RowSchema>,
}
impl Iterator for SharedSpillReader {
type Item = ExecResult<Batch>;
fn next(&mut self) -> Option<Self::Item> {
if self.failed {
return None;
}
match &mut self.reader {
SharedSpillReaderSource::Memory { next_batch } => {
let source = self
.source
.as_ref()
.expect("shared memory reader retains its source");
let SharedSpillStorage::Memory(batches) = &source.storage else {
unreachable!("shared materialization reader/storage mismatch")
};
let batch = batches.get(*next_batch)?.clone();
*next_batch += 1;
Some(validate_decoded_schema(
batch,
self.expected_schema.as_ref(),
))
}
SharedSpillReaderSource::OwnedMemory(batches) => batches
.next()
.map(|batch| validate_decoded_schema(batch, self.expected_schema.as_ref())),
SharedSpillReaderSource::Disk(reader) => {
match read_bounded_spill_record(reader, self.max_record_bytes, "shared spill batch")
{
Ok(None) => None,
Ok(Some(record)) => {
let decoded = decode_batch(&record).and_then(|batch| {
validate_decoded_schema(batch, self.expected_schema.as_ref())
});
if decoded.is_err() {
self.failed = true;
}
Some(decoded)
}
Err(error) => {
self.failed = true;
Some(Err(spill_error(format!(
"failed to read shared spill batch: {error}"
))))
}
}
}
}
}
}
pub struct SpillDrain {
reader: Option<BufReader<File>>,
spill_file: Option<NamedTempFile>,
memory: std::vec::IntoIter<Batch>,
disk_finished: bool,
failed: bool,
max_record_bytes: usize,
expected_schema: Option<RowSchema>,
}
impl Iterator for SpillDrain {
type Item = ExecResult<Batch>;
fn next(&mut self) -> Option<Self::Item> {
if self.failed {
return None;
}
if !self.disk_finished {
let Some(reader) = self.reader.as_mut() else {
self.failed = true;
self.disk_finished = true;
return Some(Err(spill_error(
"spill drain entered disk phase without a reader",
)));
};
match read_bounded_spill_record(reader, self.max_record_bytes, "spill batch") {
Ok(None) => {
self.disk_finished = true;
self.reader = None;
self.spill_file = None;
}
Ok(Some(record)) => {
let decoded = decode_batch(&record).and_then(|batch| {
validate_decoded_schema(batch, self.expected_schema.as_ref())
});
if decoded.is_err() {
self.failed = true;
}
return Some(decoded);
}
Err(error) => {
self.failed = true;
return Some(Err(spill_error(format!(
"failed to read spill batch: {error}"
))));
}
}
}
self.memory
.next()
.map(|batch| validate_decoded_schema(batch, self.expected_schema.as_ref()))
}
fn size_hint(&self) -> (usize, Option<usize>) {
let lower = if self.disk_finished {
self.memory.len()
} else {
0
};
(lower, None)
}
}
pub struct SpillReader<'a> {
reader: Option<BufReader<File>>,
memory: std::slice::Iter<'a, Batch>,
disk_finished: bool,
failed: bool,
max_record_bytes: usize,
expected_schema: Option<RowSchema>,
}
impl Iterator for SpillReader<'_> {
type Item = ExecResult<Batch>;
fn next(&mut self) -> Option<Self::Item> {
if self.failed {
return None;
}
if !self.disk_finished {
let Some(reader) = self.reader.as_mut() else {
self.failed = true;
self.disk_finished = true;
return Some(Err(spill_error(
"spill reader entered disk phase without a file reader",
)));
};
match read_bounded_spill_record(reader, self.max_record_bytes, "spill batch") {
Ok(None) => {
self.disk_finished = true;
self.reader = None;
}
Ok(Some(record)) => {
let decoded = decode_batch(&record).and_then(|batch| {
validate_decoded_schema(batch, self.expected_schema.as_ref())
});
if decoded.is_err() {
self.failed = true;
}
return Some(decoded);
}
Err(error) => {
self.failed = true;
return Some(Err(spill_error(format!(
"failed to read spill batch: {error}"
))));
}
}
}
self.memory
.next()
.cloned()
.map(|batch| validate_decoded_schema(batch, self.expected_schema.as_ref()))
}
fn size_hint(&self) -> (usize, Option<usize>) {
let lower = if self.disk_finished {
self.memory.len()
} else {
0
};
(lower, None)
}
}
pub struct SpillRows<I> {
batches: I,
current_schema: Option<RowSchema>,
current: std::vec::IntoIter<PhysicalRow>,
}
pub struct IndexedSpill {
schema: RowSchema,
data: NamedTempFile,
offsets: NamedTempFile,
rows: u64,
encoded_bytes: u64,
}
impl IndexedSpill {
pub fn new(input_schema: RowSchema) -> ExecResult<Self> {
Ok(Self {
schema: input_schema,
data: NamedTempFile::new().map_err(|error| {
spill_error(format!("failed to create indexed spill data: {error}"))
})?,
offsets: NamedTempFile::new().map_err(|error| {
spill_error(format!("failed to create indexed spill offsets: {error}"))
})?,
rows: 0,
encoded_bytes: 0,
})
}
pub fn len(&self) -> u64 {
self.rows
}
pub fn is_empty(&self) -> bool {
self.rows == 0
}
pub fn encoded_bytes(&self) -> u64 {
self.encoded_bytes
}
pub fn row_schema(&self) -> &RowSchema {
&self.schema
}
pub(crate) fn encoded_row_size(schema: &RowSchema, row: &PhysicalRow) -> ExecResult<usize> {
encoded_physical_row_record_size(row, schema.physical_width())?
.checked_add(RECORD_PREFIX_BYTES)
.ok_or_else(|| spill_error("indexed spill row size overflow"))
}
pub fn push(&mut self, row: &PhysicalRow) -> ExecResult<()> {
let payload = encode_physical_row_record(row, self.schema.physical_width())?;
let length = u64::try_from(payload.len())
.map_err(|_| spill_error("indexed spill row is too large"))?;
let next_rows = self
.rows
.checked_add(1)
.ok_or_else(|| spill_error("indexed spill row count overflow"))?;
let record_bytes = length
.checked_add(8)
.ok_or_else(|| spill_error("indexed spill row length overflow"))?;
let next_encoded_bytes = self
.encoded_bytes
.checked_add(record_bytes)
.ok_or_else(|| spill_error("indexed spill byte count overflow"))?;
let data_length = self
.data
.as_file_mut()
.seek(SeekFrom::End(0))
.map_err(|error| spill_error(format!("failed to seek indexed spill data: {error}")))?;
let offsets_length =
self.offsets
.as_file_mut()
.seek(SeekFrom::End(0))
.map_err(|error| {
spill_error(format!("failed to seek indexed spill offsets: {error}"))
})?;
let write_result = (|| -> std::io::Result<()> {
self.data.as_file_mut().write_all(&length.to_le_bytes())?;
self.data.as_file_mut().write_all(&payload)?;
self.data.as_file_mut().flush()?;
self.offsets
.as_file_mut()
.write_all(&data_length.to_le_bytes())?;
self.offsets.as_file_mut().flush()
})();
if let Err(error) = write_result {
let data_rollback = self.data.as_file_mut().set_len(data_length);
let offsets_rollback = self.offsets.as_file_mut().set_len(offsets_length);
let rollback_error = match (data_rollback, offsets_rollback) {
(Ok(()), Ok(())) => None,
(Err(data), Ok(())) => Some(format!("data rollback failed: {data}")),
(Ok(()), Err(offsets)) => Some(format!("offset rollback failed: {offsets}")),
(Err(data), Err(offsets)) => Some(format!(
"data rollback failed: {data}; offset rollback failed: {offsets}"
)),
};
if let Some(rollback) = rollback_error {
return Err(spill_error(format!(
"failed to append indexed spill row: {error}; {rollback}"
)));
}
return Err(spill_error(format!(
"failed to append indexed spill row: {error}"
)));
}
self.rows = next_rows;
self.encoded_bytes = next_encoded_bytes;
Ok(())
}
pub fn get(&mut self, index: u64) -> ExecResult<PhysicalRow> {
if index >= self.rows {
return Err(spill_error(format!(
"indexed spill row {index} is outside 0..{}",
self.rows
)));
}
let expected_offsets_length = self
.rows
.checked_mul(8)
.ok_or_else(|| spill_error("indexed spill offsets length overflow"))?;
let actual_offsets_length = self
.offsets
.as_file()
.metadata()
.map_err(|error| {
spill_error(format!("failed to inspect indexed spill offsets: {error}"))
})?
.len();
if actual_offsets_length != expected_offsets_length {
return Err(spill_error(format!(
"indexed spill offsets length {actual_offsets_length} does not match expected {expected_offsets_length}"
)));
}
let data_length = self
.data
.as_file()
.metadata()
.map_err(|error| spill_error(format!("failed to inspect indexed spill data: {error}")))?
.len();
let offset_position = index
.checked_mul(8)
.ok_or_else(|| spill_error("indexed spill offset overflow"))?;
let offset = read_indexed_offset(self.offsets.as_file_mut(), offset_position)?;
let record_end = if index
.checked_add(1)
.ok_or_else(|| spill_error("indexed spill row index overflow"))?
< self.rows
{
read_indexed_offset(
self.offsets.as_file_mut(),
offset_position
.checked_add(8)
.ok_or_else(|| spill_error("indexed spill next offset overflow"))?,
)?
} else {
data_length
};
let payload_start = offset
.checked_add(8)
.ok_or_else(|| spill_error("indexed spill payload offset overflow"))?;
if payload_start > record_end || record_end > data_length {
return Err(spill_error(format!(
"indexed spill record bounds {offset}..{record_end} are outside data length {data_length}"
)));
}
self.data
.as_file_mut()
.seek(SeekFrom::Start(offset))
.map_err(|error| spill_error(format!("failed to seek indexed spill row: {error}")))?;
let mut length = [0_u8; 8];
self.data
.as_file_mut()
.read_exact(&mut length)
.map_err(|error| {
spill_error(format!("failed to read indexed spill length: {error}"))
})?;
let declared_length = u64::from_le_bytes(length);
let available_length = record_end - payload_start;
if declared_length != available_length {
return Err(spill_error(format!(
"indexed spill row length {declared_length} does not match record payload {available_length}"
)));
}
let length = usize::try_from(declared_length)
.map_err(|_| spill_error("indexed spill row length is outside address space"))?;
let mut payload = Vec::new();
payload.try_reserve_exact(length).map_err(|error| {
spill_error(format!(
"unable to allocate indexed spill row payload of {length} bytes: {error}"
))
})?;
payload.resize(length, 0);
self.data
.as_file_mut()
.read_exact(&mut payload)
.map_err(|error| spill_error(format!("failed to read indexed spill row: {error}")))?;
decode_physical_row_record(&payload, self.schema.physical_width())
}
}
fn read_indexed_offset(file: &mut File, position: u64) -> ExecResult<u64> {
file.seek(SeekFrom::Start(position))
.map_err(|error| spill_error(format!("failed to seek indexed spill offset: {error}")))?;
let mut encoded = [0_u8; 8];
file.read_exact(&mut encoded)
.map_err(|error| spill_error(format!("failed to read indexed spill offset: {error}")))?;
Ok(u64::from_le_bytes(encoded))
}
impl<I> SpillRows<I> {
fn new(batches: I) -> Self {
Self {
batches,
current_schema: None,
current: Vec::new().into_iter(),
}
}
}
impl<I> Iterator for SpillRows<I>
where
I: Iterator<Item = ExecResult<Batch>>,
{
type Item = ExecResult<OwnedPhysicalRow>;
fn next(&mut self) -> Option<Self::Item> {
loop {
if let Some(row) = self.current.next() {
let schema = self
.current_schema
.as_ref()
.expect("spill row iterator retains the current batch schema")
.clone();
return Some(Ok(OwnedPhysicalRow::new(schema, row)));
}
match self.batches.next()? {
Ok(batch) => {
self.current_schema = Some(batch.schema);
self.current = batch.rows.into_iter();
}
Err(error) => return Some(Err(error)),
}
}
}
}
#[cfg(test)]
mod tests;