use std::sync::Arc;
use acta::{
Array, Column, LogicalType, PrimitiveArray, Reader, RecordBatch, Schema, Writer, WriterOptions,
};
fn int64_batch(schema: &Schema, values: Vec<i64>) -> RecordBatch {
let row_count = values.len();
RecordBatch::try_new(
Arc::new(schema.clone()),
vec![Array::Int64(PrimitiveArray::new(values, None))],
row_count,
)
.expect("well-formed batch")
}
fn main() -> acta::Result<()> {
let path = std::env::temp_dir().join("acta-refresh-tail-example.acta");
let _ = std::fs::remove_file(&path);
let schema = Schema::new(
1,
vec![Column::new(1, "value", LogicalType::Int64, false)],
None,
);
let mut writer = Writer::create(&path, schema.clone(), WriterOptions::default())?;
writer.append(int64_batch(&schema, vec![1, 2, 3]))?;
let mut reader = Reader::open(&path)?;
println!(
"snapshot before flush: {} block(s), {} row(s)",
reader.blocks().len(),
reader.total_rows()
);
writer.flush()?;
let report = reader.refresh()?;
println!(
"refresh added {} block(s) / {} row(s), tail {}",
report.blocks_added(),
report.rows_added(),
report.incomplete_tail()
);
let mut scanned = 0;
for batch in reader.scan() {
scanned += batch?.row_count();
}
println!("scan over the refreshed snapshot: {scanned} row(s)");
let mut tail = reader.tail();
writer.append(int64_batch(&schema, vec![4, 5]))?;
let _summary = writer.finish()?;
let mut idle_polls = 0;
while idle_polls < 3 {
match tail.poll_next()? {
Some(batch) => {
println!("tailed a batch of {} row(s)", batch.row_count());
idle_polls = 0;
}
None => idle_polls += 1,
}
}
println!("tail is pending; dropping it cancels the follow");
let _ = std::fs::remove_file(&path);
Ok(())
}