Skip to content
Closed
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
16 changes: 8 additions & 8 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -33,10 +33,10 @@ rust-version = "1.89"

[workspace.dependencies]
anyhow = "1.0"
arrow-array = "58.0.0"
arrow-cast = "58.0.0"
arrow-json = "58.0.0"
arrow-schema = "58.0.0"
arrow-array = "59.0.0"
arrow-cast = "59.0.0"
arrow-json = "59.0.0"
arrow-schema = "59.0.0"
assert-json-diff = "2.0"
assert_cmd = "2.1"
async-recursion = "1.1.1"
Expand All @@ -58,10 +58,10 @@ futures-util = "0.3.31"
geo = "0.33.0"
geo-traits = "0.3.0"
geo-types = "0.7.16"
geoarrow-array = "0.8.0"
geoarrow-schema = "0.8.0"
geoarrow-array = "0.9.0"
geoarrow-schema = "0.9.0"
geojson = "1.0.0"
geoparquet = "0.8.0"
geoparquet = "0.9.0"
getrandom = { version = "0.4.0", features = ["wasm_js"] }
http = "1.1"
indexmap = { version = "2.10.0", features = ["serde"] }
Expand All @@ -73,7 +73,7 @@ log = "0.4.25"
mime = "0.3.17"
mockito = "1.5"
object_store = "0.13.0"
parquet = { version = "58.0.0" }
parquet = { version = "59.0.0" }
pgstac = "0.4.9"
quote = "1.0"
referencing = { version = "0.58.2", features = ["retrieve-async"] }
Expand Down
4 changes: 2 additions & 2 deletions crates/duckdb/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -17,10 +17,10 @@ bundled = ["duckdb/bundled"]
async = ["dep:futures", "dep:futures-core", "stac/async"]

[dependencies]
arrow-array.workspace = true
arrow-array = { workspace = true, features = ["ffi"] }
futures = { workspace = true, optional = true }
futures-core = { workspace = true, optional = true }
arrow-schema.workspace = true
arrow-schema = { workspace = true, features = ["ffi"] }
chrono.workspace = true
cql2.workspace = true
duckdb.workspace = true
Expand Down
62 changes: 48 additions & 14 deletions crates/duckdb/src/client.rs
Original file line number Diff line number Diff line change
@@ -1,6 +1,9 @@
use crate::{Error, Extension, Result};
use arrow_array::{RecordBatch, RecordBatchIterator};
use arrow_schema::{ArrowError, SchemaRef};
use arrow_array::{
RecordBatch, RecordBatchIterator, StructArray,
ffi::{FFI_ArrowArray, FFI_ArrowSchema},
};
use arrow_schema::{ArrowError, Schema, SchemaRef};
use chrono::DateTime;
use cql2::{Expr, ToDuckSQL};
use duckdb::{Connection, Statement, types::Value};
Expand All @@ -13,7 +16,7 @@ use stac::api::{
};
use stac::{Collection, SpatialExtent, TemporalExtent, geoarrow::DATETIME_COLUMNS};
use std::ops::{Deref, DerefMut};
use std::sync::Mutex;
use std::sync::{Arc, Mutex};

/// Default hive partitioning value
pub const DEFAULT_USE_HIVE_PARTITIONING: bool = false;
Expand Down Expand Up @@ -250,11 +253,7 @@ impl Client {
let mut statement = self.prepare(&sql)?;
statement.execute(duckdb::params_from_iter(params))?;
log::debug!("query complete");
Ok(SearchArrowBatchIter::new(
statement,
self.convert_wkb,
self.remove_filename_column,
))
SearchArrowBatchIter::new(statement, self.convert_wkb, self.remove_filename_column)
} else {
Ok(SearchArrowBatchIter::empty(
self.convert_wkb,
Expand Down Expand Up @@ -720,14 +719,18 @@ pub struct SearchArrowBatchIter<'conn> {
}

impl<'conn> SearchArrowBatchIter<'conn> {
fn new(statement: Statement<'conn>, convert_wkb: bool, remove_filename_column: bool) -> Self {
let schema = Some(statement.schema());
Self {
fn new(
statement: Statement<'conn>,
convert_wkb: bool,
remove_filename_column: bool,
) -> Result<Self> {
let schema = Some(schema_from_duckdb(&statement.schema())?);
Ok(Self {
statement: Some(statement),
convert_wkb,
remove_filename_column,
schema,
}
})
}

fn empty(convert_wkb: bool, remove_filename_column: bool) -> Self {
Expand Down Expand Up @@ -764,8 +767,9 @@ impl<'conn> Iterator for SearchArrowBatchIter<'conn> {

match statement.step() {
Ok(Some(struct_array)) => {
let record_batch = RecordBatch::from(&struct_array);
match self.finalize_batch(record_batch) {
match struct_array_from_duckdb(&struct_array)
.and_then(|struct_array| self.finalize_batch(RecordBatch::from(struct_array)))
{
Ok(batch) => Some(Ok(batch)),
Err(err) => {
self.statement = None;
Expand All @@ -785,6 +789,36 @@ impl<'conn> Iterator for SearchArrowBatchIter<'conn> {
}
}

/// Converts a [`duckdb::arrow`] struct array into our version of arrow via the
/// [C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html).
///
/// DuckDB may depend on a different version of arrow than we do. The C Data
/// Interface structs share a `#[repr(C)]` layout across arrow versions, so we
/// can move the exported structs between versions without copying any data.
fn struct_array_from_duckdb(
struct_array: &duckdb::arrow::array::StructArray,
) -> Result<StructArray> {
use duckdb::arrow::array::Array;
let (mut array, mut schema) = duckdb::arrow::ffi::to_ffi(&struct_array.to_data())?;
// SAFETY: both pointers are valid, aligned, and initialized C Data Interface
// structs. `from_raw` leaves empty structs behind, so the originals are
// released only once, by the moved values.
let array = unsafe { FFI_ArrowArray::from_raw((&raw mut array).cast()) };
let schema = unsafe { FFI_ArrowSchema::from_raw((&raw mut schema).cast()) };
// SAFETY: `array` and `schema` were produced together by `to_ffi`.
let data = unsafe { arrow_array::ffi::from_ffi(array, &schema)? };
Ok(StructArray::from(data))
}

/// Converts a [`duckdb::arrow`] schema into our version of arrow via the
/// [C Data Interface](https://arrow.apache.org/docs/format/CDataInterface.html).
fn schema_from_duckdb(schema: &duckdb::arrow::datatypes::Schema) -> Result<SchemaRef> {
let mut schema = duckdb::arrow::ffi::FFI_ArrowSchema::try_from(schema)?;
// SAFETY: see `struct_array_from_duckdb`.
let schema = unsafe { FFI_ArrowSchema::from_raw((&raw mut schema).cast()) };
Ok(Arc::new(Schema::try_from(&schema)?))
}

fn remove_column(mut record_batch: RecordBatch, name: &str) -> RecordBatch {
if let Some((index, _)) = record_batch.schema().column_with_name(name) {
record_batch.remove_column(index);
Expand Down
8 changes: 8 additions & 0 deletions crates/duckdb/src/error.rs
Original file line number Diff line number Diff line change
Expand Up @@ -4,6 +4,10 @@ use thiserror::Error;
#[derive(Debug, Error)]
#[non_exhaustive]
pub enum Error {
/// [arrow_schema::ArrowError]
#[error(transparent)]
Arrow(#[from] arrow_schema::ArrowError),

/// [chrono::format::ParseError]
#[error(transparent)]
ChronoParse(#[from] chrono::format::ParseError),
Expand All @@ -16,6 +20,10 @@ pub enum Error {
#[error(transparent)]
DuckDB(#[from] duckdb::Error),

/// [duckdb::arrow::error::ArrowError], from the version of arrow used by duckdb
#[error(transparent)]
DuckDBArrow(#[from] duckdb::arrow::error::ArrowError),

/// [geoarrow_schema::error::GeoArrowError]
#[error(transparent)]
GeoArrow(#[from] geoarrow_schema::error::GeoArrowError),
Expand Down
2 changes: 1 addition & 1 deletion crates/wasm/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,7 @@ debug = ["console_error_panic_hook"]
console_error_panic_hook = { version = "0.1.6", optional = true }
arrow-array.workspace = true
arrow-schema.workspace = true
arrow-wasm = { git = "https://github.com/kylebarron/arrow-wasm", rev = "4da0aa2b45c7ffd8c7e7449274a4b6d84f10cf94" }
arrow-wasm = { git = "https://github.com/kylebarron/arrow-wasm", rev = "ee163ba9cfa1c5d078e7625ac8bcd3cd1e6bda8e" }
getrandom = { version = "0.4", features = ["wasm_js"] }
getrandom_03 = { package = "getrandom", version = "0.3", features = ["wasm_js"] }
serde.workspace = true
Expand Down
Loading