How to Configure Distributed Execution in Cognee: A Complete Guide with Modal
Enable distributed execution in Cognee by setting the environment variable COGNEE_DISTRIBUTED=True and launching the Modal application, which automatically routes storage adapter operations through scalable remote workers instead of executing them locally.
Cognee supports distributed execution to offload heavy graph and data-point operations to Modal workers. When you configure distributed execution in Cognee, core storage adapters for Neo4j, PGVector, and other backends automatically route write operations through Modal queues rather than processing them in the main process.
How Distributed Execution Works in Cognee
The distributed architecture in topoteretes/cognee uses Modal for containerized worker orchestration. The system intercepts storage adapter method calls using a decorator-based approach and redirects them to persistent queues when distributed mode is active.
The COGNEE_DISTRIBUTED environment variable serves as the global toggle. When set to "True", the override_distributed decorator in distributed/utils.py (lines 5-16) forwards calls to asynchronous Modal tasks instead of executing local synchronous methods. These tasks push payloads to dedicated queues consumed by specialized workers that handle the actual database writes with distributed=False to prevent recursive distribution.
Prerequisites for Distributed Mode
Before configuring distributed execution, ensure you have:
- Modal credentials configured locally or in your deployment environment
- Access to Modal's infrastructure for running containers and queues
- Python 3.9+ with Cognee installed
- Existing storage adapters (Neo4j, PostgreSQL with PGVector, etc.) properly configured
Step-by-Step Configuration
Set the COGNEE_DISTRIBUTED Environment Variable
Activate distributed mode by exporting the environment variable before running any Cognee operations. In distributed/entrypoint.py at line 18, the system explicitly sets this variable to ensure workers run in the correct mode:
export COGNEE_DISTRIBUTED=True
You can also set this programmatically in Python before importing Cognee modules:
import os
os.environ["COGNEE_DISTRIBUTED"] = "True"
Launch the Modal Application
Deploy the Modal application defined in distributed/app.py (lines 1-4) to instantiate the queues and worker infrastructure:
modal run cognee/distributed/entrypoint.py
This command launches the Modal App instance and spins up the graph_saving_worker and data_point_saving_worker processes that listen on the queues defined in distributed/queues.py (lines 1-5).
Run Your Cognee Pipeline
Once the environment variable is set and workers are running, execute your pipelines normally. The override_distributed decorator automatically intercepts storage operations:
from cognee.modules.pipelines.operations.run_tasks import run_tasks
# Pipeline calls storage adapters; @override_distributed redirects
# writes to Modal queues because COGNEE_DISTRIBUTED=True
await run_tasks(tasks=[your_task])
When COGNEE_DISTRIBUTED is false (the default), all calls execute directly against local adapters, preserving single-process behavior.
Core Components and Architecture
The override_distributed Decorator
Located in distributed/utils.py at lines 5-16, this decorator wraps storage adapter methods like add_nodes, add_edges, and add_data_points. It checks the COGNEE_DISTRIBUTED flag and conditionally routes execution:
- Local mode: Executes the original synchronous method directly
- Distributed mode: Invokes the corresponding Modal task (
queued_add_nodes,queued_add_edges, orqueued_add_data_points)
Storage adapters such as the Neo4j driver and PGVector adapter apply this decorator to their write methods, making distribution transparent to the rest of the codebase.
Modal Queues and Workers
The distributed/queues.py file instantiates persistent Modal queues that hold batched payloads:
add_nodes_and_edges_queue: Handles graph structure operationsadd_data_points_queue: Manages vector store embeddings
Workers defined in distributed/workers/graph_saving_worker.py (line 51) and distributed/workers/data_point_saving_worker.py poll these queues continuously. They deserialize payloads and call the regular storage adapters with distributed=False to ensure writes execute locally within the worker container.
The queued tasks in distributed/tasks/queued_add_nodes.py (lines 1-13) implement retry logic with batch splitting to handle GRPC errors and prevent queue overflow.
Adapter Integration
Storage adapters integrate distribution through method decoration. For example, the Neo4j adapter in cognee/infrastructure/databases/graph/neo4j_driver/adapter.py annotates its add_nodes method with @override_distributed(queued_add_nodes). Similarly, cognee/infrastructure/databases/vector/pgvector/PGVectorAdapter.py wraps its data-point methods. This is the only modification required in core storage code to support distribution.
Advanced Configuration and Debugging
Manual Queue Operations
For debugging or custom batching, you can manually push data to the queues:
from distributed.queues import add_nodes_and_edges_queue
import asyncio
batch = [
{"id": "node1", "type": "Person"},
{"id": "node2", "type": "Company"}
]
# Directly push a batch onto the queue
await add_nodes_and_edges_queue.put.aio((batch, []))
Customizing Queue Parameters
Adjust queue behavior in your Modal configuration by modifying the queue instantiation in distributed/queues.py:
from modal import Queue
# Create a queue with increased capacity for high-throughput scenarios
large_queue = Queue.from_name(
"large_queue",
create_if_missing=True,
max_size=10_000
)
Summary
- Set
COGNEE_DISTRIBUTED=Trueto enable distributed execution across all storage operations - Launch the Modal application using
modal run cognee/distributed/entrypoint.pyto start worker processes - Use
@override_distributeddecorator pattern to make any storage adapter distributed-aware - Queues handle batching of nodes, edges, and data-points with automatic retry and splitting logic
- Workers consume queues and execute writes locally within containerized Modal environments
- Default behavior remains local when the environment variable is unset or false
Frequently Asked Questions
What is the default execution mode in Cognee?
By default, Cognee operates in local synchronous mode. When the COGNEE_DISTRIBUTED environment variable is unset or set to "False", the override_distributed decorator in distributed/utils.py passes through to the original method implementation without invoking Modal, ensuring no external dependencies are required for standard operation.
How does Cognee handle failures in distributed mode?
The queued tasks implemented in files like distributed/tasks/queued_add_nodes.py include recursive batch-splitting logic that activates on GRPC errors. If a batch fails to enqueue due to size limits or network issues, the system automatically splits the payload into smaller chunks and retries, ensuring eventual delivery without blocking the main pipeline.
Can I use distributed execution with any storage adapter?
Yes, any storage adapter can participate in distributed execution by applying the @override_distributed decorator to its write methods. The repository already implements this for Neo4j in cognee/infrastructure/databases/graph/neo4j_driver/adapter.py and PGVector in cognee/infrastructure/databases/vector/pgvector/PGVectorAdapter.py, but the pattern works for custom adapters as well.
Do I need to modify my existing pipeline code to use distributed mode?
No, existing pipelines require zero code changes to benefit from distributed execution. Because the intercept logic resides in the @override_distributed decorator applied to storage adapters, calling await run_tasks() from cognee/modules/pipelines/operations/run_tasks.py automatically respects the COGNEE_DISTRIBUTED environment variable and routes operations accordingly.
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 →