Async-native message processing inspired by Dramatiq.
Project description
Fluxera
Fluxera is an async-native Python task runtime inspired by Dramatiq.
It is built for workloads where a worker should keep a lot of I/O in flight without buying concurrency through large worker-thread pools, while still handling synchronous and CPU-bound work through dedicated execution lanes.
Why Fluxera
async defactors run as realasynciotasks on the worker event loop.defactors still work through a bounded thread lane.- CPU-heavy actors can be isolated in a separate process lane.
- Redis Streams is supported as an at-least-once transport with lease renewal, stale reclaim, deduplication, and idempotency primitives.
- Rolling deploys can hand off unstarted backlog between old and new worker revisions without rotating namespaces.
Status
0.1.2 is the current public alpha.
The runtime, Redis transport v2, revision management, benchmark harnesses, and release packaging are in place, but APIs may still change as the project hardens.
Install
pip install fluxera
For local Redis development:
docker compose up -d
Quick Start
import asyncio
import fluxera
broker = fluxera.RedisBroker(
"redis://127.0.0.1:6379/15",
namespace="hello-fluxera",
)
@fluxera.actor(broker=broker, queue_name="default")
async def fetch_user(user_id: str) -> None:
await asyncio.sleep(0.1)
print("fetched", user_id)
async def main() -> None:
async with fluxera.Worker(
broker,
concurrency=128,
thread_concurrency=16,
process_concurrency=4,
):
await fetch_user.send("user-123")
await broker.join(fetch_user.queue_name)
asyncio.run(main())
Worker CLI
Fluxera can now start workers directly from the CLI, similar to the way
projects used dramatiq ... before.
If your project has:
- a setup module that creates and registers a broker
- a worker registry that lists actor modules
you can run it like this:
fluxera worker \
your_project.fluxera_setup \
--module-registry your_project.worker_registry:WORKER_MODULES \
--broker your_project.fluxera_setup:broker \
--uvloop \
--concurrency 64 \
--thread-concurrency 8
For smoke tests and one-shot local runs, add --exit-when-idle.
Producer APIs
Fluxera separates message production from message execution.
actor.send(...)andactor.send_sync(...)produce messages- worker execution lanes consume and execute those messages later
This is why send_sync() is not the same thing as the worker thread lane.
For RedisBroker, send_sync() now uses a real synchronous Redis producer
path. It is safe for normal blocking contexts such as schedulers, CLI tools,
and plain threads, but it still intentionally rejects calls made from inside an
already-running event loop.
Execution Model
Fluxera has three execution lanes:
async: default forasync defactorsthread: default for regulardefactorsprocess: opt-in for CPU-heavy actors
The process lane defaults to spawn for safe multithreaded startup. You can still override it through Worker(process_start_method=...) or FLUXERA_PROCESS_START_METHOD when needed.
Example CPU actor:
import fluxera
broker = fluxera.RedisBroker("redis://127.0.0.1:6379/15", namespace="cpu-example")
def score_document(text: str) -> int:
return sum(ord(ch) for ch in text)
score_document_actor = fluxera.actor(
broker=broker,
actor_name="score_document",
queue_name="cpu",
execution="process",
)(score_document)
Serving Revision Admin
Fluxera keeps namespace as the broker identity boundary and uses worker_revision and serving_revision for rollout control.
Read the current serving revision:
fluxera revision get \
--redis-url redis://127.0.0.1:6379/15 \
--namespace hello-fluxera \
--queue default
Promote a new serving revision with a CAS guard:
fluxera revision promote \
--redis-url redis://127.0.0.1:6379/15 \
--namespace hello-fluxera \
--queue default \
--revision 20260329153000 \
--expected-revision 20260329140000
Use --format json when the command is called by deployment automation.
Fluxera does not auto-promote serving_revision at worker startup by default.
That is intentional: startup only proves a worker booted, not that the new
revision should already receive queue traffic. Simple deployments may still
choose to auto-promote in their entrypoint or release automation.
Runtime Monitoring And Admin Dashboard
Fluxera includes runtime monitoring commands and a lightweight /admin dashboard.
Get a JSON snapshot of worker and queue state:
fluxera monitor snapshot \
--redis-url redis://127.0.0.1:6379/15 \
--namespace hello-fluxera \
--format json
Run a local admin dashboard:
fluxera monitor serve \
--redis-url redis://127.0.0.1:6379/15 \
--namespace hello-fluxera \
--host 0.0.0.0 \
--port 8090
Dashboard endpoints:
/admin: auto-refreshing queue/worker dashboard/admin/snapshot: full runtime JSON payload/healthz: health check for probes and deployment gates
The runtime snapshot includes:
- online/stale workers, revision, queue acceptance state
- queue backlog (
stream_ready,delayed) and pending deliveries pending_stalecount (idle deliveries likely stuck)waiting_not_runningcount for requests received but not yet executing
If you already run a root ASGI server (FastAPI/Starlette), mount Fluxera admin into the existing router instead of using a separate port:
from fastapi import FastAPI
import fluxera
app = FastAPI()
fluxera.mount_admin_asgi(
app,
mount_path="/admin/fluxera",
redis_url="redis://127.0.0.1:6379/15",
namespace="hello-fluxera",
)
Mounted endpoints become:
/admin/fluxera//admin/fluxera/snapshot/admin/fluxera/healthz
Delivery Semantics
- Transport delivery is at-least-once.
- Deduplication is an enqueue-time admission policy, not exactly-once execution.
- Effectively-once side effects require idempotency keys or application-level dedupe.
- Redis workers renew leases for long-running tasks and reclaim stale pending deliveries.
Retry And Callbacks
Fluxera now supports Dramatiq-compatible retry controls plus async-native callbacks.
Retry example:
import fluxera
broker = fluxera.RedisBroker("redis://127.0.0.1:6379/15", namespace="retry-example")
def retry_when(attempt: int, exc: BaseException, record: fluxera.TaskRecord) -> bool:
del record
return attempt < 3 and not isinstance(exc, ValueError)
@fluxera.actor(
broker=broker,
queue_name="default",
max_retries=5,
min_backoff=15_000,
max_backoff=300_000,
jitter="full",
retry_when=retry_when,
throws=(ValueError,),
)
async def fetch_report(report_id: str) -> None:
raise RuntimeError(f"temporary upstream failure for {report_id}")
Callback example:
import fluxera
broker = fluxera.RedisBroker("redis://127.0.0.1:6379/15", namespace="callback-example")
def on_success(context: fluxera.OutcomeContext) -> None:
print("finished", context.actor_name, context.message.message_id, context.result)
async def on_failure(context: fluxera.OutcomeContext) -> None:
print("failed", context.failure_kind, context.exception_type)
@fluxera.actor(broker=broker, queue_name="callbacks")
async def notify_exhausted(payload: dict[str, object]) -> None:
print("retry exhausted", payload["dead_letter_id"])
@fluxera.actor(
broker=broker,
queue_name="default",
on_success=on_success,
on_failure=on_failure,
on_retry_exhausted="notify_exhausted",
)
async def generate_summary() -> dict[str, str]:
return {"status": "ok"}
Rules of thumb:
throwsskips retries and goes straight to terminal handlingretry_whenoverrides the simpler retry count ruleon_successandon_failurecan be sync or async callableson_retry_exhaustedandon_dead_letteredcan be callables or actor names- actor callbacks receive JSON-safe payloads, not raw exception objects
Distributed Concurrency Limits
Fluxera now ships a Redis-backed ConcurrentRateLimiter for application-level
distributed mutexes and small concurrency caps.
import redis
import fluxera
client = redis.Redis.from_url("redis://127.0.0.1:6379/15")
limiter = fluxera.ConcurrentRateLimiter(client, "report:123", limit=1)
with limiter.acquire(raise_on_failure=False) as acquired:
if not acquired:
return
print("exclusive section")
Use aacquire() when the limiter is created from an async Redis client or a
fluxera.RedisBroker.
The default limiter TTL is now aligned with the previous production wrappers:
2 hours, or WORKER_CONCURRENCY_LOCK_TTL_MS when that environment variable
is set.
The limiter can also be used from the CLI.
Probe a key without holding it:
fluxera rate-limit probe \
--redis-url redis://127.0.0.1:6379/15 \
--key report:123 \
--format json
Run a command under a distributed mutex:
fluxera rate-limit run \
--redis-url redis://127.0.0.1:6379/15 \
--key report:123 \
-- python3 scripts/generate_report.py
Benchmark Snapshot
Latest local measurements were taken on 2026-03-29 on macOS 26.3.1, Python 3.12.10, Apple M5 Pro (15 cores).
Benchmark label legend:
c=: Fluxera worker concurrency setting used by the benchmark runnert=: Dramatiqworker_threads
Headline results against the current local Dramatiq checkout:
| Scenario | Fluxera | Dramatiq | Takeaway |
|---|---|---|---|
| production-shaped async fanout | 0.258s |
0.385s (t=8) / 0.319s (t=32) |
Fluxera is faster with 2 threads instead of 12 or 36 |
| single-worker CPU-bound | 1.270s |
3.893s (t=8) / 3.704s (t=32) |
process lane still gives Fluxera a large single-worker win |
| mixed long I/O + short work | short_drain=0.040s |
6.038s (t=8) / 0.098s (t=32) |
long I/O does not starve short work |
| Redis mixed long/short | wall=1.542s, short_drain=0.089s |
3.125s, 1.660s (t=8) / 1.681s, 0.094s (t=32) |
transport advantage remains on real Redis |
See BENCHMARK.md for the full methodology and numbers.
One nuance matters: with the safer default spawn process policy, cluster-scale CPU throughput is no longer universally faster than Dramatiq. Fluxera's strongest advantage is still async-heavy and mixed I/O workloads.
Verification
The current release candidate was checked with:
python3 -m unittest discover -s tests -vpython3 benchmarks/production_compare.py --profile smokepython3 benchmarks/redis_transport_compare.py --repeat 5 --long-io-secs 1.5/tmp/fluxera-release-venv/bin/python -m build --sdist --wheel/tmp/fluxera-release-venv/bin/python -m twine check dist/*
Documentation
- Getting Started
- FAQ
- Benchmark Results
- Dead Letter and Retry
- Revision Management
- System Design
- Deduplication and Idempotency
- Redis Lua Contract
Current Limits
- public APIs may still change during the alpha period
- result backends are not implemented yet
- message registry garbage collection is still intentionally simple
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 Distribution
Built Distribution
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 fluxera-0.1.2.tar.gz.
File metadata
- Download URL: fluxera-0.1.2.tar.gz
- Upload date:
- Size: 66.0 kB
- Tags: Source
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
4bbfc12444436ef3cf762d7044acc60a717f1a1f310273837ef32d331231b2f8
|
|
| MD5 |
f02087707619b647a67bef161671533c
|
|
| BLAKE2b-256 |
4441688181f780fb90c8f0ab2f8528528c5284a8c896827835280ba9d76239ea
|
File details
Details for the file fluxera-0.1.2-py3-none-any.whl.
File metadata
- Download URL: fluxera-0.1.2-py3-none-any.whl
- Upload date:
- Size: 58.7 kB
- Tags: Python 3
- Uploaded using Trusted Publishing? No
- Uploaded via: twine/6.2.0 CPython/3.12.10
File hashes
| Algorithm | Hash digest | |
|---|---|---|
| SHA256 |
28947f105904e5bc3ff6f941b18d5620410081f6a1727859cf7e8d8713e630f9
|
|
| MD5 |
229131acbb3da8419c2fee0bca4541a0
|
|
| BLAKE2b-256 |
3b792d9c0b4842615e9340d4768a9b4160a43c64ce09010ba72033d701a14488
|