How Celery Tasks Handle Async Enricher Processing in Flowsint
Flowsint uses Celery's distributed task queue to execute enrichment jobs asynchronously, bridging synchronous Celery workers with async enricher implementations via asyncio.run() while persisting results to a Scan table indexed by the task UUID.
Flowsint leverages Celery to offload compute-intensive enrichment operations from the web request cycle. This architecture, implemented in the reconurge/flowsint repository, allows the platform to execute both Python-based and template-driven enrichers asynchronously while maintaining full traceability through a dedicated database model.
The Celery Task Architecture
Flowsint defines two primary Celery tasks in flowsint-core/src/flowsint_core/tasks/enricher.py to handle different enrichment paradigms. Both tasks follow a similar execution pattern but differ in how they instantiate the enricher logic.
Python-Based Enrichers
The run_enricher task executes enrichers that implement the Enricher interface defined in flowsint_core.core.enricher_base. When invoked, the task performs the following steps:
- Creates a new
Scandatabase row using a UUID that matches the Celery task ID (self.request.id) - Optionally instantiates a Vault service when an
owner_idis provided - Retrieves the enricher class from the global
ENRICHER_REGISTRY - Passes JSON-serialized input objects to the enricher via
asyncio.run(enricher.execute(...)) - Updates the
Scanstatus toCOMPLETEDand stores JSON-serializable results inscan.detailson success - Rolls back the transaction, marks the
ScanasFAILED, and updates the Celery state toFAILUREon error
Reference: Task definition and execution flow—lines 26-73 of enricher.py.
Template-Driven Enrichers
The run_template_enricher task handles enrichers defined as YAML templates stored in the database. This task mirrors the Python-based flow but instantiates a TemplateEnricher class instead:
- Loads the template configuration from the database
- Creates a
Scanrecord with the Celery task UUID - Resolves vault credentials if
owner_idis present - Executes the template logic via
asyncio.run(enricher.execute(...)) - Persists results to the same
Scantable structure
Reference: Template task implementation—lines 95-152 of enricher.py.
Bridging Synchronous Workers with Async Enrichers
Celery workers operate on synchronous Python code, while Flowsint enrichers implement asynchronous execution patterns for I/O-bound operations like HTTP requests.
The bridge occurs through the asyncio.run() method:
# From flowsint_core/tasks/enricher.py (lines 69-70, 45-46)
results = asyncio.run(enricher.execute(serialized_objects))
This pattern starts a temporary event loop within the synchronous Celery worker, executes the async execute method to completion, and returns the results to the worker process. The Enricher base class requires all implementations to define async def execute(self, objects) according to the interface in enricher_base.py.
Results are converted to JSON-serializable format using to_json_serializable() before storage, ensuring compatibility with Celery's JSON serializer and PostgreSQL JSON columns.
Celery Configuration and Time Limits
The Celery application is configured in flowsint-core/src/flowsint_core/core/celery.py:
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",
],
)
celery.conf.update(
task_serializer="json",
accept_content=["json"],
result_serializer="json",
timezone="UTC",
enable_utc=True,
task_track_started=True,
task_time_limit=3600, # 1 hour per task
worker_max_tasks_per_child=1000,
worker_prefetch_multiplier=4,
)
Key configuration parameters:
- Broker and Backend: URLs are sourced from
settings.CELERY_BROKER_URLandsettings.CELERY_RESULT_BACKENDin the config module - Task Discovery: The
includelist ensures both enrichment tasks are auto-discovered by workers - Time Limits:
task_time_limit=3600prevents runaway enrichers from consuming resources indefinitely - Prefetch Control:
worker_prefetch_multiplier=4limits task pre-fetching, improving fairness when many long-running enrichers are queued
Error Handling and State Management
Both tasks implement comprehensive error handling through try/except/finally blocks:
- Database Transactions: On any exception,
session.rollback()is invoked and theScanrecord updates withstatus = EventLevel.FAILEDand the error message - Celery State:
self.update_state(state=states.FAILURE)notifies the Celery monitor of task failure - Audit Logging: Errors are printed and logged via the internal
Loggerfor debugging and compliance purposes
This dual-layer state tracking ensures the frontend can poll either Celery's result backend or the Scan table to determine task status.
End-to-End Implementation Examples
Enqueuing Python Enrichers
To trigger a Python-based enrichment job:
from flowsint_core.tasks.enricher import run_enricher
# objects is a list of dicts from the API payload
task = run_enricher.delay(
enricher_name="github_user_lookup",
serialized_objects=objects,
sketch_id="a1b2c3d4-e5f6-7g8h-9i0j-k1l2m3n4o5p6",
owner_id="d4e5f6a7-b8c9-0d1e-f2g3-h4i5j6k7l8m9",
)
scan_id = task.id # UUID shared between Celery and Scan table
Enqueuing Template Enrichers
For YAML template-based enrichments:
from flowsint_core.tasks.enricher import run_template_enricher
task = run_template_enricher.delay(
template_name="whois_lookup",
serialized_objects=objects,
sketch_id=None, # optional
owner_id="123e4567-e89b-12d3-a456-426614174000",
)
scan_id = task.id
Retrieving Results
Since the Celery task ID matches the Scan table UUID, clients can poll for results:
from flowsint_core.core.models import Scan
from flowsint_core.core.postgre_db import SessionLocal
def get_scan_result(scan_id: str):
with SessionLocal() as db:
scan = db.query(Scan).filter(Scan.id == scan_id).first()
if not scan:
raise ValueError("Scan not found")
return {
"status": scan.status,
"details": scan.details,
"error": scan.error,
}
The status field reflects Celery task states: PENDING, COMPLETED, or FAILED.
Summary
- Flowsint uses two Celery tasks—
run_enricherandrun_template_enricher—to process enrichment jobs asynchronously outside the request-response cycle - The
asyncio.run()bridge inflowsint_core/tasks/enricher.pyallows synchronous Celery workers to execute async enricher implementations - Task status and results are persisted to a
Scantable using the Celery task UUID as the primary key, enabling reliable polling mechanisms - Configuration in
flowsint_core/core/celery.pyenforces 1-hour time limits and JSON serialization for task safety - Comprehensive error handling updates both the database transaction state and Celery's internal failure state when enrichers encounter exceptions
Frequently Asked Questions
Why does Flowsint use asyncio.run() instead of native async Celery tasks?
Celery workers historically operate in synchronous contexts. Rather than implementing complex async worker pools, Flowsint uses asyncio.run() as a bridge to execute async enricher code within standard synchronous Celery tasks. This pattern (found in lines 69-70 of run_enricher) creates a temporary event loop, runs the async I/O operations to completion, and returns control to the Celery worker, simplifying deployment and compatibility.
How does Flowsint track the status of long-running enrichment tasks?
Flowsint implements a dual-tracking system using the Scan model in flowsint_core/core/models.py. When a task starts, it creates a Scan row with an ID matching self.request.id (the Celery task UUID). The task updates this record with COMPLETED or FAILED status upon completion. Clients can query the PostgreSQL Scan table directly using the task ID, avoiding the need to maintain Celery result backend connections for status checks.
What happens when an enricher fails during execution?
Both tasks wrap execution in try/except/finally blocks. On exception, the code calls session.rollback() to revert database changes, updates the Scan record with status=EventLevel.FAILED and the error message, and invokes self.update_state(state=states.FAILURE) to notify Celery's monitoring system. This ensures atomicity between the database state and Celery's internal task state while preserving error details for debugging.
Where is the Celery broker configuration defined in Flowsint?
The broker URL and result backend are defined in flowsint_core/core/config.py as CELERY_BROKER_URL and CELERY_RESULT_BACKEND. These settings are consumed by the Celery app instantiation in flowsint_core/core/celery.py. The configuration supports standard Celery broker protocols (Redis, RabbitMQ, etc.) via environment variables, allowing flexible deployment across different infrastructure environments.
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 →