Skip to content
Merged
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
1 change: 1 addition & 0 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

19 changes: 19 additions & 0 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -63,6 +63,24 @@ COPY --from=hotblocks-retain-builder /out/sqd-hotblocks-retain .
ENTRYPOINT ["/app/sqd-hotblocks-retain"]


FROM builder AS flush-bench-builder
ARG TARGETARCH
# `cargo bench --no-run` emits target/release/deps/flush_spill-<hash>; the cache mount may
# hold stale hashes, so take the newest non-.d artifact.
RUN --mount=type=cache,target=/usr/local/cargo/registry,sharing=locked \
--mount=type=cache,target=/usr/local/cargo/git,sharing=locked \
--mount=type=cache,target=/app/target,id=cargo-target-${TARGETARCH},sharing=locked \
cargo bench -p sqd-data --bench flush_spill --no-run \
&& mkdir -p /out \
&& cp "$(ls -t target/release/deps/flush_spill-* | grep -v '\.d$' | head -1)" /out/flush_spill


FROM debian:bookworm-slim AS flush-bench
WORKDIR /app
COPY --from=flush-bench-builder /out/flush_spill .
ENTRYPOINT ["/app/flush_spill"]


FROM builder AS archive-builder
ARG TARGETARCH
RUN --mount=type=cache,target=/usr/local/cargo/registry,sharing=locked \
Expand All @@ -73,6 +91,7 @@ RUN --mount=type=cache,target=/usr/local/cargo/registry,sharing=locked \
&& cp target/release/sqd-archive /out/


# keep this stage last: it is the default target of a bare `docker build .`
FROM debian:bookworm-slim AS sqd-archive
RUN apt-get update && apt-get install ca-certificates -y
WORKDIR /app
Expand Down
38 changes: 37 additions & 1 deletion crates/data-core/src/chunk_builder.rs
Original file line number Diff line number Diff line change
Expand Up @@ -2,7 +2,7 @@ use std::collections::BTreeMap;

use sqd_array::slice::AnyTableSlice;

use crate::ChunkProcessor;
use crate::{ChunkProcessor, PreparedChunk};

pub trait BlockChunkBuilder: ChunkBuilder {
type Block;
Expand All @@ -24,6 +24,10 @@ pub trait ChunkBuilder {
fn new_chunk_processor(&self) -> anyhow::Result<ChunkProcessor>;

fn submit_to_processor(&self, processor: &mut ChunkProcessor) -> anyhow::Result<()>;

/// In-memory equivalent of `new_chunk_processor()` + submit + `finish()`; clears the
/// builder on success. No temp files.
fn prepare_in_memory(&mut self) -> anyhow::Result<PreparedChunk>;
}

#[macro_export]
Expand Down Expand Up @@ -98,6 +102,34 @@ macro_rules! chunk_builder {
Ok(())
}

pub fn prepare_in_memory(&mut self) -> anyhow::Result<sqd_data_core::PreparedChunk> {
use sqd_array::slice::*;
let downcast = sqd_data_core::Downcast::new();
// downcast is chunk-wide: register all tables before building any
$(
sqd_data_core::register_downcast(
&self.$table.as_slice(),
&self.$table.schema(),
$builder::table_description(),
&downcast
)?;
)*
let mut tables = std::collections::BTreeMap::new();
$(
tables.insert(
stringify!($table),
sqd_data_core::PreparedTable::from_slice(
&self.$table.as_slice(),
self.$table.schema(),
$builder::table_description(),
downcast.clone()
)?
);
)*
self.clear();
Ok(tables)
}

pub fn dataset_description() -> sqd_dataset::DatasetDescriptionRef {
use sqd_dataset::*;
use std::sync::{Arc, LazyLock};
Expand Down Expand Up @@ -145,6 +177,10 @@ macro_rules! chunk_builder {
fn submit_to_processor(&self, processor: &mut sqd_data_core::ChunkProcessor) -> anyhow::Result<()> {
self.submit_to_processor(processor)
}

fn prepare_in_memory(&mut self) -> anyhow::Result<sqd_data_core::PreparedChunk> {
self.prepare_in_memory()
}
}

impl Default for $name {
Expand Down
105 changes: 102 additions & 3 deletions crates/data-core/src/table_processor.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,6 @@
use std::{collections::HashMap, sync::Arc};

use anyhow::bail;
use arrow::{
array::RecordBatch,
datatypes::{DataType, Field, SchemaRef}
Expand All @@ -9,6 +10,7 @@ use sqd_array::{
item_index_cast::cast_item_index,
schema_patch::SchemaPatch,
slice::{AnyTableSlice, AsSlice, Slice},
sort::sort_table_to_indexes,
util::build_field_offsets,
writer::ArrayWriter
};
Expand Down Expand Up @@ -43,25 +45,39 @@ impl TableWriter {

enum TableReader {
Plain(TableFile),
Sort(SortedTable)
Sort(SortedTable),
Mem(MemTable)
}

impl TableReader {
fn read_column(&mut self, dst: &mut impl ArrayWriter, i: usize, offset: usize, len: usize) -> anyhow::Result<()> {
match self {
TableReader::Plain(reader) => reader.read_column(dst, i, offset, len),
TableReader::Sort(reader) => reader.read_column(dst, i, offset, len)
TableReader::Sort(reader) => reader.read_column(dst, i, offset, len),
TableReader::Mem(reader) => reader.read_column(dst, i, offset, len)
}
}

fn into_writer(self) -> anyhow::Result<TableWriter> {
match self {
TableReader::Plain(reader) => reader.into_writer().map(TableWriter::Plain),
TableReader::Sort(reader) => reader.into_sorter().map(TableWriter::Sort)
TableReader::Sort(reader) => reader.into_sorter().map(TableWriter::Sort),
TableReader::Mem(_) => bail!("processor reuse is not supported for in-memory prepared tables")
}
}
}

/// Columns materialized in output (sort-key) order.
struct MemTable {
columns: Vec<AnyBuilder>
}

impl MemTable {
fn read_column(&self, dst: &mut impl ArrayWriter, i: usize, offset: usize, len: usize) -> anyhow::Result<()> {
self.columns[i].as_slice().slice(offset, len).write(dst)
}
}

pub struct TableProcessor {
downcast: Downcast,
schema: SchemaRef,
Expand Down Expand Up @@ -178,6 +194,72 @@ impl PreparedTable {
})
}

/// One-batch, in-memory equivalent of `TableProcessor` push + finish — no temp files.
///
/// `downcast` must have every table of the chunk registered ([`register_downcast`])
/// before any table is built. Row order within equal full sort keys is unspecified and
/// may differ from the spill path (unstable sort); group contents are identical, and
/// queries re-sort output by a row-unique primary key, so it is not client-visible.
pub fn from_slice(
records: &AnyTableSlice<'_>,
schema: SchemaRef,
desc: &TableDescription,
downcast: Downcast
) -> anyhow::Result<Self> {
let block_number_columns = desc
.downcast
.block_number
.iter()
.map(|name| schema.index_of(name))
.collect::<Result<Vec<_>, _>>()?;

let item_index_columns = desc
.downcast
.item_index
.iter()
.map(|name| schema.index_of(name))
.collect::<Result<Vec<_>, _>>()?;

let sort_key = desc
.sort_key
.iter()
.map(|name| schema.index_of(name))
.collect::<Result<Vec<_>, _>>()?;

let order = (!sort_key.is_empty() && records.len() > 0).then(|| sort_table_to_indexes(records, &sort_key));

let columns = (0..records.num_columns())
.map(|i| {
let mut b = AnyBuilder::new(schema.field(i).data_type());
match &order {
Some(order) => records.column(i).write_indexes(&mut b, order.iter().copied())?,
None => records.column(i).write(&mut b)?
}
Ok(b)
})
.collect::<anyhow::Result<Vec<_>>>()?;

let prepared_schema = downcast_schema(
schema.clone(),
&block_number_columns,
&item_index_columns,
downcast.get_block_number_type(),
downcast.get_item_index_type()
);

Ok(Self {
downcast,
block_number_columns,
item_index_columns,
column_offsets: build_field_offsets(0, schema.fields()),
writer_schema: schema,
prepared_schema,
reader: TableReader::Mem(MemTable { columns }),
buffers: HashMap::with_capacity(3),
num_rows: records.len()
})
}

pub fn into_processor(self) -> anyhow::Result<TableProcessor> {
self.downcast.reset();
Ok(TableProcessor {
Expand Down Expand Up @@ -262,6 +344,23 @@ impl PreparedTable {
}
}

/// Register a table's downcast columns without processing it — the in-memory path must
/// register the whole chunk before building any [`PreparedTable`].
pub fn register_downcast(
records: &AnyTableSlice<'_>,
schema: &SchemaRef,
desc: &TableDescription,
downcast: &Downcast
) -> anyhow::Result<()> {
for name in desc.downcast.block_number.iter() {
downcast.reg_block_number(&records.column(schema.index_of(name)?));
}
for name in desc.downcast.item_index.iter() {
downcast.reg_item_index(&records.column(schema.index_of(name)?));
}
Ok(())
}

fn downcast_schema(
schema: SchemaRef,
block_number_columns: &[usize],
Expand Down
Loading
Loading