How to Run Custom Pipelines in Cognee: A Complete Guide
Run custom pipelines in Cognee by wrapping functions in Task objects and passing them to cognee.run_custom_pipeline(), which orchestrates execution through the core run_pipeline engine with support for blocking or background modes.
Cognee is an open-source knowledge graph engine that enables developers to build complex data processing workflows. While the framework provides built-in CLI commands for common operations, the run_custom_pipeline interface in topoteretes/cognee allows you to compose arbitrary sequences of tasks for specialized ingestion, enrichment, or graph construction needs. This guide explains how to construct and execute custom pipelines using the underlying Task API and pipeline execution layers.
Understanding the Custom Pipeline Architecture
Custom pipelines in Cognee are built on three foundational components: the Task wrapper, the orchestration layer, and the core execution engine. Understanding how these pieces fit together enables you to construct everything from simple text ingestion to complex LLM-driven knowledge graph creation.
Task Definition with the Task Class
Every step in a Cognee pipeline is encapsulated by a Task object defined in cognee/modules/pipelines/tasks/task.py. This wrapper records the callable function, default arguments, batch size, and determines execution strategy—whether the underlying code is a standard function, coroutine, generator, or async generator.
The Task class abstracts away execution complexity, allowing you to focus on business logic. When you instantiate a Task, you provide the function reference and any default parameters, which the pipeline engine later invokes with runtime data.
Pipeline Orchestration Layer
The run_custom_pipeline function in cognee/modules/run_custom_pipeline/run_custom_pipeline.py serves as the primary entry point for custom workflows. This function performs three critical operations:
- Task normalization – Validates and prepares the list of
Taskobjects for execution - Executor selection – Retrieves a pipeline executor via
get_pipeline_executorthat determines whether the pipeline runs blocking or in the background (defined incognee/modules/pipelines/layers/pipeline_execution_mode.py) - Core delegation – Forwards all parameters (tasks, data, dataset, user, configuration flags) to the central
run_pipelinefunction
Core Runner and Execution Flow
The run_pipeline function in cognee/modules/pipelines/operations/pipeline.py implements the actual execution logic. It validates the task list through validate_pipeline_tasks, pre-checks the environment via setup_and_check_environment, and resolves authorized datasets for the supplied user through resolve_authorized_user_datasets.
For each authorized dataset, the runner streams PipelineRunInfo objects while delegating work to run_tasks. This architecture yields progress updates, start/completion events, and error states, enabling real-time monitoring of long-running operations.
Execution Modes and Configuration Options
Cognee supports two primary execution strategies for custom pipelines, controlled through the orchestration layer in pipeline_execution_mode.py.
Blocking vs. Background Execution
Blocking execution (run_pipeline_blocking) runs the coroutine directly, yielding results to the caller until completion. This mode is suitable for scripts and notebooks where you need immediate confirmation of success or failure.
Background execution (run_pipeline_as_background_process) launches the pipeline in a separate async task, immediately returning a pipeline_run_id that you can use to monitor progress later. Enable this mode by setting run_in_background=True when calling run_custom_pipeline.
Performance and Caching Controls
The pipeline engine exposes several optimization flags:
use_pipeline_cache– Skips re-processing if a pipeline with the same ID has already run or is currently runningincremental_loading– Processes only newly added or changed data, requiring Cognee’sDatamodel to track statedata_per_batch– Controls parallel batch size for task execution, tuning throughput for I/O-bound operations
Code Examples
Minimal Custom Pipeline for Text Ingestion
The following example demonstrates adding plain text to Cognee using the Task API:
import asyncio
import cognee
from cognee.modules.pipelines import Task
from cognee.modules.users.methods import get_default_user
async def main() -> None:
user = await get_default_user()
# Define the tasks that compose the pipeline
add_tasks = [
Task(cognee.modules.tasks.ingestion.resolve_data_directories, include_subdirectories=True),
Task(cognee.modules.tasks.ingestion.ingest_data, "main_dataset", user),
]
# Run the pipeline – blocking mode (default)
await cognee.run_custom_pipeline(
tasks=add_tasks,
data="Natural language processing (NLP) is a subfield of AI.",
user=user,
dataset="main_dataset",
pipeline_name="add_text",
)
print("✅ Text added to Cognee")
if __name__ == "__main__":
asyncio.run(main())
Key implementation details:
Taskobjects wrap low-level ingestion helpers fromcognee.modules.tasks.ingestionrun_custom_pipelineforwards thedataparameter to the first task (resolve_data_directories)- The call blocks until completion because
run_in_backgrounddefaults toFalse
Background Execution with Caching and Incremental Loading
For production workflows processing large datasets, use background execution with caching:
import asyncio, cognee
from cognee.modules.pipelines import Task
from cognee.modules.users.methods import get_default_user
async def main():
user = await get_default_user()
# Re‑use the default Cognify tasks (LLM‑driven KG creation)
cognify_tasks = await cognee.api.v1.cognify.cognify.get_default_tasks(user=user)
# Launch the pipeline in the background, enable cache & incremental mode
pipeline_info = await cognee.run_custom_pipeline(
tasks=cognify_tasks,
user=user,
dataset="main_dataset",
pipeline_name="cognify_pipeline",
use_pipeline_cache=True,
incremental_loading=True,
run_in_background=True, # ← fire‑and‑forget
)
print(f"🚀 Background pipeline started – run_id={pipeline_info.pipeline_run_id}")
# You can later inspect status with `cognee.get_pipeline_status(pipeline_info.pipeline_run_id)`
if __name__ == "__main__":
asyncio.run(main())
This configuration:
- Returns immediately with a
PipelineRunStartedobject containing thepipeline_run_id - Avoids duplicate work through
use_pipeline_cache - Processes only new documents via
incremental_loading=True
Full-Featured Example
For a complete end-to-end demonstration, reference examples/custom_pipelines/custom_cognify_pipeline_example.py in the Cognee repository. This script illustrates the full workflow:
- Reset Cognee state using
cognee.prune.prune_data() - Setup a relational database via
cognee.modules.engine.operations.setup.setup - Create and run an "add" pipeline using
Taskobjects - Execute built-in Cognify tasks inside a custom pipeline
- Query the resulting graph with
cognee.search
Copy this example, adjust the LLM_API_KEY in your .env file, and run it directly to observe custom pipeline behavior in action.
Summary
- Custom pipelines in Cognee are constructed by chaining
Taskobjects and executing them viacognee.run_custom_pipeline() - The Task class in
cognee/modules/pipelines/tasks/task.pywraps functions, coroutines, and generators with execution metadata - Execution mode is determined by
get_pipeline_executorinpipeline_execution_mode.py, supporting both blocking and background processing - Core orchestration happens in
cognee/modules/pipelines/operations/pipeline.py, which streamsPipelineRunInfoobjects for progress monitoring - Performance features include
use_pipeline_cachefor deduplication,incremental_loadingfor delta processing, anddata_per_batchfor parallelism tuning
Frequently Asked Questions
What is the difference between run_custom_pipeline and the built-in CLI commands?
run_custom_pipeline in cognee/modules/run_custom_pipeline/run_custom_pipeline.py provides programmatic access to Cognee's pipeline engine, allowing you to compose arbitrary sequences of Task objects for specific data processing needs. CLI commands wrap predefined task sequences for standard workflows, while the custom pipeline API enables ad-hoc ingestion, custom LLM-based enrichment, or partial graph reconstruction without modifying the framework's core code.
How do I monitor a pipeline running in background mode?
When you set run_in_background=True, run_custom_pipeline returns a PipelineRunStarted object containing a pipeline_run_id. You can later query the status of this execution using cognee.get_pipeline_status(pipeline_run_id) to retrieve progress updates, completion status, or error states. The underlying run_pipeline function yields PipelineRunInfo objects that track the pipeline lifecycle through the streaming interface implemented in cognee/modules/pipelines/operations/pipeline.py.
Can I mix built-in Cognify tasks with custom functions in the same pipeline?
Yes. Retrieve the default Cognify task list via cognee.api.v1.cognify.cognify.get_default_tasks(user=user), then append or prepend your own Task objects before passing the combined list to run_custom_pipeline. This approach allows you to insert custom preprocessing steps before LLM-driven knowledge graph construction or add post-processing logic after the standard Cognify workflow completes.
What happens when use_pipeline_cache is enabled?
When use_pipeline_cache=True, the pipeline engine checks whether a pipeline with the same ID has already run or is currently running. If a cached result exists, the engine skips re-processing and returns the existing output. This optimization is particularly valuable for expensive operations like LLM-based entity extraction or large-scale document ingestion where duplicate work would waste computational resources and API costs.
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 →