Skip to content
brunolnettoPublic

About

A DuckDB-based queue manager

Resources

Stars

10 stars

Watchers

0 watching

Forks

Latest commit

Β 

History

140 Commits

Folders and files

Repository files navigation

Queuack Logo Queuack - DuckDB-Powered Job Queue & DAG Engine

Version codecov downloads

Queuack Mascot

Queuack is a pragmatic, single-node job queue that stores jobs in a DuckDB table. It’s built for dev/test and small-to-medium production workloads where you want durability without the operational overhead of Redis/RabbitMQ/Celery.

Perfect for dev/test environments and small-to-medium production workloads where you want:

  • βœ… Persistent queues without Redis/RabbitMQ complexity
  • βœ… DAG workflows without Airflow's operational overhead
  • βœ… Memory-efficient streaming for processing massive datasets
  • βœ… Beautiful visualizations with customizable Mermaid diagrams
  • βœ… Zero external dependencies (just DuckDB + stdlib)

Table of Contents


🎯 Why Queuack?

Feature Queuack Celery + Redis Airflow MLflow + Kubeflow
Setup pip install queuack Install Redis, configure Celery Docker compose, PostgreSQL, webserver K8s cluster, multiple servers
DAG Workflows βœ… Built-in ❌ Separate tools βœ… Core feature βœ… Complex setup
Streaming ETL βœ… O(1) memory ❌ Load all data ⚠️ Manual batching ⚠️ Manual setup
ML Pipelines βœ… Native support ⚠️ Manual ⚠️ Complex βœ… Core feature
Visualization βœ… Mermaid (6 themes) ❌ None βœ… Web UI (complex) βœ… Web UI (complex)
Local Development βœ… Single file ⚠️ Need Redis ⚠️ Need full stack ❌ Need K8s
Memory Footprint ~50 MB ~200 MB ~2 GB ~4 GB

Perfect for:

  • πŸ”¬ ML Engineers - Train models without Kubernetes
  • πŸ“Š Data Engineers - Build ETL pipelines without Airflow overhead
  • πŸš€ Startups - Ship fast without infrastructure complexity
  • πŸ’» Solo Developers - Full workflow engine on your laptop

πŸš€ Quick Start

Installation

pip install queuack

# Optional: For Parquet support
pip install queuack[parquet]

Basic Queue (30 seconds)

from queuack import DuckQueue, Worker

# Create queue
queue = DuckQueue("jobs.db")  # or ":memory:" for testing

# Enqueue a job
def process_data(x):
    return x * 2

job_id = queue.enqueue(process_data, args=(42,))

# Process jobs
worker = Worker(queue, concurrency=4)
worker.run()  # Blocks and processes jobs

DAG Workflow (1 minute)

from queuack import DAG, DuckQueue

queue = DuckQueue(":memory:")
dag = DAG("etl_pipeline", queue=queue)

# Define tasks
def extract():
    return {"records": 1000}

def transform(context):
    data = context.upstream("extract")
    return {"processed": data["records"] * 2}

def load(context):
    result = context.upstream("transform")
    print(f"Loaded {result['processed']} records")
    return {"status": "success"}

# Build pipeline
dag.add_node(extract, name="extract")
dag.add_node(transform, name="transform", depends_on=["extract"])
dag.add_node(load, name="load", depends_on=["transform"])

# Execute
dag.execute()

Streaming ETL (2 minutes)

from queuack import generator_task, StreamReader

# Process 1 million records with ~50MB memory
@generator_task(format="parquet")  # or csv, jsonl, pickle
def extract_data():
    for i in range(1_000_000):
        yield {"id": i, "value": i * 2}

# Returns path to Parquet file
output_path = extract_data()

# Read lazily - one row at a time!
reader = StreamReader(output_path)
for row in reader:
    process(row)  # Memory stays constant

Async I/O (1 minute)

from queuack import async_task
import asyncio

# 10-100x speedup for I/O-bound tasks
@async_task
async def fetch_data(urls):
    async with aiohttp.ClientSession() as session:
        tasks = [session.get(url) for url in urls]
        responses = await asyncio.gather(*tasks)
        return [await r.json() for r in responses]

# Call synchronously - decorator handles event loop
results = fetch_data(["url1", "url2", ..., "url50"])
# Completes in ~0.1s instead of ~5s (50x faster!)

ML Pipeline (3 minutes)

from queuack import DAG, DuckQueue

# Replace Airflow + MLflow + Kubeflow with 50 MB
queue = DuckQueue("ml_pipeline.db")
dag = DAG("model_training", queue=queue)

# Build ML pipeline
dag.add_node(ingest_data, name="ingest")
dag.add_node(validate_data, name="validate", depends_on=["ingest"])
dag.add_node(engineer_features, name="features", depends_on=["validate"])
dag.add_node(train_model, name="train", depends_on=["features"])
dag.add_node(evaluate_model, name="evaluate", depends_on=["train"])
dag.add_node(deploy_model, name="deploy", depends_on=["evaluate"])

# Execute pipeline - all state tracked in SQLite
dag.execute()

# Parallel hyperparameter search
for params in param_grid:  # 54 combinations
    queue.enqueue(train_model, args=(params,))
# Trains 4x faster with 4 workers, no Ray/Dask needed

✨ Key Features

1. Job Queue - Redis-free persistence

# Priorities, delays, retries, timeouts
queue.enqueue(
    send_email,
    args=("user@example.com",),
    priority=90,           # 0-100 (higher = sooner)
    delay_seconds=3600,    # Schedule for later
    max_attempts=5,        # Retry failed jobs
    timeout_seconds=300    # Job timeout
)

# Batch operations
job_ids = queue.enqueue_batch([
    (task1, (arg1,), {}),
    (task2, (arg2,), {}),
])

# Monitor
stats = queue.stats()
# {'pending': 42, 'claimed': 3, 'done': 1250, 'failed': 5}

2. DAG Workflows - Complex pipelines made simple

# Fan-out/fan-in pattern
dag.add_node(extract, name="extract")
dag.add_node(transform_a, name="transform_a", depends_on=["extract"])
dag.add_node(transform_b, name="transform_b", depends_on=["extract"])
dag.add_node(load, name="load", depends_on=["transform_a", "transform_b"])

# Conditional execution (ANY mode)
dag.add_node(
    validate,
    name="validate",
    depends_on=["source_a", "source_b"],
    dependency_mode="any"  # Run when ANY parent completes
)

# Sub-DAGs for reusability
preprocessing_dag = create_preprocessing_dag()
dag.add_subdag(preprocessing_dag, name="preprocess")

3. Memory-Efficient Streaming - Process billions of rows

from queuack import generator_task, StreamReader, StreamWriter

# Write generator β†’ file (O(1) memory)
@generator_task(format="parquet")
def extract():
    for i in range(100_000_000):  # 100M rows!
        yield {"id": i, "data": process(i)}

# Supports 4 formats:
# - JSONL: Human-readable, universal
# - CSV: Excel-compatible
# - Parquet: Analytics, Spark/Pandas
# - Pickle: Complex Python objects

Memory comparison:

  • Traditional: Load all β†’ 14 GB RAM ❌
  • Streaming: ~50 MB RAM βœ…

4. Beautiful Visualizations - 6 customizable themes

from queuack import MermaidColorScheme

# Pre-built themes
dark = MermaidColorScheme.dark_mode()
professional = MermaidColorScheme.blue_professional()
accessible = MermaidColorScheme.high_contrast()

# Generate diagram
mermaid = dag.export_mermaid(color_scheme=dark)

# Paste into GitHub/GitLab Markdown:
# ```mermaid
# [paste here]
# ```

Available themes: default, blue_professional, dark_mode, pastel, high_contrast, grayscale

5. Production-Ready Features

  • βœ… Backpressure control - Automatic throttling at 10k pending jobs
  • βœ… Graceful shutdown - SIGTERM/SIGINT handling
  • βœ… Dead letter queue - Failed job inspection
  • βœ… Claim recovery - Auto-recover stuck jobs
  • βœ… Multi-queue workers - Priority-based claiming
  • βœ… Concurrent execution - Thread pool workers
  • βœ… Transaction safety - ACID guarantees via DuckDB

πŸ“š Documentation & Examples

Examples Structure

Our examples follow a progressive learning path:

01_basic/ - Core Concepts

  • Simple queue operations
  • Priority and delayed jobs
  • Batch operations

02_workers/ - Worker Patterns

  • Single and concurrent workers
  • Multi-queue processing
  • Graceful shutdown

03_dag_workflows/ - DAG Patterns

  • Linear pipelines
  • Fan-out/fan-in
  • Conditional execution
  • Diamond dependencies
  • Sub-DAGs

04_advanced/ - Advanced Techniques

  • Custom backpressure
  • Monitoring dashboards
  • Distributed workers
  • Custom Mermaid color schemes

05_real_world/ - Production Use Cases

  • ETL pipelines
  • Web scraping
  • Image processing
  • ML training pipelines (see 05_real_world/03_ml/)
  • Streaming ETL (1M+ records)
  • Multi-format exports (JSONL/CSV/Parquet)
  • Async API fetching (50x faster)

06_integration/ - Framework Integration

  • Flask, FastAPI, Django
  • CLI tools

Run any example:

cd examples/05_real_world
python 01_web_scraper.py

πŸ—οΈ Architecture

Queue Storage

  • Engine: DuckDB (embedded OLAP database)
  • Schema: Single jobs table with indexes
  • Locking: File-based for multi-process safety
  • Transactions: ACID compliance for atomic operations

Job Execution

  • Serialization: Pickle (functions + arguments)
  • Concurrency: ThreadPoolExecutor per worker
  • Claim Semantics: Visibility timeout with stale recovery
  • Retry Logic: Exponential backoff (configurable)

DAG Execution

  • Graph Engine: NetworkX for topological sorting
  • Scheduling: Level-based parallel execution
  • Dependencies: ALL (default) or ANY mode
  • Status Tracking: Real-time job status monitoring

Streaming Engine

  • Memory Model: O(1) constant memory usage
  • Batch Processing: 10k row batches for Parquet
  • Format Support: JSONL, CSV, Parquet, Pickle
  • Lazy Reading: Generator-based iteration

βš™οΈ Configuration & Tuning

Queue Configuration

queue = DuckQueue(
    db_path="jobs.db",          # or ":memory:"
    default_queue="default",
    workers_num=4,               # Auto-start workers
    worker_concurrency=2,        # Threads per worker
    poll_timeout=1.0,            # Claim polling interval
    serialization="pickle"       # or "json_ref"
)

Worker Configuration

worker = Worker(
    queue,
    queues=[
        ("high_priority", 100),
        ("normal", 50),
        ("low", 10)
    ],
    concurrency=8,               # Thread pool size
    worker_id="worker-01"        # For distributed setups
)

DAG Configuration

dag = DAG(
    name="pipeline",
    queue=queue,
    max_retries=3,
    retry_delay=60,
    timeout_per_task=600,
    show_progress=True           # Progress bar
)

Backpressure Thresholds

# Customize in subclass
class MyQueue(DuckQueue):
    @classmethod
    def backpressure_warning_threshold(cls):
        return 5000  # Warn at 5k pending

    @classmethod
    def backpressure_block_threshold(cls):
        return 50000  # Block at 50k pending

πŸ”’ Security Considerations

Pickle Serialization

  • ⚠️ Not safe for untrusted input - Pickle can execute arbitrary code
  • ⚠️ Not portable across refactors - Function signature changes break old pickles
  • βœ… Fast and convenient - Works with any Python object

Mitigations:

  1. Use serialization="json_ref" mode (functions by reference only)
  2. Validate all job inputs before enqueueing
  3. Run workers in sandboxed environments
  4. Keep function signatures stable

Multi-Process Safety

  • βœ… File-based locking for concurrent workers
  • βœ… Automatic stale claim recovery
  • ⚠️ Avoid :memory: with multiple workers (use temp file instead)

πŸ“Š Performance & Scaling

Throughput Benchmarks

  • Enqueue: ~5,000 jobs/second (batch mode)
  • Claim: ~1,000 claims/second (single worker)
  • Execute: Limited by job duration + thread pool

Scaling Guidelines

Jobs/Day Workers Concurrency DB Size
< 10k 1 2-4 < 100 MB
10k - 100k 2-4 4-8 100 MB - 1 GB
100k - 1M 4-8 8-16 1-10 GB
> 1M 8+ 16+ 10+ GB

Tips:

  • Purge completed jobs regularly (queue.purge())
  • Use multiple queues for different priorities
  • Run workers on same host as DB file (avoid network file systems)
  • For CPU-bound jobs, use process-based workers
  • Monitor with queue.stats() and dag.get_progress()

πŸ§ͺ Testing

# Run all tests
pytest

# With coverage
pytest --cov=queuack --cov-report=html

# Specific test suite
pytest tests/test_dag.py -v

# Fast tests only (skip slow integration tests)
pytest -k "not test_large"

πŸ—ΊοΈ Roadmap

Completed βœ…

  • Basic queue with priorities and delays
  • DAG workflow engine
  • Generator streaming (O(1) memory)
  • Multi-format support (CSV, Parquet, JSONL, Pickle)
  • Mermaid visualization with themes
  • Sub-DAG support
  • Async/await support for I/O-heavy tasks
  • Job priorities within DAGs

In Progress 🚧

  • Terminal UI for monitoring

Planned πŸ“‹

  • Storage backend abstraction (SQLite, PostgreSQL, Redis, S3)
  • Scheduled/cron jobs
  • Dynamic DAG generation
  • Result caching
  • Job pause/resume
  • Prometheus metrics

🀝 Contributing

We welcome contributions! Please:

  1. Fork the repository
  2. Create a feature branch (git checkout -b feat/amazing-feature)
  3. Write tests for your changes
  4. Ensure tests pass (pytest)
  5. Commit your changes (git commit -m 'feat: Add amazing feature')
  6. Push to the branch (git push origin feat/amazing-feature)
  7. Open a Pull Request

Development setup:

# Clone repo
git clone https://github.com/brunolnetto/queuack.git
cd queuack

# Install in dev mode
pip install -e ".[dev]"

# Run tests
pytest -v

πŸ“„ License

MIT Β© 2025 Bruno Peixoto


πŸ™ Acknowledgments

  • Built with DuckDB - Fast in-process analytical database
  • DAG engine inspired by Airflow, simplified for single-node use
  • Visualization powered by Mermaid

πŸ“ž Support


Made with πŸ¦† and ❀️ by the Queuack team

About

A DuckDB-based queue manager

Resources

Stars

10 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages