use std::fs::File;
use arrow::array::Array;
use arrow::chunk::Chunk;
use arrow::datatypes::Schema;
use arrow::error::Result;
use std::io::Write;
use strawboat::{write, CommonCompression};
fn write_batches(path: &str, schema: Schema, chunks: &[Chunk<Box<dyn Array>>]) -> Result<()> {
let file = File::create(path)?;
let options = write::WriteOptions {
default_compression: CommonCompression::Lz4,
default_compress_ratio: None,
max_page_size: Some(8192),
forbidden_compressions: vec![],
};
let mut writer = write::NativeWriter::new(file, schema, options);
writer.start()?;
for chunk in chunks {
writer.write(chunk)?
}
writer.finish()?;
let metas = serde_json::to_vec(&writer.metas).unwrap();
let mut meta_file = File::options()
.create(true)
.write(true)
.truncate(true)
.open("/tmp/pa.st")?;
meta_file.write_all(&metas)?;
meta_file.flush()?;
Ok(())
}
fn main() -> Result<()> {
use std::env;
let args: Vec<String> = env::args().collect();
let file_path = &args[1];
let (chunk, schema) = read_chunk();
write_batches(file_path, schema, &[chunk])?;
Ok(())
}
fn read_chunk() -> (Chunk<Box<dyn Array>>, Schema) {
let file_path = "/tmp/input.parquet";
let mut reader = File::open(file_path).unwrap();
let metadata = arrow::io::parquet::read::read_metadata(&mut reader).unwrap();
let schema = arrow::io::parquet::read::infer_schema(&metadata).unwrap();
let schema = schema.filter(|_index, _field| true);
for field in &schema.fields {
let statistics =
arrow::io::parquet::read::statistics::deserialize(field, &metadata.row_groups).unwrap();
println!("{statistics:#?}");
}
let row_groups = metadata
.row_groups
.into_iter()
.enumerate()
.filter(|(index, _)| *index == 0 || *index == 1)
.map(|(_, row_group)| row_group)
.collect();
let mut chunks = arrow::io::parquet::read::FileReader::new(
reader,
row_groups,
schema.clone(),
Some(usize::MAX),
None,
None,
);
if let Some(maybe_chunk) = chunks.next() {
let chunk = maybe_chunk.unwrap();
println!("chunk len -> {:?}", chunk.len());
return (chunk, schema);
}
unreachable!()
}