use std::collections::BTreeMap;
use crate::archive::{Archive, ArchiveMut, SegmentMeta, Transaction};
use crate::error::{Error, Result};
use crate::segment::SegmentEncoder;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum SchemaPolicy {
#[default]
StopAtChange,
UnionFields,
}
pub struct CompactSpec {
pub target_rows: u64,
pub writer_props: Option<parquet::file::properties::WriterProperties>,
pub schema: SchemaPolicy,
}
impl CompactSpec {
pub fn to_rows(target_rows: u64) -> Self {
CompactSpec {
target_rows,
writer_props: None,
schema: SchemaPolicy::StopAtChange,
}
}
pub fn unioning_fields(mut self) -> Self {
self.schema = SchemaPolicy::UnionFields;
self
}
fn props(&self) -> parquet::file::properties::WriterProperties {
self.writer_props
.clone()
.unwrap_or_else(crate::segment::writer_props)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
#[non_exhaustive]
pub struct Compacted {
pub before: usize,
pub after: usize,
pub merges: usize,
}
pub fn compact_stream(
db: &mut ArchiveMut,
source_id: i64,
stream: &str,
spec: &CompactSpec,
) -> Result<Compacted> {
let metas = db.read_segment_meta(source_id, stream)?;
let mut done = Compacted {
before: metas.len(),
after: metas.len(),
merges: 0,
};
let mut at = 0usize;
while at < metas.len() {
let mut end = at;
let mut rows = 0u64;
while end < metas.len() && rows + metas[end].1.rows <= spec.target_rows.max(1) {
rows += metas[end].1.rows;
end += 1;
}
if end - at < 2 {
at = (at + 1).max(end);
continue;
}
let mut blobs = Vec::with_capacity(end - at);
for (seq, _) in &metas[at..end] {
match db.read_segment_bytes(source_id, stream, *seq)? {
Some(b) => blobs.push(b),
None => break,
}
}
let (merged, merged_rows, consumed) =
concat_parquet(&blobs, spec.props(), spec.schema, stream)?;
if consumed < 2 {
at += 1;
continue;
}
let run = &metas[at..at + consumed];
let meta = SegmentMeta {
rows: merged_rows,
first_ts: run[0].1.first_ts,
last_ts: run[consumed - 1].1.last_ts,
};
let seq = run[0].0;
db.transaction(|tx| {
for (s, _) in run {
tx.delete_segment(source_id, stream, *s)?;
}
tx.insert_segment(source_id, stream, seq, &meta, &merged)
})?;
done.merges += 1;
done.after -= consumed - 1;
at += consumed;
}
Ok(done)
}
pub fn compact(db: &mut ArchiveMut, spec: &CompactSpec) -> Result<Compacted> {
let mut total = Compacted::default();
let sources = db.read_sources()?;
for src in sources {
for stream in db.all_streams(src.id)? {
let one = compact_stream(db, src.id, &stream, spec)?;
total.before += one.before;
total.after += one.after;
total.merges += one.merges;
}
}
if total.merges > 0 {
db.incremental_vacuum(u32::MAX)?;
}
Ok(total)
}
fn concat_parquet(
blobs: &[Vec<u8>],
props: parquet::file::properties::WriterProperties,
policy: SchemaPolicy,
stream: &str,
) -> Result<(Vec<u8>, u64, usize)> {
use arrow::array::new_null_array;
use arrow::datatypes::{Field, Schema};
use arrow::record_batch::RecordBatch;
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use parquet::arrow::ArrowWriter;
use std::sync::Arc;
let open = |b: &Vec<u8>| {
ParquetRecordBatchReaderBuilder::try_new(bytes::Bytes::from(b.clone()))
.map_err(|e| Error::Message(format!("failed to open a {stream} segment to merge: {e}")))
};
let first = open(&blobs[0])?;
let (consumed, schema) = match policy {
SchemaPolicy::StopAtChange => {
let schema = first.schema().clone();
let mut consumed = 1usize;
while consumed < blobs.len() && open(&blobs[consumed])?.schema() == &schema {
consumed += 1;
}
(consumed, schema)
}
SchemaPolicy::UnionFields => {
let mut order: Vec<Field> = Vec::new();
let mut seen: BTreeMap<String, usize> = BTreeMap::new();
let mut consumed = 0usize;
'segments: while consumed < blobs.len() {
let s = open(&blobs[consumed])?.schema().clone();
for f in s.fields() {
if let Some(i) = seen.get(f.name()) {
if &order[*i] != f.as_ref() {
break 'segments;
}
}
}
for f in s.fields() {
if !seen.contains_key(f.name()) {
seen.insert(f.name().clone(), order.len());
order.push(f.as_ref().clone());
}
}
consumed += 1;
}
if consumed == 0 {
(1, first.schema().clone())
} else {
let mut in_every: BTreeMap<String, usize> = BTreeMap::new();
for b in &blobs[..consumed] {
for f in open(b)?.schema().fields() {
*in_every.entry(f.name().clone()).or_insert(0) += 1;
}
}
let fields: Vec<Field> = order
.into_iter()
.map(|f| {
if in_every.get(f.name()) == Some(&consumed) {
f
} else {
let md = f.metadata().clone();
Field::new(f.name(), f.data_type().clone(), true).with_metadata(md)
}
})
.collect();
(consumed, Arc::new(Schema::new(fields)))
}
}
};
let mut buf: Vec<u8> = Vec::new();
let mut rows = 0u64;
{
let mut writer = ArrowWriter::try_new(&mut buf, schema.clone(), Some(props))
.map_err(|e| Error::Message(format!("failed to open a merged {stream} writer: {e}")))?;
for b in &blobs[..consumed] {
for batch in open(b)?.build().map_err(|e| {
Error::Message(format!("failed to read a {stream} segment to merge: {e}"))
})? {
let batch = batch.map_err(|e| {
Error::Message(format!("failed to read a {stream} batch to merge: {e}"))
})?;
rows += batch.num_rows() as u64;
let batch = if batch.schema().fields() == schema.fields() {
batch
} else {
let n = batch.num_rows();
let columns = schema
.fields()
.iter()
.map(|f| match batch.schema().index_of(f.name()) {
Ok(i) => batch.column(i).clone(),
Err(_) => new_null_array(f.data_type(), n),
})
.collect();
RecordBatch::try_new(schema.clone(), columns).map_err(|e| {
Error::Message(format!("failed to widen a {stream} batch: {e}"))
})?
};
writer.write(&batch).map_err(|e| {
Error::Message(format!("failed to write a merged {stream} batch: {e}"))
})?;
}
}
writer.close().map_err(|e| {
Error::Message(format!("failed to finish a merged {stream} segment: {e}"))
})?;
}
Ok((buf, rows, consumed))
}
pub trait ColumnFilter {
fn keep(&self, field: &arrow::datatypes::Field) -> bool;
fn is_data(&self, field: &arrow::datatypes::Field) -> bool;
}
pub struct CopySpec<'a> {
pub start: i64,
pub end: i64,
pub keep_streams: Option<&'a dyn Fn(&str) -> bool>,
pub metadata_extra: Option<&'a BTreeMap<String, String>>,
pub keep_columns: Option<&'a dyn ColumnFilter>,
pub writer_props: Option<parquet::file::properties::WriterProperties>,
}
impl CopySpec<'_> {
fn props(&self) -> parquet::file::properties::WriterProperties {
self.writer_props
.clone()
.unwrap_or_else(crate::segment::writer_props)
}
pub fn everything() -> Self {
CopySpec {
start: i64::MIN,
end: i64::MAX,
keep_streams: None,
metadata_extra: None,
keep_columns: None,
writer_props: None,
}
}
}
pub fn shared_sources(a: &Archive, b: &Archive) -> Result<Vec<String>> {
let in_a: std::collections::BTreeSet<String> = a
.read_sources()?
.into_iter()
.filter_map(|s| s.uuid)
.collect();
Ok(b.read_sources()?
.into_iter()
.filter_map(|s| s.uuid)
.filter(|u| in_a.contains(u))
.collect())
}
pub fn copy_sources_into(
src: &Archive,
tx: &Transaction<'_>,
spec: &CopySpec<'_>,
encoder: &dyn SegmentEncoder,
) -> Result<usize> {
src.read_snapshot(|src| copy_sources_snapshotted(src, tx, spec, encoder))
}
fn copy_sources_snapshotted(
src: &Archive,
tx: &Transaction<'_>,
spec: &CopySpec<'_>,
encoder: &dyn SegmentEncoder,
) -> Result<usize> {
let sources = src.read_sources()?;
let mut copied = 0usize;
for rec in &sources {
crate::segment::check_encoder(rec.id, &rec.meta.metadata, encoder)?;
let mut meta = rec.meta.clone();
if let Some(extra) = spec.metadata_extra {
for (k, v) in extra {
meta.metadata.insert(k.clone(), v.clone());
}
}
let id = tx.insert_source_with_uuid(&meta, rec.uuid.as_deref())?;
if rec.complete {
tx.mark_complete(id)?;
}
copied += 1;
for table in src.all_streams(rec.id)? {
if let Some(keep) = spec.keep_streams {
if !keep(table.as_str()) {
continue;
}
}
let mut seq = 0u64;
for segment in src.segments_overlapping(rec.id, &table, spec.start, spec.end)? {
match spec.keep_columns {
Some(keep) => {
if let Some(projected) =
project_segment_columns(&segment.bytes, keep, spec.props())?
{
tx.insert_segment(id, &table, seq, &segment.meta, &projected)?;
seq += 1;
}
}
None => {
tx.insert_segment_with_index(
id,
&table,
seq,
&segment.meta,
&segment.bytes,
segment.caller_index.as_deref(),
)?;
seq += 1;
}
}
}
let tail = src.live_wal(rec.id, &table)?;
let (Some(first), Some(last)) = (tail.first(), tail.last()) else {
continue;
};
if last.ts < spec.start || first.ts > spec.end {
continue;
}
let materialized = crate::segment::materialize(encoder, &table, &tail)?;
if let Some(materialized) = materialized {
let meta = SegmentMeta {
rows: materialized.rows,
first_ts: materialized.first_ts,
last_ts: materialized.last_ts,
};
match spec.keep_columns {
Some(keep) => {
if let Some(projected) =
project_segment_columns(&materialized.bytes, keep, spec.props())?
{
tx.insert_segment(id, &table, seq, &meta, &projected)?;
}
}
None => {
tx.insert_segment_with_index(
id,
&table,
seq,
&meta,
&materialized.bytes,
materialized.index.as_deref(),
)?;
}
}
}
}
for (ts, offset) in src.read_clock_offsets(rec.id)? {
tx.insert_clock_offset(id, ts, offset)?;
}
for name in src.caller_row_streams(rec.id)? {
if let Some(keep) = spec.keep_streams {
if !keep(name.as_str()) {
continue;
}
}
let rows = src.read_caller_rows(rec.id, &name, spec.start, spec.end)?;
tx.insert_caller_rows(id, &name, &rows)?;
}
}
Ok(copied)
}
pub fn project_segment_columns(
bytes: &[u8],
keep: &dyn ColumnFilter,
props: parquet::file::properties::WriterProperties,
) -> Result<Option<Vec<u8>>> {
use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder;
use parquet::arrow::ArrowWriter;
let builder = ParquetRecordBatchReaderBuilder::try_new(bytes::Bytes::copy_from_slice(bytes))
.map_err(|e| Error::Message(format!("failed to open a segment for projection: {e}")))?;
let schema = builder.schema().clone();
let mut indices: Vec<usize> = Vec::new();
let mut has_value = false;
for (i, f) in schema.fields().iter().enumerate() {
if keep.keep(f) {
indices.push(i);
has_value |= keep.is_data(f);
}
}
if !has_value {
return Ok(None);
}
let projected_schema = std::sync::Arc::new(
schema
.project(&indices)
.map_err(|e| Error::Message(format!("failed to project a segment schema: {e}")))?,
);
let reader = builder
.build()
.map_err(|e| Error::Message(format!("failed to read a segment for projection: {e}")))?;
let mut buf: Vec<u8> = Vec::new();
{
let mut writer =
ArrowWriter::try_new(&mut buf, projected_schema, Some(props)).map_err(|e| {
Error::Message(format!("failed to open a projected segment writer: {e}"))
})?;
for batch in reader {
let batch = batch
.map_err(|e| Error::Message(format!("failed to read a segment batch: {e}")))?;
let projected = batch
.project(&indices)
.map_err(|e| Error::Message(format!("failed to project a segment batch: {e}")))?;
writer.write(&projected).map_err(|e| {
Error::Message(format!("failed to write a projected segment batch: {e}"))
})?;
}
writer
.close()
.map_err(|e| Error::Message(format!("failed to finalize a projected segment: {e}")))?;
}
Ok(Some(buf))
}
#[cfg(test)]
mod tests {
use crate::archive::ArchiveMut;
#[test]
fn every_schema_table_is_either_copied_or_deliberately_dropped() {
const COPIED: &[&str] = &["sources", "segments", "wal", "clock_offsets", "caller_rows"];
const NOT_CARRIED: &[&str] = &[];
let dir = tempfile::tempdir().unwrap();
let path = dir.path().join("schema.dendro");
let db = ArchiveMut::create(&path).unwrap();
let mut actual = db.user_table_names().unwrap();
actual.sort();
let mut expected: Vec<String> = COPIED
.iter()
.chain(NOT_CARRIED)
.map(|s| s.to_string())
.collect();
expected.sort();
assert_eq!(
actual, expected,
"the schema changed. [`copy_sources_into`] copies a fixed set of tables, so a \
new one is silently dropped from every combined/filtered/dumped archive until it \
is handled. Copy it, or list it in NOT_CARRIED with the reason."
);
}
}