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
19 changes: 18 additions & 1 deletion README.md
Original file line number Diff line number Diff line change
Expand Up @@ -93,6 +93,23 @@ An appender is refused inside a transaction. Its batches are commits of their ow

The wrapper costs 5 microseconds for an empty transaction, so what it costs is what the engine charges. On this machine that is more rather than less: 200 `INSERT`s cost 2.2 seconds each committing on its own and 3.3 seconds inside one transaction, and reads cost the same either way. A transaction here is worth taking for the span it holds and not for the time it saves, and the v0 write path is where that number has to change.

## Bringing a DataFrame in

A frame a program already has becomes something a statement can match on, under a name the program picks.

```python
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.

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 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.

`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.

## Reading a result as columns

A result is rows to iterate and columns to hand to something else. The columns go out over the Arrow C Data Interface, so pyarrow, pandas and polars each read the same buffers and none of them gets a Python object per cell.
Expand Down Expand Up @@ -135,7 +152,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, 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. `register` 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, 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
10 changes: 10 additions & 0 deletions python/zudb/_zudb.pyi
Original file line number Diff line number Diff line change
Expand Up @@ -86,6 +86,16 @@ class Connection:
def appender(self, table: str) -> Appender:
"""Opens an appender on `table`, for loading rows into a database that already exists."""

def register(self, name: str, data: Any) -> int:
"""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."""

@property
def registered(self) -> list[str]:
"""The names this connection has registered frames under, sorted."""

def close(self) -> None:
"""Closes the connection and frees what it held."""

Expand Down
22 changes: 22 additions & 0 deletions src/buffer.rs
Original file line number Diff line number Diff line change
Expand Up @@ -21,6 +21,7 @@ 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 @@ -285,6 +286,27 @@ 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
211 changes: 175 additions & 36 deletions src/conn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -7,6 +7,7 @@
//! 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 All @@ -21,6 +22,7 @@ use crate::appender::Appender;
use crate::columns;
use crate::error::{closed, programming, to_py_err};
use crate::interrupt;
use crate::register;
use crate::txn::Transaction;
use crate::value::{Names, from_py, to_py};

Expand Down Expand Up @@ -66,6 +68,18 @@ 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 All @@ -87,36 +101,7 @@ impl Connection {
statement: &str,
params: Option<&Bound<'_, PyDict>>,
) -> PyResult<Result> {
let params = bind(params)?;
// Owned rather than borrowed, both of them, because the thread
// that runs the statement is not this one and the borrow would
// have to say so. It is two allocations against a statement,
// which is nothing beside the parse it is about to have.
let statement = statement.to_string();
// The GIL goes down for the whole statement, waiting for the
// connection's own lock included. That is the point of a
// compiled engine in a Python process: another thread runs
// while this one is inside the executor. Where a `Ctrl-C` can
// arrive the statement goes on the connection's own thread and
// this one waits for it, which is the only way a press is felt
// before the statement ends.
let (result, names) =
interrupt::watched(py, &self.runner, &self.inner, &self.stop, move |conn| {
let borrowed: Vec<(&str, Value)> = params
.iter()
.map(|(name, value)| (name.as_str(), value.clone()))
.collect();
let names = Names::of(conn.session_mut().catalog());
conn.query_with(&statement, &borrowed)
.map(|result| (result, names))
})
.map_err(|stopped| stopped.raise(py))?
.map_err(|err| to_py_err(py, err))?;
Ok(Result {
result,
names,
next: Mutex::new(0),
})
self.query(py, statement, bind(params)?)
}

/// The same call, named for the way it reads in a notebook:
Expand Down Expand Up @@ -164,7 +149,7 @@ impl Connection {
/// statement, since the answer belongs to the session and the
/// session is what the statement is holding.
#[getter]
fn in_transaction(&self, py: Python<'_>) -> PyResult<bool> {
pub(crate) fn in_transaction(&self, py: Python<'_>) -> PyResult<bool> {
if !self.alive.load(Ordering::Acquire) {
return Err(closed(py, "this connection"));
}
Expand Down Expand Up @@ -202,6 +187,59 @@ impl Connection {
Appender::open(py, slf, table)
}

/// Puts a DataFrame under a name a statement can match on.
///
/// ```python
/// conn.register("people", frame)
/// conn.execute("MATCH (p:people) WHERE p.age > 40 RETURN p.name AS name")
/// ```
///
/// 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.
///
/// 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)
}

/// Takes the rows of a registered frame back out.
///
/// 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.
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()
}

/// Asks the statement running on this connection to stop.
///
/// The one call meant to be made from another thread while the
Expand Down Expand Up @@ -323,20 +361,121 @@ impl Connection {
inner: Arc::new(Mutex::new(Some(opened))),
runner: OnceLock::new(),
alive: AtomicBool::new(true),
frames: Mutex::new(BTreeMap::new()),
path,
read_only,
})
}

/// Runs one statement, with its parameters already engine values.
///
/// The one path every statement takes, whether the parameters came
/// from a caller's dictionary or were built here.
pub(crate) fn query(
&self,
py: Python<'_>,
statement: &str,
params: Vec<(String, Value)>,
) -> PyResult<Result> {
// Owned rather than borrowed, because the thread that runs the
// statement is not this one and the borrow would have to say
// so. It is one allocation against a statement, which is
// nothing beside the parse it is about to have.
let statement = statement.to_string();
// The GIL goes down for the whole statement, waiting for the
// connection's own lock included. That is the point of a
// compiled engine in a Python process: another thread runs
// while this one is inside the executor. Where a `Ctrl-C` can
// arrive the statement goes on the connection's own thread and
// this one waits for it, which is the only way a press is felt
// before the statement ends.
let (result, names) =
interrupt::watched(py, &self.runner, &self.inner, &self.stop, move |conn| {
let borrowed: Vec<(&str, Value)> = params
.iter()
.map(|(name, value)| (name.as_str(), value.clone()))
.collect();
let names = Names::of(conn.session_mut().catalog());
conn.query_with(&statement, &borrowed)
.map(|result| (result, names))
})
.map_err(|stopped| stopped.raise(py))?
.map_err(|err| to_py_err(py, err))?;
Ok(Result {
result,
names,
next: Mutex::new(0),
})
}

/// Runs a statement that takes no parameters and gives back
/// nothing, which is what the three transaction words are.
///
/// It goes through `execute` rather than around it, so a `COMMIT`
/// waits for the connection's lock, releases the GIL and feels a
/// `Ctrl-C` exactly as every other statement does. The result is
/// dropped, since the words return no columns.
/// It goes through the same path as every other statement, so a
/// `COMMIT` waits for the connection's lock, releases the GIL and
/// feels a `Ctrl-C` exactly as they do. The result is dropped,
/// since the words return no columns.
pub(crate) fn run(&self, py: Python<'_>, statement: &str) -> PyResult<()> {
self.execute(py, statement, None).map(|_| ())
self.query(py, statement, Vec::new()).map(|_| ())
}

/// Work done on the engine connection itself, with the GIL down.
///
/// For the things a statement cannot say: reading the catalog,
/// opening the engine's own appender. The two failures are kept
/// apart because they belong to different people. A connection that
/// is closed or locked is this call's own answer and comes back as
/// the outer error; anything the work itself decides is wrong comes
/// back as the inner one, for the caller to turn into an exception
/// now that there is an interpreter to build one with.
pub(crate) fn engine<T, S, F>(
&self,
py: Python<'_>,
work: F,
) -> PyResult<std::result::Result<T, S>>
where
T: Send,
S: Send,
F: FnOnce(&mut zudb::Connection) -> std::result::Result<T, S> + Send,
{
if !self.alive.load(Ordering::Acquire) {
return Err(closed(py, "this connection"));
}
py.detach(|| {
let mut held = self.inner.lock().map_err(|_| ())?;
let conn = held.as_mut().ok_or(())?;
Ok(work(conn))
})
.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;
}
}
}

Expand Down
Loading
Loading