diff --git a/Cargo.lock b/Cargo.lock index d2f6993..2ffbe2b 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -609,6 +609,7 @@ dependencies = [ "clap", "comfy-table", "nu-ansi-term", + "parquet", "reedline", "syntect", "tempfile", @@ -1096,6 +1097,32 @@ dependencies = [ "windows-link 0.2.1", ] +[[package]] +name = "parquet" +version = "59.3.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "ff322f54b1a0f9288e614ed1f2d329b380af5476420db19f46ffb865e1163d73" +dependencies = [ + "ahash", + "arrow-array", + "arrow-buffer", + "arrow-data", + "arrow-ipc", + "arrow-schema", + "arrow-select", + "base64", + "bytes", + "chrono", + "half", + "hashbrown", + "num-bigint", + "num-integer", + "num-traits", + "seq-macro", + "snap", + "twox-hash", +] + [[package]] name = "path-slash" version = "0.2.1" @@ -1299,6 +1326,12 @@ version = "1.0.28" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "8a7852d02fc848982e0c167ef163aaff9cd91dc640ba85e263cb1ce46fae51cd" +[[package]] +name = "seq-macro" +version = "0.3.6" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "1bc711410fbe7399f390ca1c3b60ad0f53f80e95c5eb935e52268a0e2cd49acc" + [[package]] name = "serde" version = "1.0.229" @@ -1418,6 +1451,12 @@ version = "1.16.2" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "f9395f0f0eee849a9b707b2f06bb92a6a422090e2123bb2ef8e87a0e61892a8e" +[[package]] +name = "snap" +version = "1.1.2" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "199905e6153d6405f9728fe44daace35f8f837bbf830bb6e85fbd5828709a886" + [[package]] name = "strip-ansi-escapes" version = "0.2.1" @@ -1632,6 +1671,12 @@ version = "1.1.2+spec-1.1.0" source = "registry+https://github.com/rust-lang/crates.io-index" checksum = "7d56353a2a665ad0f41a421187180aab746c8c325620617ad883a99a1cbe66d2" +[[package]] +name = "twox-hash" +version = "2.1.4" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "5283634e518fe9e82c7b20520bb4bc209009fd16c82077c802f8111ecbb0117a" + [[package]] name = "unicode-ident" version = "1.0.26" diff --git a/Cargo.toml b/Cargo.toml index 303b6d3..18874b9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -20,6 +20,7 @@ arrow-schema = "59.3.0" clap = { version = "4.6.7", features = ["derive"] } comfy-table = "8.0.1" nu-ansi-term = "0.50.3" +parquet = { version = "59.3.0", default-features = false, features = ["arrow", "snap"] } reedline = "0.52.0" syntect = "5.3.0" terminal-colorsaurus = "1.0.3" diff --git a/README.md b/README.md index 961fb4d..1a0cdc3 100644 --- a/README.md +++ b/README.md @@ -21,7 +21,7 @@ A command-line tool for querying databases via [ADBC](https://arrow.apache.org/a - **Interactive SQL shell** - Execute SQL queries with command history and intuitive navigation - **Syntax highlighting** - SQL queries highlighted for improved readability - **Formatted output** - Results displayed in clean, aligned tables with dynamic column width -- **File export** - Export query results to JSON, CSV, or Arrow IPC files +- **File export** - Export query results to JSON, JSON lines, CSV, Arrow IPC, or Parquet files - **Fast and lightweight** - Built in Rust for high performance and minimal resource usage ## Installation @@ -107,8 +107,11 @@ Execute a query and output the result to a file: ```sh databow --driver duckdb --query "SELECT 42 AS the_answer" --output result.json +databow --driver duckdb --query "SELECT 42 AS the_answer" --output result.jsonl databow --driver duckdb --query "SELECT 42 AS the_answer" --output result.csv databow --driver duckdb --query "SELECT 42 AS the_answer" --output result.arrow +databow --driver duckdb --query "SELECT 42 AS the_answer" --output result.arrows +databow --driver duckdb --query "SELECT 42 AS the_answer" --output result.parquet ``` ## Reference diff --git a/docs/index.md b/docs/index.md index 06e5a75..2954f7e 100644 --- a/docs/index.md +++ b/docs/index.md @@ -19,5 +19,5 @@ databow is a command-line tool for querying databases. - **Interactive SQL shell** - Execute SQL queries with command history and intuitive navigation - **Syntax highlighting** - SQL queries highlighted for improved readability - **Formatted output** - Results displayed in clean, aligned tables with dynamic column width -- **File export** - Export query results to JSON, CSV, or Arrow IPC files +- **File export** - Export query results to JSON, JSON lines, CSV, Arrow IPC, or Parquet files - **Fast and lightweight** - Built in Rust for high performance and minimal resource usage diff --git a/docs/reference.md b/docs/reference.md index 25c203a..14bcd52 100644 --- a/docs/reference.md +++ b/docs/reference.md @@ -108,11 +108,14 @@ databow --driver duckdb --query "SELECT 42 AS the_answer" --output result.json The output format is inferred from the file extension: -| Extension | Format | -|-----------------|-----------| -| `.json` | JSON | -| `.csv` | CSV | -| `.arrow`, `.ipc`| Arrow IPC | +| Extension | Format | +|-----------------|------------------| +| `.json` | JSON | +| `.jsonl` | JSON lines | +| `.csv` | CSV | +| `.arrow`, `.ipc`| Arrow IPC file | +| `.arrows` | Arrow IPC stream | +| `.parquet` | Parquet | ## --help diff --git a/docs/tutorial.md b/docs/tutorial.md index 59ec6ae..e38cb24 100644 --- a/docs/tutorial.md +++ b/docs/tutorial.md @@ -143,7 +143,7 @@ $ databow --profile warehouse --file query.sql └───────────────────┘ ``` -Instead of printing query results to stdout, the [`--output` argument](/reference/#-output) can be used to write results to JSON, CSV, or Arrow IPC files: +Instead of printing query results to stdout, the [`--output` argument](/reference/#-output) can be used to write results to JSON, JSON lines, CSV, Arrow IPC file, Arrow IPC stream, or Parquet files: ```console $ databow --profile warehouse --query "SELECT * FROM penguins" --output penguins.csv diff --git a/src/cli.rs b/src/cli.rs index 8dac594..9dd9061 100644 --- a/src/cli.rs +++ b/src/cli.rs @@ -1,6 +1,7 @@ // Copyright 2026 Columnar Technologies Inc. // SPDX-License-Identifier: Apache-2.0 +use crate::output::OutputFormat; use crate::table::TableMode; use clap::{Arg, ArgAction, Command, value_parser}; use std::path::PathBuf; @@ -84,7 +85,8 @@ pub fn parse_args() -> AppConfig { Arg::new("output") .long("output") .help("Write result to file") - .value_name("file"), + .value_name("file") + .value_parser(parse_output_path), ]; let command = Command::new("databow") .version(env!("CARGO_PKG_VERSION")) @@ -163,7 +165,7 @@ pub fn parse_args() -> AppConfig { .copied() .unwrap_or_default(); - let output_path = matches.get_one::("output").map(PathBuf::from); + let output_path = matches.get_one::("output").cloned(); if output_path.is_some() && matches!(query_source, QuerySource::Interactive) { eprintln!("Error: --output cannot be used in interactive mode"); exit(1); @@ -177,6 +179,12 @@ pub fn parse_args() -> AppConfig { } } +fn parse_output_path(s: &str) -> Result { + let path = PathBuf::from(s); + OutputFormat::from_path(&path)?; + Ok(path) +} + fn uri_has_driver_scheme(uri: &str) -> bool { let Some(idx) = uri.find(':') else { return false; @@ -202,6 +210,26 @@ fn parse_option(option: &str) -> Result<(String, String), String> { mod tests { use super::*; + #[test] + fn test_parse_output_path_valid() { + assert_eq!( + parse_output_path("out.parquet").unwrap(), + PathBuf::from("out.parquet") + ); + } + + #[test] + fn test_parse_output_path_unsupported_extension() { + let err = parse_output_path("out.xyz").unwrap_err(); + assert!(err.contains("Unsupported file extension")); + } + + #[test] + fn test_parse_output_path_no_extension() { + let err = parse_output_path("out").unwrap_err(); + assert!(err.contains("no file extension")); + } + #[test] fn test_connection_source_direct() { let source = ConnectionSource::Direct { diff --git a/src/output.rs b/src/output.rs index 690335a..7264591 100644 --- a/src/output.rs +++ b/src/output.rs @@ -2,26 +2,35 @@ // SPDX-License-Identifier: Apache-2.0 use arrow::csv::writer::Writer as CsvWriter; -use arrow::ipc::writer::FileWriter as IpcWriter; -use arrow::json::writer::{JsonArray, Writer as JsonWriter}; +use arrow::ipc::writer::{FileWriter as IpcWriter, StreamWriter as IpcStreamWriter}; +use arrow::json::writer::{JsonArray, LineDelimited, Writer as JsonWriter}; use arrow_array::RecordBatch; use arrow_schema::ArrowError; +use parquet::arrow::ArrowWriter; +use parquet::basic::Compression; +use parquet::file::properties::WriterProperties; use std::fs::File; use std::path::Path; #[derive(Debug, Clone, Copy, PartialEq)] pub enum OutputFormat { Json, + Jsonl, Csv, Arrow, + ArrowStream, + Parquet, } impl OutputFormat { pub fn from_path(path: &Path) -> Result { match path.extension().and_then(|ext| ext.to_str()) { Some("json") => Ok(OutputFormat::Json), + Some("jsonl") => Ok(OutputFormat::Jsonl), Some("csv") => Ok(OutputFormat::Csv), Some("arrow" | "ipc") => Ok(OutputFormat::Arrow), + Some("arrows") => Ok(OutputFormat::ArrowStream), + Some("parquet") => Ok(OutputFormat::Parquet), Some(ext) => Err(format!("Unsupported file extension: '.{ext}'")), None => Err("Cannot infer format: no file extension".to_string()), } @@ -38,8 +47,11 @@ pub fn write_batches_to_file(batches: &[RecordBatch], path: &Path) -> Result<(), match format { OutputFormat::Json => write_json(batches, file), + OutputFormat::Jsonl => write_jsonl(batches, file), OutputFormat::Csv => write_csv(batches, file), OutputFormat::Arrow => write_arrow_ipc(batches, file), + OutputFormat::ArrowStream => write_arrow_stream(batches, file), + OutputFormat::Parquet => write_parquet(batches, file), } } @@ -53,6 +65,16 @@ fn write_json(batches: &[RecordBatch], file: File) -> Result<(), ArrowError> { Ok(()) } +fn write_jsonl(batches: &[RecordBatch], file: File) -> Result<(), ArrowError> { + let mut writer = JsonWriter::<_, LineDelimited>::new(file); + for batch in batches { + writer.write(batch)?; + } + writer.finish()?; + + Ok(()) +} + fn write_csv(batches: &[RecordBatch], file: File) -> Result<(), ArrowError> { let mut writer = CsvWriter::new(file); for batch in batches { @@ -76,6 +98,37 @@ fn write_arrow_ipc(batches: &[RecordBatch], file: File) -> Result<(), ArrowError Ok(()) } +fn write_arrow_stream(batches: &[RecordBatch], file: File) -> Result<(), ArrowError> { + if batches.is_empty() { + return Ok(()); + } + let schema = batches[0].schema(); + let mut writer = IpcStreamWriter::try_new(file, &schema)?; + for batch in batches { + writer.write(batch)?; + } + writer.finish()?; + + Ok(()) +} + +fn write_parquet(batches: &[RecordBatch], file: File) -> Result<(), ArrowError> { + if batches.is_empty() { + return Ok(()); + } + let schema = batches[0].schema(); + let props = WriterProperties::builder() + .set_compression(Compression::SNAPPY) + .build(); + let mut writer = ArrowWriter::try_new(file, schema, Some(props))?; + for batch in batches { + writer.write(batch)?; + } + writer.close()?; + + Ok(()) +} + #[cfg(test)] mod tests { use super::*; @@ -243,6 +296,119 @@ mod tests { assert!(!path.exists()); } + #[test] + fn test_output_format_from_path_arrows() { + let path = Path::new("output.arrows"); + assert_eq!( + OutputFormat::from_path(path).unwrap(), + OutputFormat::ArrowStream + ); + } + + #[test] + fn test_output_format_from_path_parquet() { + let path = Path::new("output.parquet"); + assert_eq!( + OutputFormat::from_path(path).unwrap(), + OutputFormat::Parquet + ); + } + + #[test] + fn test_output_format_from_path_jsonl() { + let path = Path::new("output.jsonl"); + assert_eq!(OutputFormat::from_path(path).unwrap(), OutputFormat::Jsonl); + } + + #[test] + fn test_write_arrow_stream() { + let dir = tempdir().unwrap(); + let path = dir.path().join("output.arrows"); + let batch = create_test_batch(); + + write_batches_to_file(std::slice::from_ref(&batch), &path).unwrap(); + + // Verify by reading it back as an IPC stream + let file = File::open(&path).unwrap(); + let reader = arrow::ipc::reader::StreamReader::try_new(file, None).unwrap(); + let read_batches: Vec = reader.map(|r| r.unwrap()).collect(); + + assert_eq!(read_batches.len(), 1); + assert_eq!(read_batches[0].num_rows(), batch.num_rows()); + assert_eq!(read_batches[0].num_columns(), batch.num_columns()); + } + + #[test] + fn test_write_arrow_stream_empty_batches_direct() { + let dir = tempdir().unwrap(); + let path = dir.path().join("output.arrows"); + let file = File::create(&path).unwrap(); + + let result = write_arrow_stream(&[], file); + + assert!(result.is_ok()); + } + + #[test] + fn test_write_parquet() { + use parquet::arrow::arrow_reader::ParquetRecordBatchReaderBuilder; + + let dir = tempdir().unwrap(); + let path = dir.path().join("output.parquet"); + let batch = create_test_batch(); + + write_batches_to_file(std::slice::from_ref(&batch), &path).unwrap(); + + // Verify by reading it back + let file = File::open(&path).unwrap(); + let builder = ParquetRecordBatchReaderBuilder::try_new(file).unwrap(); + for column in builder.metadata().row_group(0).columns() { + assert_eq!(column.compression(), Compression::SNAPPY); + } + let reader = builder.build().unwrap(); + let read_batches: Vec = reader.map(|r| r.unwrap()).collect(); + + let total_rows: usize = read_batches.iter().map(|b| b.num_rows()).sum(); + assert_eq!(total_rows, batch.num_rows()); + assert_eq!(read_batches[0].num_columns(), batch.num_columns()); + } + + #[test] + fn test_write_parquet_empty_batches_direct() { + let dir = tempdir().unwrap(); + let path = dir.path().join("output.parquet"); + let file = File::create(&path).unwrap(); + + let result = write_parquet(&[], file); + + assert!(result.is_ok()); + } + + #[test] + fn test_write_jsonl() { + let dir = tempdir().unwrap(); + let path = dir.path().join("output.jsonl"); + let batch = create_test_batch(); + + write_batches_to_file(std::slice::from_ref(&batch), &path).unwrap(); + + let mut file = File::open(&path).unwrap(); + let mut contents = String::new(); + file.read_to_string(&mut contents).unwrap(); + + // JSON lines: one object per line, not a wrapping array + assert!(!contents.trim_start().starts_with('[')); + let lines: Vec<&str> = contents.lines().filter(|l| !l.is_empty()).collect(); + assert_eq!(lines.len(), 3); + for line in &lines { + assert!(line.trim_start().starts_with('{')); + assert!(line.trim_end().ends_with('}')); + } + assert!(contents.contains("Alice")); + assert!(contents.contains("Bob")); + assert!(contents.contains("Charlie")); + } + #[test] fn test_write_multiple_batches() { let dir = tempdir().unwrap(); diff --git a/tests/integration_test.rs b/tests/integration_test.rs index 302c35c..e68795c 100644 --- a/tests/integration_test.rs +++ b/tests/integration_test.rs @@ -479,3 +479,27 @@ fn test_timestamp_with_time_zone() { stdout ); } + +#[test] +fn test_invalid_output_extension_fails_before_connect() { + let output = Command::new("cargo") + .args([ + "run", + "--", + "--driver", + "nonexistent_driver", + "--query", + "SELECT 1", + "--output", + "out.xyz", + ]) + .output() + .expect("Failed to execute command"); + + // Exit code 2 and "invalid value" come from clap argument parsing, + // which happens before any connection attempt. + assert_eq!(output.status.code(), Some(2)); + let stderr = String::from_utf8_lossy(&output.stderr); + assert!(stderr.contains("invalid value 'out.xyz' for '--output '")); + assert!(stderr.contains("Unsupported file extension")); +}