Fast, stateful, ephemeral stream functions with a Rust core
Project description
velo ⚡
Stateful stream processing without the infrastructure overhead.
pip install velo-stream
The problem
For one stream, a Python variable is fine. No argument there.
prev_frame = None
def process_frame(frame):
global prev_frame
diff = compare(frame, prev_frame)
prev_frame = frame
return diff
The problem starts when you have many streams simultaneously — hundreds of users, sessions, devices, clips — each needing their own isolated state.
What happens when you scale with a dict
# Step 1: one dict per state variable
prev_frames = {}
def process_frame(user_id, frame):
diff = compare(frame, prev_frames.get(user_id))
prev_frames[user_id] = frame
return diff
Fine. Now the problems:
When do you delete prev_frames[user_id]? The user disconnected. Or did they time out? Or crash? prev_frames grows forever. Memory leak.
# Step 2: add timeout cleanup
last_seen = {}
def cleanup():
stale = [k for k, v in last_seen.items() if time.time() - v > 30]
for k in stale:
del prev_frames[k]
del last_seen[k]
Two events from the same user arrive simultaneously. Race condition.
# Step 3: add locks
lock = threading.Lock()
def process_frame(user_id, frame):
with lock:
diff = compare(frame, prev_frames.get(user_id))
prev_frames[user_id] = frame
last_seen[user_id] = time.time()
cleanup()
return diff
Your state is more than one variable. Now every new piece of state needs its own dict, its own cleanup entry, its own lock path.
# Step 4: five state variables = five dicts to manage
prev_frames = {}
frame_counts = {}
motion_scores = {}
last_seen = {}
alert_thresholds = {}
# all need cleanup. all need locking. all need the same boilerplate.
You've spent 50 lines building a fragile lifecycle manager instead of writing business logic.
Velo replaces all of that:
@stream_fn
async def process_frames(frames):
prev = None
count = 0
motion_scores = []
async for frame in frames:
diff = compare(frame, prev) if prev else 0
count += 1
motion_scores.append(diff)
prev = frame
yield {"frame": count, "motion": diff, "avg": sum(motion_scores) / count}
# State lifecycle is automatic.
# Stream closes → all variables are garbage collected.
# No dicts. No cleanup. No locks. No boilerplate.
One worker per stream. All state is just local variables. Worker lives exactly as long as the stream — then disappears.
What Velo is (and isn't)
Velo is: A library you run inside your existing server or container. It manages stateful worker lifecycles so you don't have to.
Velo is not: A serverless platform. Velo runs in a long-lived process — the same way your web server does. You deploy it like any other service.
Velo competes with: Flink, Faust, Bytewax — for the specific case of short-lived, bursty, stateful streams.
vs the alternatives
| Tool | Startup | Has state | Idle cost | Complexity |
|---|---|---|---|---|
| Velo | ~350μs | ✅ local vars | Near zero — workers drop on idle | Low — just write generators |
| Redis + functions | ~ms + RTT | ✅ external | Redis always running | Medium — manage keys + TTLs |
| Apache Flink | 2–10 seconds | ✅ | High — always on | High — JVM, cluster setup |
| Faust | ~seconds | ✅ | Medium — always on | Medium — Kafka required |
| Bytewax | ~ms | ✅ | Medium — always on | Medium — continuous pipeline |
Quickstart
from velo import stream_fn
# Define — just write an async generator
@stream_fn
async def running_average(events):
total, count = 0.0, 0
async for event in events:
total += event
count += 1
yield total / count
# Batch — process a list, get results
results = await running_average.run([1, 2, 3, 4, 5])
# → [1.0, 1.5, 2.0, 2.5, 3.0]
# Live stream — open, send events, receive results
async with running_average.open() as stream:
await stream.send(10)
print(await stream.recv()) # 10.0
await stream.send(20)
print(await stream.recv()) # 15.0
# Compose — chain functions with |
pipeline = normalize | running_average | alert_if_high
async with pipeline.open() as stream:
async for result in stream.feed(sensor_data):
print(result)
The API
Four exports. That's the entire public surface.
from velo import stream_fn # the decorator
from velo import Stream # type hint for stream handles
from velo import StreamMetrics # per-stream metrics
from velo import StreamConfig # optional config
@stream_fn — define a stream function
Write it exactly like a Python async generator. events is an async iterable.
@stream_fn
async def my_fn(events):
state = {} # any Python state you want
async for event in events:
state = update(state, event)
yield result(state)
.run(iterable) — batch mode
results = await my_fn.run([e1, e2, e3])
.open() — live stream mode
async with my_fn.open() as stream:
await stream.send(event)
result = await stream.recv()
# or iterate results:
async for result in stream.feed(source):
handle(result)
| — pipe composition
pipeline = fn_a | fn_b | fn_c
results = await pipeline.run(data)
Optional config
@stream_fn(
buffer=256, # events buffered before backpressure (default: 256)
timeout=30.0, # idle seconds before auto-close (default: 30.0)
max_concurrent=1000, # max parallel instances (default: 1000)
)
async def my_fn(events):
...
Real examples
Video — motion detection across frames
@stream_fn
async def detect_motion(frames):
prev = None
async for frame in frames:
if prev is not None:
diff = abs(frame.astype(int) - prev.astype(int)).mean()
yield {"frame": frame.id, "motion": diff > 5.0, "score": diff}
prev = frame
results = await detect_motion.run(video.frames())
IoT — rolling window over sensor bursts
from collections import deque
@stream_fn
async def rolling_stats(events):
window = deque(maxlen=10)
async for reading in events:
window.append(reading["value"])
yield {
"mean": sum(window) / len(window),
"min": min(window),
"max": max(window),
}
LLM — stateful token stream post-processing
import json
@stream_fn
async def extract_json(tokens):
"""Accumulate tokens until a complete JSON object forms."""
buffer, depth = "", 0
async for token in tokens:
buffer += token
depth += token.count("{") - token.count("}")
if depth == 0 and buffer.strip().startswith("{"):
yield json.loads(buffer)
buffer = ""
Session — per-user fraud detection
@stream_fn
async def detect_fraud(events):
seen_ips = set()
total_spend = 0.0
async for event in events:
seen_ips.add(event["ip"])
total_spend += event.get("amount", 0)
risk = len(seen_ips) > 3 or total_spend > 1000
yield {"event": event, "risk_score": risk}
# One stream per user — isolated state, auto cleanup on disconnect
async with detect_fraud.open() as stream:
async for result in stream.feed(user_events):
if result["risk_score"]:
flag_for_review(result)
More examples in examples/.
Performance
Benchmarked on real hardware (Rust core, crossbeam SPSC channels):
| Metric | Result |
|---|---|
| Stream startup | 0.35ms |
| Inter-event P99 latency | 0.48ms |
| 1000 concurrent streams | ✅ stable |
| Throughput (current) | ~6K ev/s |
Note on throughput: The
asyncio.to_threadbridge adds ~200μs overhead per event at the Python/Rust boundary. The Rust core itself handles >500K ev/s — the bottleneck is OS thread scheduling, not the channels. A batch API (in roadmap) will close this gap significantly. Seedocs/throughput-v2-design.mdfor the full analysis and options.
Run the benchmarks yourself:
python benchmarks/runner.py --scenario all
python benchmarks/runner.py --scenario all --format markdown
How it works
send(event)
│
▼ serialize (msgpack/pickle)
Rust crossbeam SPSC input channel ← lock-free, GIL released
│
▼
Python worker (asyncio task)
reads from Rust channel via to_thread()
│
▼
Python async generator (your code)
│
▼ result
Rust crossbeam SPSC output channel ← lock-free, GIL released
│
▼ deserialize
recv() → caller
Rust handles: stream lifecycle, channel buffering (crossbeam, lock-free), backpressure (bounded channels), concurrency limits, metrics.
Python handles: your stream function logic (async generators), serialization boundary.
For a deep dive: docs/architecture.md · docs/rust-data-path.md
Installation
From PyPI
pip install velo-stream
Wheels available for Python 3.8–3.12 on Linux (x86_64, aarch64), macOS (universal), and Windows (x86_64). No Rust required.
From source (requires Rust)
curl --proto '=https' --tlsv1.2 -sSf https://sh.rustup.rs | sh
git clone https://github.com/sahilmalik27/velo.git
cd velo
pip install maturin
maturin develop --release
Contributing
Contributions are welcome.
git clone https://github.com/sahilmalik27/velo.git
cd velo
pip install maturin && maturin develop
pip install -e ".[dev]"
pytest tests/ -v
Good first issues:
- New adapters (Kafka, Redis Streams, WebSocket)
- Batch API (
send_batch) for higher throughput - JavaScript / Node.js bindings
Rules:
- Fork → branch → change → test → PR
- All Rust changes need a before/after benchmark
- Keep the public API surface small — resist adding to the 4 exports
Project structure:
velo/
├── velo-core/ # Rust runtime (tokio, crossbeam, PyO3)
├── velo/ # Python API (@stream_fn, Stream, config)
├── benchmarks/ # Performance suite
├── examples/ # Real-world usage demos
├── docs/ # Architecture, design decisions
└── tests/ # Unit + integration
Roadmap
- Full Rust data path (crossbeam SPSC, GIL-released send/recv)
- PyPI release (
pip install velo-stream) - Multi-platform wheels via GitHub Actions
- Batch API —
send_batch([e1, e2, ...])for 10-50× throughput - Dedicated worker thread per stream — eliminate
asyncio.to_threadoverhead - Kafka adapter
- Redis Streams adapter
- Persistent state (checkpoint to disk)
- Prometheus / OpenTelemetry metrics export
License
Apache 2.0 — see LICENSE. Free for commercial use. Includes explicit patent grant.
Project details
Release history Release notifications | RSS feed
Download files
Download the file for your platform. If you're not sure which to choose, learn more about installing packages.
Source Distributions
Built Distributions
Filter files by name, interpreter, ABI, and platform.
If you're not sure about the file name format, learn more about wheel file names.
Copy a direct link to the current filters
File details
Details for the file velo_stream-0.2.0-cp38-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl.
File metadata
- Download URL: velo_stream-0.2.0-cp38-abi3-manylinux_2_17_x86_64.manylinux2014_x86_64.whl
- Upload date:
- Size: 800.7 kB
- Tags: CPython 3.8+, manylinux: glibc 2.17+ x86-64
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
1e75841ae9fe5f703be1e4f8f58f383027c9f4b2aaba24d5046339a089d641c8
|
|
| MD5 |
e6f568b45b62d5297cf363179e7d064a
|
|
| BLAKE2b-256 |
d98d2c5109073de6c27308b9a9564e7ecc71f579b43a9e2566cb7d6c3c7ee740
|
File details
Details for the file velo_stream-0.2.0-cp38-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl.
File metadata
- Download URL: velo_stream-0.2.0-cp38-abi3-manylinux_2_17_aarch64.manylinux2014_aarch64.whl
- Upload date:
- Size: 785.0 kB
- Tags: CPython 3.8+, manylinux: glibc 2.17+ ARM64
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
68764e11c731613635d56b1e6b397f5947d5b697f751353760d3768b86c93767
|
|
| MD5 |
70ec0fd4111efffda3d9d5e337f12ea1
|
|
| BLAKE2b-256 |
119fbbfd0238fdb0bf66736ffe28d4421d574ff5141edcd4b6d67c1bbf5720f4
|
File details
Details for the file velo_stream-0.2.0-cp38-abi3-macosx_10_12_x86_64.macosx_11_0_arm64.macosx_10_12_universal2.whl.
File metadata
- Download URL: velo_stream-0.2.0-cp38-abi3-macosx_10_12_x86_64.macosx_11_0_arm64.macosx_10_12_universal2.whl
- Upload date:
- Size: 1.4 MB
- Tags: CPython 3.8+, macOS 10.12+ universal2 (ARM64, x86-64), macOS 10.12+ x86-64, macOS 11.0+ ARM64
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.1.0 CPython/3.13.7
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
5bdff280c5367db42597c60a8599127be7638a39e530f21b6173613e31da415c
|
|
| MD5 |
71a7d3280a2cdcc4c5086f3fffd206cf
|
|
| BLAKE2b-256 |
88a4473922e713e1b87a9d8e3187290fb45b97f42b0bcbb837704f2a30b54490
|