How Flowsint Uses Celery for Task Distribution: Architecture and Implementation
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. This file constructs the application object using broker and result backend URLs pulled from the project settings.
# 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 orchestrates complete scan workflows. It creates a database Scan entry, executes the orchestration logic, persists results, and explicitly manages task state using bind=True.
# 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. 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. 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 constructs a payload and dispatches it to the run_flow task via send_task.
# 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 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:
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, 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.pywith broker and backend URLs sourced from project settings. - Modular task discovery: The
includelist registersevent,enricher, andflowmodules 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=Trueandupdate_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. The 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.
Have a question about this repo?
These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →