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 Task objects for execution
  • Executor selection – Retrieves a pipeline executor via get_pipeline_executor that determines whether the pipeline runs blocking or in the background (defined in cognee/modules/pipelines/layers/pipeline_execution_mode.py)
  • Core delegation – Forwards all parameters (tasks, data, dataset, user, configuration flags) to the central run_pipeline function

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 running
  • incremental_loading – Processes only newly added or changed data, requiring Cognee’s Data model to track state
  • data_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:

  • Task objects wrap low-level ingestion helpers from cognee.modules.tasks.ingestion
  • run_custom_pipeline forwards the data parameter to the first task (resolve_data_directories)
  • The call blocks until completion because run_in_background defaults to False

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 PipelineRunStarted object containing the pipeline_run_id
  • Avoids duplicate work through use_pipeline_cache
  • Processes only new documents via incremental_loading=True

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:

  1. Reset Cognee state using cognee.prune.prune_data()
  2. Setup a relational database via cognee.modules.engine.operations.setup.setup
  3. Create and run an "add" pipeline using Task objects
  4. Execute built-in Cognify tasks inside a custom pipeline
  5. 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 Task objects and executing them via cognee.run_custom_pipeline()
  • The Task class in cognee/modules/pipelines/tasks/task.py wraps functions, coroutines, and generators with execution metadata
  • Execution mode is determined by get_pipeline_executor in pipeline_execution_mode.py, supporting both blocking and background processing
  • Core orchestration happens in cognee/modules/pipelines/operations/pipeline.py, which streams PipelineRunInfo objects for progress monitoring
  • Performance features include use_pipeline_cache for deduplication, incremental_loading for delta processing, and data_per_batch for 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:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →