Embedded streaming-batch SQL engine. Same query, static Parquet or live Kafka. No cluster.
Problem. Streaming SQL today forces a choice between heavy clusters (Flink, Spark Structured Streaming) and complete databases (Materialize, RisingWave). Embedded batch is solved — DuckDB and Polars own that lane. But there is no DuckDB-equivalent for incremental computation: a Python developer who wants to run the same SQL over a static Parquet file and a live Kafka topic still has to stand up Flink or buy into a separate database.
Solution. Riverbed is a single embeddable wheel (pip install riverbed) that runs SQL queries with two execution modes — batch (one-shot) and incremental (continuous, sub-second latency, exactly-once on supported sinks). The same query, the same Arrow output schema, the same Python API. State is local (RocksDB) or remote (S3) but never a separate service.
Who benefits.
- Python data teams who find Flink overkill for what is really a "rolling aggregation over a few topics" problem.
- ML feature stores that need online + offline parity without two codebases.
- Edge / IoT and small-fleet observability where a JVM cluster is not viable.
- Notebook-first researchers (incl. astronomy alert brokers) who want to prototype on Parquet, deploy on Kafka, and not rewrite the query.
💡 Design rationale. The bet is that Apache Arrow + DataFusion now handle ~80% of what a streaming engine needs (planner, vectorized operators, expressions, file formats), and the remaining 20% — incremental operator semantics, watermarks, checkpointing, exactly-once sinks — is a self-contained custom layer. Arroyo proved this technically; Cloudflare's acquisition of Arroyo proved it commercially. The embedded form factor is the open lane.
┌────────────────────────────────────────────────────────────────────┐
│ Python API / CLI / Server (HTTP+ADBC) │
└──────────────────────────────┬─────────────────────────────────────┘
│
┌────────────▼────────────┐
│ SQL Frontend │ ← DataFusion sqlparser
│ (CREATE STREAM ...) │
└────────────┬────────────┘
│
┌────────────▼────────────┐
│ Logical Planner │ ← DataFusion + Riverbed extensions
│ (watermark, window, │ (UDFs for STREAM, EMIT, WATERMARK)
│ emit-policy) │
└────────────┬────────────┘
│
┌────────────▼────────────┐
│ Mode Selector │ batch vs. incremental
└────┬───────────────┬────┘
│ │
┌──────────────▼─┐ ┌──────▼──────────────────────────┐
│ Batch Runner │ │ Incremental Runner │
│ (DataFusion │ │ (Riverbed delta operators │
│ physical plan)│ │ over Arrow RecordBatch deltas; │
│ │ │ state in RocksDB / S3) │
└──────┬─────────┘ └──────┬──────────────────────────┘
│ │
│ ┌────────▼─────────┐
│ │ Checkpointer │ Chandy-Lamport
│ │ (per-source │ barriers, async
│ │ offset commits) │ to S3/local
│ └────────┬─────────┘
│ │
┌──────▼──────────────────────▼──────────────────────────────┐
│ Connector Layer (TableProvider / SinkProvider traits) │
│ Sources: Parquet · Iceberg · Delta · Kafka · Pulsar · │
│ Postgres CDC (logical replication) · Arrow IPC │
│ Sinks: Parquet · Iceberg (CoW) · Kafka · Postgres · │
│ Arrow Flight · Lance │
└────────────────────────────────────────────────────────────┘
Key invariants.
- Arrow
RecordBatchis the only data type crossing operator boundaries. - Every operator has a
process_batch(batch mode) andprocess_delta(incremental mode) implementation. The compiler picks one based on the source. - State is keyed by
(operator_id, partition, key)and lives in a pluggableStateStoretrait. - Watermarks travel as Arrow schema metadata on RecordBatches — no separate control plane.
| Layer | Chosen | Stars (Apr 2026) | Why chosen | Rejected |
|---|---|---|---|---|
| Core language | Rust 1.78+ | — | Memory safety, perf, cargo, the de facto language for new data infra | C++ (build complexity), Zig (immature ecosystem) |
| Query engine | Apache DataFusion | ~7k (apache/datafusion) | Library-first; proven streaming use by Arroyo, InfluxDB, ParadeDB; Arrow-native; ASF governance | DuckDB (C++, "complete DB" model, harder to embed as a library), Velox (no SQL frontend) |
| Memory format | Apache Arrow | ~14k (apache/arrow-rs) | Zero-copy across Python/Rust, columnar, SIMD, FFI ABI | Custom format (loses ecosystem), Polars-internal (locks in) |
| Incremental layer | Custom over DataFusion physical operators | — | DBSP is the right theory but heavyweight to embed; Arroyo proved DataFusion physical operators are reusable for streaming with reasonable effort | DBSP/Feldera direct (heavier dep, GPL'd compiler), differential-dataflow (timely is overkill for single-node) |
| State store (local) | RocksDB via rocksdb crate |
~6k (rust-rocksdb) | Battle-tested, log-structured, snapshot-friendly | Sled (single maintainer, less mature), redb (newer, fewer features) |
| State store (remote) | S3 + manifest checkpoints | — | Cheap, durable, matches lakehouse model | DynamoDB (cost), FoundationDB (op overhead) |
| Python bindings | PyO3 + maturin | ~13k / ~5k | Standard for Rust→Python; zero-copy Arrow via pyo3-arrow |
cffi (slower, no Arrow FFI), nanobind (C++) |
| SQL parser | datafusion-sqlparser-rs |
~3k | Already a DataFusion dep; battle-tested; permissive license | Custom (don't) |
| Kafka client | rdkafka |
~2k | librdkafka bindings; exactly-once protocol support | kafka-rust (pure Rust but lags on transactional protocol) |
| Iceberg | iceberg-rust |
~1k (apache/iceberg-rust) | Apache project, Arrow-native; pyiceberg interop | delta-rs only (Iceberg has wider lakehouse traction) |
| CLI framework | clap v4 |
~16k | Standard | argh (less featureful) |
| Async runtime | tokio |
~28k | Standard, DataFusion uses it | smol (smaller community) |
| Test framework | cargo test + insta + proptest |
— | Snapshot + property tests catch operator regressions | pytest only (Rust core needs Rust tests) |
Sources cited:
- Bauplan: Duck Hunt — moving from DuckDB to DataFusion — production validation of DataFusion as embedded engine.
- Arroyo: We built a new SQL engine on Arrow and DataFusion — proves DataFusion physical operators are streaming-capable with adaptation.
- DBSP paper, VLDB '23 (best paper) — formalism we draw on for incremental operator design (without taking the dep).
- Composable query engines with Polars and DataFusion (ThinhDA, Nov 2025) — landscape of embeddable analytics in 2026.
riverbed/
├── crates/
│ ├── riverbed-core/ # Rust: planner extensions, mode selector, watermarks
│ ├── riverbed-incremental/ # Rust: delta operators (join, agg, window, distinct)
│ ├── riverbed-state/ # Rust: StateStore trait + RocksDB + S3 impls
│ ├── riverbed-connectors/ # Rust: Kafka, Postgres CDC, Iceberg, Parquet, Delta, Lance
│ ├── riverbed-checkpoint/ # Rust: Chandy-Lamport barriers, manifest writers
│ ├── riverbed-server/ # Rust: optional HTTP + Arrow Flight server
│ ├── riverbed-cli/ # Rust: `riverbed` binary
│ └── riverbed-py/ # Rust: PyO3 bindings → builds the `riverbed` wheel
├── python/
│ └── riverbed/
│ ├── __init__.py # public API: Session, Stream, Table
│ ├── _native.pyi # type stubs for the .so
│ └── connectors/ # thin Python wrappers (pandas/polars converters)
├── tests/
│ ├── batch/ # SQL-on-Parquet correctness vs DuckDB golden
│ ├── incremental/ # Operator-level: insert→retract→insert sequences
│ ├── exactly_once/ # Kill-recover tests for Kafka↔Postgres pipelines
│ └── property/ # proptest: batch result == replay of incremental
├── benches/ # criterion benches: TPC-H, NEXMark
├── examples/
│ ├── parquet_to_kafka.py
│ ├── kafka_join_postgres_cdc.py
│ ├── astronomy_alert_filter.py # GCN/Kafka → ranked candidate stream
│ └── feature_store_online_offline.py
├── docs/
│ ├── architecture.md
│ ├── sql-reference.md
│ ├── operator-semantics.md # which ops are incremental, which are blocking
│ └── adr/ # Architecture Decision Records
├── AGENTS.md
├── CONTRIBUTING.md
├── DEVELOPER.md
├── RELEASE.md
├── LICENSE # Apache-2.0
├── Cargo.toml # workspace
├── pyproject.toml # maturin
└── Makefile
| Requirement | Version | Notes |
|---|---|---|
| Rust | 1.78+ | rustup install stable |
| Python | 3.11+ | for the wheel |
uv |
latest | curl -LsSf https://astral.sh/uv/install.sh | sh |
maturin |
1.7+ | uv tool install maturin |
| Docker | 24+ | only for integration tests (Kafka, Postgres) |
protoc |
3.20+ | Arrow Flight |
import riverbed as rb
session = rb.Session()
# --- BATCH: same as DuckDB / Polars ---
df = session.sql("SELECT symbol, AVG(price) FROM 's3://bucket/trades/*.parquet' GROUP BY symbol").to_arrow()
# --- INCREMENTAL: same SQL, live source ---
session.create_stream(
"trades",
source=rb.Kafka(
brokers="localhost:9092",
topic="trades",
format="json",
watermark="ts - INTERVAL '5 seconds'",
),
)
q = session.sql("""
SELECT
symbol,
TUMBLE_START(ts, INTERVAL '1 minute') AS window_start,
AVG(price) AS avg_price
FROM trades
GROUP BY symbol, TUMBLE(ts, INTERVAL '1 minute')
EMIT ON WATERMARK
""")
# Run forever, sink to Iceberg
q.sink(rb.Iceberg(catalog="glue", table="analytics.trades_1m")).run()CLI:
riverbed run pipeline.sql --checkpoint-dir s3://my-bucket/ckpt/
riverbed inspect --pipeline-id abc123| Env var | Type | Default | Required | Purpose |
|---|---|---|---|---|
RIVERBED_STATE_DIR |
path | ./.riverbed-state |
no | Local RocksDB path |
RIVERBED_CHECKPOINT_URL |
URL | — | no (req. for incremental + durability) | s3://... or file://... |
RIVERBED_CHECKPOINT_INTERVAL_MS |
int | 30000 |
no | Checkpoint cadence |
RIVERBED_MAX_WATERMARK_LAG_MS |
int | 60000 |
no | Drop late events beyond this |
RIVERBED_PARALLELISM |
int | num_cpus |
no | Operator parallelism |
RIVERBED_LOG_LEVEL |
enum | info |
no | trace / debug / info / warn / error |
KAFKA_BROKERS |
csv | — | only if Kafka source | |
AWS_REGION |
str | — | only if S3/Iceberg | |
RIVERBED_TELEMETRY |
bool | false |
no | Anonymous usage stats (opt-in) |
class Session:
def __init__(self, *, config: dict | None = None) -> None: ...
def sql(self, query: str) -> Query: ...
def create_stream(self, name: str, *, source: Source, schema: pa.Schema | None = None) -> None: ...
def create_table(self, name: str, *, source: Source) -> None: ...
def register_udf(self, fn: Callable, *, name: str, return_type: pa.DataType) -> None: ...
class Query:
def to_arrow(self) -> pa.Table: ... # batch only
def to_polars(self) -> "pl.DataFrame": ... # batch only
def explain(self, *, mode: Literal["batch", "incremental", "auto"] = "auto") -> str: ...
def sink(self, sink: Sink) -> Pipeline: ... # incremental
def subscribe(self) -> Iterator[pa.RecordBatch]: ... # incremental, in-process
class Pipeline:
def run(self, *, blocking: bool = True) -> PipelineHandle: ...
def stop(self) -> None: ...| Construct | Purpose |
|---|---|
CREATE STREAM |
Source registration with watermark expression |
TUMBLE(col, INTERVAL ...) / HOP(...) / SESSION(...) |
Windowing |
EMIT ON WATERMARK / EMIT IMMEDIATE / EMIT EVERY '5 seconds' |
Output policy |
MATCH_RECOGNIZE |
Pattern detection (v1 limited subset) |
LATERAL JOIN ... AS OF SYSTEM TIME |
Versioned joins (Iceberg time travel) |
| Operator | Batch | Incremental | Notes |
|---|---|---|---|
| Projection / Filter | ✅ | ✅ | Stateless |
| Hash Join (inner, left, right) | ✅ | ✅ | Symmetric hash, retract-aware |
| Hash Join (full outer) | ✅ | Requires retraction propagation | |
| Aggregation (group-by) | ✅ | ✅ | Differential update |
| Window aggregation (tumble/hop) | ✅ | ✅ | Watermark-driven emit |
| Window aggregation (session) | ✅ | Merging state is hard | |
| TopK | ✅ | ✅ | Insert/retract via heap |
| DISTINCT | ✅ | ✅ | Count-tracked Z-set |
| Recursive CTE | ✅ | ❌ | Out of scope v1 |
| Error class | When | Behavior |
|---|---|---|
RiverbedSchemaError |
SQL refers to unknown column or types don't unify | Raise at planning, before run |
RiverbedSourceError |
Connector cannot read (auth, missing file, broker down) | Retry with backoff; after N attempts, fail pipeline |
RiverbedStateError |
RocksDB corruption / S3 manifest missing | Fail-fast; pipeline must be restarted from prior checkpoint |
RiverbedWatermarkError |
Watermark went backwards | Log + drop event (configurable to fail) |
RiverbedCheckpointError |
Checkpoint write failed | Pipeline pauses, retries; after N failures, terminate |
RiverbedExactlyOnceError |
Sink does not support 2PC but pipeline declared exactly_once=True |
Refuse to start |
Retry policy: exponential backoff, jitter, max 5 attempts on transient errors.
make test # all
make test-batch # SQL correctness vs DuckDB golden
make test-incr # operator-level delta tests
make test-eo # exactly-once kill-recover (docker compose required)
make test-prop # property: batch(query, data) == replay_incremental(query, data)
make bench # criterion benches; outputs to target/criterion/The property test is the core correctness guarantee: for any supported query and any partition of the input into a sequence of batches, running the query in incremental mode and emitting the final state must equal running the query in batch mode on the union. This is the contract.
- Distributed execution. v1 is single-node. Multi-node is v2 via Arrow Flight + a coordinator (likely Ballista-style). Don't try to build it now.
- A storage format. Riverbed reads/writes Parquet, Iceberg, Delta, Lance — it does not invent its own.
- A catalog. Use Iceberg's catalog or pass Arrow tables directly.
- A web UI. CLI + Python + JSON pipeline manifests. UI is a separate project.
- Recursive SQL / Datalog. DBSP territory; out of scope.
- JVM / non-Python bindings in v1. Add later via the same Rust core.
- Replacing dbt. Riverbed is the engine; orchestration is someone else's job.
- State format compat. Should checkpoint format be versioned with a stability promise from v1, or should v1 explicitly declare "no migration guarantees"? Default: declare unstable in v1; freeze at v1.0.
- Exactly-once across heterogeneous sinks. Kafka has 2PC; Iceberg has atomic commits; Postgres has txns. Do we expose a unified
exactly_once=Trueflag, or per-sink semantics? Default: per-sink, with--strictflag that refuses non-EO sinks. - Python GIL contention. With async Kafka consumers, when does the GIL bite? Need a benchmark before committing to in-process pipelines vs subprocess. Default: in-process, document GIL-bound throughput ceiling.
- Iceberg branching integration. Should incremental writes go to a branch by default for review? Default: no, opt-in via
branch="..."inIceberg(...)sink. - Window allowed lateness vs hard watermark. What's the default? Default: 0; require explicit opt-in for late event handling.
- License. Apache-2.0 (matches DataFusion / Arrow ecosystem) vs MIT. Default: Apache-2.0.
Implement end-to-end using only this README. Resolve the relevant Open Questions before each phase touches them.
| Phase | Deliverable | Done when |
|---|---|---|
| 0 | Cargo workspace + maturin scaffold + CI (lint, test, build wheel for linux/macos/win) | cargo test + maturin develop + pytest all green on empty stubs |
| 1 | riverbed-core: Session, batch SQL via DataFusion, Parquet/Arrow IPC sources |
Session().sql("SELECT 1").to_arrow() works; TPC-H Q1, Q3, Q6 match DuckDB |
| 2 | riverbed-state: StateStore trait + RocksDB impl + S3 manifest checkpointer |
Property test: write→checkpoint→kill→restore→read returns same bytes |
| 3 | riverbed-incremental: delta operators (filter, project, hash agg, hash join, tumble window) |
All operators pass insert/retract/insert property tests vs batch oracle |
| 4 | riverbed-connectors: Kafka source/sink, Iceberg sink, Postgres CDC source |
End-to-end: Kafka→agg→Iceberg pipeline runs for 1h on NEXMark Q5 |
| 5 | riverbed-checkpoint: Chandy-Lamport barriers, async S3 commit |
Kill-recover test passes: pipeline killed mid-window resumes with no duplicates and no lost events |
| 6 | riverbed-py: PyO3 wrapper, type stubs, examples |
All 4 example scripts in examples/ run; wheel installs from PyPI test index |
| 7 | riverbed-cli + optional riverbed-server |
riverbed run pipeline.sql works; Arrow Flight endpoint streams results |
| File | Purpose | Key symbols |
|---|---|---|
crates/riverbed-core/src/session.rs |
Top-level entry; holds DataFusion SessionContext + Riverbed extensions |
pub struct Session, Session::sql, Session::create_stream |
crates/riverbed-core/src/planner.rs |
Hooks into DataFusion logical planner; injects mode selector | RiverbedPlanner, select_mode(plan) -> ExecMode |
crates/riverbed-incremental/src/operator.rs |
Trait every incremental operator implements | trait IncrementalOp { fn process_delta(&mut self, batch: &RecordBatch, sign: Sign) -> Vec<DeltaBatch>; } |
crates/riverbed-incremental/src/agg.rs |
Differential group-by | DeltaGroupBy; uses Z-set semantics |
crates/riverbed-incremental/src/join.rs |
Symmetric hash join with retraction | DeltaHashJoin |
crates/riverbed-incremental/src/window.rs |
Tumble/hop window with watermark trigger | WindowOp, Watermark |
crates/riverbed-state/src/lib.rs |
StateStore trait |
trait StateStore { get/put/range/snapshot/restore } |
crates/riverbed-checkpoint/src/barrier.rs |
Chandy-Lamport barriers in Arrow metadata | inject_barrier, align_barriers |
crates/riverbed-connectors/src/kafka.rs |
Source + transactional sink | KafkaSource, KafkaTransactionalSink |
crates/riverbed-py/src/lib.rs |
PyO3 module root | #[pymodule] fn _native(...) |
python/riverbed/__init__.py |
Public Python API; re-exports + ergonomic wrappers | Session, Stream, Table, Kafka, Iceberg, Postgres |
tests/property/batch_eq_incremental.rs |
The core correctness test | prop_batch_equals_replay |
- Rust 1.78+;
rustfmt+clippy::pedanticclean - Python 3.11+; typed signatures;
mypy --strictonpython/riverbed/ - No
unsafeoutside FFI boundaries (riverbed-py) - Secrets via env vars only; never logged
- All
RecordBatchflows must be zero-copy across the PyO3 boundary (usepyo3-arrow) - Every new operator ships with: a batch-mode equivalent, a property test, a NEXMark-derived integration test, and a benchmark
- AIDEV-* comments mandatory at:
- all
unsafeblocks - all panicking paths
- all places where watermark or checkpoint correctness is non-obvious
- all
-
make testpasses; coverage ≥ 80% onriverbed-incrementalandriverbed-checkpoint -
make benchproduces NEXMark Q1, Q5, Q8 results; document them indocs/benchmarks.md - Property test
batch == replay(incremental)runs ≥ 10k cases per operator without failure - Kill-recover test: SIGKILL the process mid-pipeline 100 times — zero duplicate or lost events on Kafka→Iceberg path
- All 4 example scripts run end-to-end on a clean machine
- PyPI wheel builds for cp311/cp312/cp313 on linux-x86_64, linux-aarch64, macos-arm64, win-amd64
- All Open Questions resolved or explicitly deferred to v1.1 in
docs/adr/ -
AGENTS.mdpopulated using https://github.com/ejoliet/claude-skills/blob/main/AGENTS.md as base
- Resolve OQ #2 and #6 (exactly-once semantics + license) — both block public API shape.
- Spike Phase 0–1 in 1 week: prove
Session().sql(parquet) -> Arrowworks and the wheel builds across platforms. If the maturin matrix fights back, restructure now. - Build the property-test harness first (Phase 3 dep, but write it now). The
batch == replay(incremental)invariant is the entire technical bet — the moment you can run 10k random cases against a stub, every operator becomes a fill-in-the-blank. - Pick one killer demo — recommend astronomy alert broker (GCN Kafka → tumble window → ranked output → Iceberg) because it (a) has a real user (you), (b) hits all the hard parts (Kafka, watermarks, late events, lakehouse sink), and (c) makes for a great launch blog post.
- Open OSS repo at v0.1 the moment Phase 1 is green. Visibility compounds; do not wait for v1.0.
- Apply for an Apache Incubator slot when v0.3 has 3+ external contributors. Governance matters for the "next k8s" framing.