Skip to content
Merged
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
20 changes: 10 additions & 10 deletions Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

4 changes: 2 additions & 2 deletions Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -18,8 +18,8 @@ crate-type = ["cdylib"]
# with (ADR 0002), so a revision is the honest way to say which one.
# A local checkout is used instead with a `paths` override in
# `.cargo/config.toml`, which is untracked on purpose.
zudb = { package = "zu", git = "https://github.com/tamnd/zu", rev = "8009a961f463d6b576509e0f752b60f44071cd33" }
zu-common = { git = "https://github.com/tamnd/zu", rev = "8009a961f463d6b576509e0f752b60f44071cd33" }
zudb = { package = "zu", git = "https://github.com/tamnd/zu", rev = "6f38950b9f1c4b6cbca37f404ee660370a945c12" }
zu-common = { git = "https://github.com/tamnd/zu", rev = "6f38950b9f1c4b6cbca37f404ee660370a945c12" }
# `extension-module` is asked for by maturin, in pyproject.toml, and
# not here. Only the build backend knows how an extension is linked on
# the platform it is building for, and a crate that turns the feature
Expand Down
12 changes: 7 additions & 5 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -102,13 +102,15 @@ conn.register("people", frame)
conn.execute("MATCH (p:people) WHERE p.age > 40 RETURN p.name AS name")
```

Anything that speaks Arrow goes in, which is a pandas or polars DataFrame, a pyarrow table or a reader over one, and a dictionary of lists is there for a caller with none of them installed. The frame arrives over the same C Data Interface a result leaves by, so a column costs a memcpy and no Python object per cell: a million rows of one integer column are read in 5 ms, and a million rows of an integer, a float and a string take 138 ms, which is the string column being the only one that allocates.
Anything that speaks Arrow goes in, which is a pandas or polars DataFrame, a pyarrow table or a reader over one, and a dictionary of lists is there for a caller with none of them installed. Nothing is copied. The frame arrives over the same C Data Interface a result leaves by, and what the engine is told is where each column is, how wide its values are and what they mean; a statement that names it builds vectors pointing straight at the caller's buffers. So registering costs what describing the columns costs and not what the rows cost: 2.6 microseconds for ten rows and 2.3 for ten million.

What this is not is a scan of the frame where it sits. DuckDB's `register` is zero copy because its executor can call back out to the Python object holding the data, and this engine has no such callback yet, so the frame is copied in and becomes a table like any other. That makes a registered frame a snapshot: changing the DataFrame afterwards changes nothing until it is registered again. The call is the one it would be either way, and the day the engine grows a scan it becomes the cheap thing under the same name.
The one column that is walked is a string column, and it is walked once. Every offset is checked at registration so that reading the frame afterwards cannot fail, which is 362 microseconds for a million strings. Two other things copy and both are said rather than hidden: a stream that arrives as several batches is concatenated into one, because a column of a table is one run of bytes and two batches are two of them, and a dictionary of Python lists is read into buffers of this client's own, because a list holds objects rather than numbers and there is nothing in it to point at.

The copy is not the slow kind, but the write is a write. Those million rows take 9.3 seconds in all, of which 0.14 is reading the frame and the rest is the engine, which is the same 0.1 to 0.2 million rows a second `load` and the appender manage on this machine. Registering is the right way to get a frame in and it is not a way to make the v0 write path faster.
Because it is not a copy, a registered frame is a view and not a snapshot. Write into the array behind it and the next statement answers what is there now, which is the thing to know about the call and the reason it is worth having. Reading one is as fast as reading a table of the database and faster where the database has to decode: over a million rows on this machine, summing an integer column takes 662 microseconds against a stored table's 833, and finding one row by a string takes 1.8 ms against 3.1.

`unregister(name)` takes the rows back out and `conn.registered` says what is registered here. Registering the same name again replaces the rows under it, which is what rerunning a cell means by it, and registering over a table this connection did not register is refused, since a frame knows nothing about the rows that were already there. Two things show the engine through and are said rather than worked around: a name that has been used keeps the columns it was used with, because a table's columns are declared by its first row and no statement alters them, and `unregister` empties the table rather than removing it, because no statement drops one. A null anywhere is refused by column and row, since a property that is null is one no row of this engine holds, and registering inside a transaction is refused for the reason an appender is.
The frame belongs to the connection it was registered on and goes when that connection does. Nothing is written to the file, so another program opening the same database has never heard of it, and nothing writes to it either: a statement that inserts into or deletes from a registered name is refused with the reason, because that memory is the caller's DataFrame. `unregister(name)` takes the name away and hands the bytes back, which is not always that instant, since a statement still reading the frame holds it until it ends. `conn.registered` says what is registered here.

Registering the same name again replaces what it stands for, columns and all, which is what rerunning a cell means by it. Registering over a table the database already holds is refused, since a statement naming it would mean the stored one. A frame with no rows is a table to match on and answers nothing, because a frame knows its columns without being told by a row. A null anywhere is refused by column and row, since a property that is null is one no row of this engine holds, and registering inside a transaction is refused because a frame is registered on the session, which is the thing the transaction is running on.

## Reading a result as columns

Expand Down Expand Up @@ -152,7 +154,7 @@ The stub is checked against the module it describes in CI: griffe reads the stub

## What works today

The list above is what this client is for. What it does so far is the core of it: `connect`, `execute` and `sql` with named parameters, results that iterate and fetch, values as Python objects both ways including dates, times, datetimes and durations, `Node`, `Rel` and `Path` as classes, `load` for building a graph with edges in it, an appender for growing one, transactions as a context manager that commits at the end of a block and rolls back when it raises, every condition as an exception class carrying its code, its position and its documentation link, results as Arrow columns and as pandas and polars frames, `register` for putting a frame under a name a statement can match on, stubs inside the wheel with a gate that keeps them true, the GIL released around every statement, every load and every copy out, and `Ctrl-C` and `interrupt()` stopping a statement without touching the connection under it. `zudb.aio` is next, and each one lands with the tests that say it works.
The list above is what this client is for. What it does so far is the core of it: `connect`, `execute` and `sql` with named parameters, results that iterate and fetch, values as Python objects both ways including dates, times, datetimes and durations, `Node`, `Rel` and `Path` as classes, `load` for building a graph with edges in it, an appender for growing one, transactions as a context manager that commits at the end of a block and rolls back when it raises, every condition as an exception class carrying its code, its position and its documentation link, results as Arrow columns and as pandas and polars frames, `register` for putting a frame under a name a statement can match on and reading it where it lies, stubs inside the wheel with a gate that keeps them true, the GIL released around every statement, every load and every copy out, and `Ctrl-C` and `interrupt()` stopping a statement without touching the connection under it. `zudb.aio` is next, and each one lands with the tests that say it works.

## Wheels

Expand Down
2 changes: 1 addition & 1 deletion python/zudb/_zudb.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -90,7 +90,7 @@ class Connection:
"""Puts a DataFrame under a name a statement can match on."""

def unregister(self, name: str) -> None:
"""Takes the rows of a registered frame back out."""
"""Takes a registered frame's name away."""

@property
def registered(self) -> list[str]:
Expand Down
22 changes: 0 additions & 22 deletions src/buffer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,7 +21,6 @@ use pyo3::types::{PyBool, PyBytes, PyDate, PyDateTime, PyDelta, PyTime};
use zu_common::temporal::days_from_civil;
use zu_common::{DurationKind, FloatBits, IntBits, LogicalType, Temporal};
use zudb::Field;
use zudb::query::Value;

use crate::value::{Duration, clock_nanos};

Expand Down Expand Up @@ -286,27 +285,6 @@ impl Column {
}
}

/// One value of this column as a statement's parameter takes it.
///
/// Owned rather than borrowed, unlike `field`, because a parameter
/// is bound for the length of a statement and the statement is
/// where the copy was always going to happen. `None` for a column
/// of byte strings, which no statement carries a literal or a
/// parameter for.
pub fn value(&self, row: usize) -> Option<Value> {
Some(match self {
Column::Int(v) => Value::Int(v[row] as i64),
Column::Float(v) => Value::Float(v[row]),
Column::Bool(v) => Value::Bool(v[row]),
Column::Str(v) => Value::Str(v[row].clone()),
Column::Bytes(_) => return None,
Column::Date(v) => Value::Temporal(Temporal::Date(v[row])),
Column::LocalTime(v) => Value::Temporal(Temporal::LocalTime(v[row])),
Column::LocalDatetime(v) => Value::Temporal(Temporal::LocalDatetime(v[row])),
Column::Duration(kind, v) => Value::Temporal(Temporal::Duration(*kind, v[row])),
})
}

/// One value of this column as the engine's appender takes it.
///
/// A field borrows rather than owning, which is the point of it on
Expand Down
97 changes: 22 additions & 75 deletions src/conn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,6 @@
//! to run statements at once want two connections. It is there so that
//! a program which shares one by accident waits rather than corrupts.

use std::collections::BTreeMap;
use std::ffi::CStr;
use std::path::PathBuf;
use std::sync::atomic::{AtomicBool, Ordering};
Expand Down Expand Up @@ -68,18 +67,6 @@ pub struct Connection {
/// asking it to stop, should not queue behind a ten second
/// statement.
alive: AtomicBool,
/// The names this connection has registered frames under, against
/// whether one is registered under them now. A name that was
/// unregistered stays here as `false`, because the empty table it
/// left behind is still a table this connection made and is still
/// the one a later `register` under the same name may fill: forget
/// it entirely and the name would be refused as somebody else's.
///
/// Kept per connection rather than per database, because a
/// registered frame is a caller's data under a caller's name and
/// another program opening the same file has no idea it is anything
/// but a table.
frames: Mutex<BTreeMap<String, bool>>,
#[pyo3(get)]
path: PathBuf,
#[pyo3(get)]
Expand Down Expand Up @@ -196,48 +183,38 @@ impl Connection {
///
/// Anything that speaks Arrow goes in: a pandas or polars
/// DataFrame, a pyarrow Table or RecordBatchReader, or a dictionary
/// of lists. The frame is copied into the database rather than
/// scanned where it sits, because the engine has no way yet to call
/// back out to a Python object mid statement, so what a program
/// registers is a snapshot: changing the DataFrame afterwards
/// changes nothing here until it is registered again.
/// of lists. Nothing is copied. The engine is told where the
/// columns are and reads them where they lie, so registering costs
/// the same whether the frame has ten rows or ten million, and a
/// frame is a view rather than a snapshot: write into the array
/// behind it and the next statement answers what is there now. A
/// dictionary of Python lists is the one exception, since a list
/// holds objects rather than numbers and there is nothing in it to
/// point at.
///
/// Registering the same name again replaces the rows under it.
/// Registering over a table this connection did not register is
/// refused, since a frame knows nothing about the rows that were
/// already there. Gives back the number of rows written.
fn register(
slf: Py<Self>,
py: Python<'_>,
name: &str,
data: &Bound<'_, PyAny>,
) -> PyResult<usize> {
register::register(py, slf, name, data)
/// The frame belongs to this connection and goes when it does.
/// Nothing is written to the database and no other program sees it.
/// Registering the same name again replaces what it stands for,
/// columns and all. Registering over a table the database already
/// holds is refused, since a statement naming it would mean the
/// stored one. Gives back the number of rows the frame has.
fn register(&self, py: Python<'_>, name: &str, data: &Bound<'_, PyAny>) -> PyResult<usize> {
register::register(py, self, name, data)
}

/// Takes the rows of a registered frame back out.
/// Takes a registered frame's name away.
///
/// The name stops being a frame of this connection's and the rows
/// go. The empty table stays in the catalog, because no statement
/// of this engine drops one yet, so the name cannot afterwards be
/// used for something that is not a table.
/// The name stops standing for anything and the bytes go back to
/// the caller, which is not necessarily this instant: a statement
/// still reading the frame holds it until it ends.
fn unregister(&self, py: Python<'_>, name: &str) -> PyResult<()> {
register::unregister(py, self, name)
}

/// The names this connection has registered frames under, sorted.
#[getter]
fn registered(&self) -> Vec<String> {
self.frames
.lock()
.map(|frames| {
frames
.iter()
.filter(|&(_, &registered)| registered)
.map(|(name, _)| name.clone())
.collect()
})
.unwrap_or_default()
fn registered(&self, py: Python<'_>) -> PyResult<Vec<String>> {
register::registered(py, self)
}

/// Asks the statement running on this connection to stop.
Expand Down Expand Up @@ -361,7 +338,6 @@ impl Connection {
inner: Arc::new(Mutex::new(Some(opened))),
runner: OnceLock::new(),
alive: AtomicBool::new(true),
frames: Mutex::new(BTreeMap::new()),
path,
read_only,
})
Expand Down Expand Up @@ -448,35 +424,6 @@ impl Connection {
})
.map_err(|()| closed(py, "this connection"))
}

/// Whether a frame is registered here under this name right now.
pub(crate) fn registered_here(&self, name: &str) -> bool {
self.frames
.lock()
.is_ok_and(|frames| frames.get(name) == Some(&true))
}

/// Whether this connection made the table of this name, whether or
/// not a frame is registered under it now.
pub(crate) fn made_here(&self, name: &str) -> bool {
self.frames
.lock()
.is_ok_and(|frames| frames.contains_key(name))
}

pub(crate) fn remember(&self, name: &str) {
if let Ok(mut frames) = self.frames.lock() {
frames.insert(name.to_string(), true);
}
}

pub(crate) fn forget(&self, name: &str) {
if let Ok(mut frames) = self.frames.lock()
&& let Some(registered) = frames.get_mut(name)
{
*registered = false;
}
}
}

/// The rows a statement gave back.
Expand Down
Loading
Loading