# How Flowsint Uses Celery for Task Distribution: Architecture and Implementation

> Discover how Flowsint employs Celery for efficient task distribution. Learn its architecture, implementation, and how workers process scan, enrichment, and event workloads from queues.

- Repository: [reconurge/flowsint](https://github.com/reconurge/flowsint)
- Tags: architecture
- Published: 2026-06-05

---

**Flowsint leverages Celery as its asynchronous task distribution backbone, allowing FastAPI endpoints to enqueue CPU-intensive scan, enrichment, and event workloads via `send_task()`, while dedicated workers process jobs from configurable queues.**

The reconurge/flowsint repository implements a decoupled microservices architecture where the FastAPI frontend never performs heavy computation directly. Instead, it uses Celery for task distribution, delegating flow orchestration, data enrichment, and event emission to background workers that communicate through a configurable broker and result backend.

## Celery Application Configuration

The Celery instance is centralized in [`flowsint-core/src/flowsint_core/core/celery.py`](https://github.com/reconurge/flowsint/blob/main/flowsint-core/src/flowsint_core/core/celery.py). This file constructs the application object using broker and result backend URLs pulled from the project settings.

```python

# flowsint-core/src/flowsint_core/core/celery.py

from celery import Celery
from flowsint_core.config import settings

celery = Celery(
    "flowsint",
    broker=settings.CELERY_BROKER_URL,
    backend=settings.CELERY_RESULT_BACKEND,
    include=["flowsint_core.tasks.event", 
             "flowsint_core.tasks.enricher", 
             "flowsint_core.tasks.flow"]
)

```

The `include` list registers three task modules—**`event`**, **`enricher`**, and **`flow`**—ensuring workers automatically discover and load these functions at startup. This configuration keeps the broker URL and backend settings externalized, allowing deployment-specific configuration without code changes.

## Task Implementation and Registration

Each unit of work is decorated with `@celery.task`, binding the function to the Celery app for remote execution. The repository organizes tasks into three functional domains.

### Flow Execution Tasks

The `run_flow` task in [`flowsint-core/src/flowsint_core/tasks/flow.py`](https://github.com/reconurge/flowsint/blob/main/flowsint-core/src/flowsint_core/tasks/flow.py) orchestrates complete scan workflows. It creates a database `Scan` entry, executes the orchestration logic, persists results, and explicitly manages task state using `bind=True`.

```python

# flowsint-core/src/flowsint_core/tasks/flow.py

from celery import states

@celery.task(name="run_flow", bind=True)
def run_flow(self, enricher_branches, serialized_objects, sketch_id, owner_id=None):
    try:
        scan = create_scan_entry(sketch_id, owner_id)
        result = orchestrator.run(enricher_branches, serialized_objects)
        scan.store_result(result)
        return {"result": scan.details}
    except Exception as exc:
        self.update_state(state=states.FAILURE, meta={'exc': str(exc)})
        raise

```

Using `bind=True` passes the task instance as `self`, enabling explicit state updates. When failures occur, the task updates its state to `FAILURE`, allowing the API layer to query the result backend and determine the exact failure mode.

### Data Enrichment Tasks

Data processing pipelines reside in [`flowsint-core/src/flowsint_core/tasks/enricher.py`](https://github.com/reconurge/flowsint/blob/main/flowsint-core/src/flowsint_core/tasks/enricher.py). The `run_enricher` and `run_template_enricher` tasks handle generic and template-based enrichment operations, isolating heavy CPU work from the web request lifecycle.

### Event Emission Tasks

Notification workflows are encapsulated in [`flowsint-core/src/flowsint_core/tasks/event.py`](https://github.com/reconurge/flowsint/blob/main/flowsint-core/src/flowsint_core/tasks/event.py). The `emit_event` and `emit_status_event` tasks push notification-style messages to external systems, ensuring non-blocking event propagation even when downstream systems experience latency.

## Dispatching Tasks from the FastAPI Layer

FastAPI route handlers do not import or call task functions directly. Instead, they use the `celery.send_task` method to enqueue work, maintaining a strict separation between the HTTP interface and processing logic.

### Enqueuing Flow Scans

The `/flows` endpoint in [`flowsint-api/app/api/routes/flows.py`](https://github.com/reconurge/flowsint/blob/main/flowsint-api/app/api/routes/flows.py) constructs a payload and dispatches it to the `run_flow` task via `send_task`.

```python

# flowsint-api/app/api/routes/flows.py

@router.post("/flows", response_model=FlowResponse)
def create_flow(request: FlowCreate, db: Session = Depends(get_db)):
    payload = {
        "enricher_branches": request.branches,
        "serialized_objects": request.objects,
        "sketch_id": request.sketch_id,
        "owner_id": request.owner_id,
    }
    task = celery.send_task(
        "run_flow",
        kwargs=payload,
        queue="flows",
        expires=3600
    )
    return {"task_id": task.id, "status": "queued"}

```

By specifying `queue="flows"`, the API routes specific workloads to dedicated queues, enabling targeted scaling of flow-processing workers independent of other task types.

### Triggering Enrichment Pipelines

Similarly, [`flowsint-api/app/api/routes/enrichers.py`](https://github.com/reconurge/flowsint/blob/main/flowsint-api/app/api/routes/enrichers.py) enriches data by dispatching `run_enricher` or `run_template_enricher` tasks. This pattern ensures the API remains stateless and lightweight, as it never waits for long-running database queries or external API calls to complete.

## Worker Architecture and State Management

Workers are launched using the standard Celery CLI, pointing to the application factory in the core module:

```bash
celery -A flowsint_core.core.celery worker \
    --loglevel=INFO \
    -Q flows,enrichers,events \
    --concurrency=4

```

The `-Q` flag assigns workers to specific queues, preventing a backlog in one domain from blocking tasks in another. Workers connect to the broker defined in [`celery.py`](https://github.com/reconurge/flowsint/blob/main/celery.py), consume messages, execute the registered Python functions, and report results back to the result backend.

State tracking relies on Celery’s built-in mechanism. Tasks report `STARTED`, `SUCCESS`, or `FAILURE` statuses automatically. The explicit `self.update_state(state=states.FAILURE)` call in `run_flow` ensures detailed error metadata is available when querying the task status via the.AsyncResult API.

## Summary

- **Centralized configuration**: The Celery app is defined in [`flowsint-core/src/flowsint_core/core/celery.py`](https://github.com/reconurge/flowsint/blob/main/flowsint-core/src/flowsint_core/core/celery.py) with broker and backend URLs sourced from project settings.
- **Modular task discovery**: The `include` list registers `event`, `enricher`, and `flow` modules for automatic worker discovery.
- **Decoupled dispatch**: FastAPI routes use `celery.send_task()` to enqueue work without importing heavy task dependencies.
- **Explicit state control**: Tasks use `bind=True` and `update_state()` to report granular failure states to the result backend.
- **Queue-based scaling**: Dedicated queues (`flows`, `enrichers`, `events`) allow independent worker pools for different workload characteristics.

## Frequently Asked Questions

### What message broker does Flowsint use with Celery?

Flowsint configures the broker via the `CELERY_BROKER_URL` setting in [`flowsint_core/config.py`](https://github.com/reconurge/flowsint/blob/main/flowsint_core/config.py). The [`celery.py`](https://github.com/reconurge/flowsint/blob/main/celery.py) file instantiates the Celery app using this URL, allowing operators to swap between Redis, RabbitMQ, or other supported brokers without modifying task code.

### How does Flowsint track the status of a running scan?

The `run_flow` task uses `bind=True` to access the task instance via `self`, explicitly calling `self.update_state(state=states.FAILURE)` when exceptions occur. This integrates with Celery’s built-in state machine, storing results in the configured `CELERY_RESULT_BACKEND` for retrieval via the task ID returned by `send_task()`.

### Why does the FastAPI layer use `send_task` instead of calling the function directly?

Using `celery.send_task()` enforces a decoupled architecture where the API service does not need to import the heavy task logic or its dependencies. This keeps the API lightweight and stateless while allowing workers to handle CPU-intensive operations like database scans and data enrichment independently.

### Can different task types be routed to specific worker pools?

Yes. The implementation supports queue-based routing by passing the `queue` parameter to `send_task()` (e.g., `queue="flows"`). When starting workers, you specify which queues to consume from using the `-Q` flag (`-Q flows,enrichers,events`), enabling independent scaling and resource allocation for different workload types.