3.2w topoteretes

cognee-custom-pipelines Skill

在构建自己的 cognee 处理流程时使用——编写自定义任务,用 run_custom_pipeline 或轻量的 run_pipeline(from cognee.pipelines import run_pipeline)将它们串成 pipeline,用 add_data_points 存储自定义 DataPoint,用 memify 对现有图谱运行自定义提取/增强,检查 pipeline 运行状态,或调试任务间的数据流(batch_size、data_per_batch、ctx、Drop、enriches)。

安装方式:把技能目录放入 ~/.claude/skills/(Claude Code)或在 claude.ai 设置中启用;也可复制右侧安装命令一键添加。

查看源码

技能指令原文(SKILL.md)

Custom tasks and pipelines

Everything cognee does runs as a pipeline: an ordered list of tasks,
each a plain Python function whose output feeds the next one. remember()
is the right tool for ordinary ingestion. Build a pipeline when you need
processing cognee does not ship: your own extraction, your own node types,
or a post-processing step over the graph.

import cognee
from cognee.modules.pipelines import Task
from cognee.tasks.storage import add_data_points
from cognee.low_level import DataPoint


class Person(DataPoint):
    name: str
    metadata: dict = {"index_fields": ["name"], "identity_fields": ["name"]}


async def extract_people(data_items: list) -> list[Person]:
    people = []
    for item in data_items:  # always a list, see below
        text = item if isinstance(item, str) else ""
        people += [Person(name=n.strip()) for n in text.split(",") if n.strip()]
    return people


result = await cognee.run_custom_pipeline(
    tasks=[
        Task(extract_people, needs_llm=False),
        Task(add_data_points, needs_llm=False),  # store in graph + vector DBs
    ],
    data=["Ada Lovelace, Alan Turing"],
    dataset="people",
)

Use it

Pick the runner

There are three, and two share the name run_pipeline:

| Runner | Import | Use it for |
|---|---|---|
| cognee.run_custom_pipeline(...) | cognee | The normal choice: runs your tasks against a dataset with permissions, a per-dataset lock, run records, and status |
| Full orchestrator run_pipeline(tasks=..., data=..., datasets=...) | cognee.modules.pipelines | What run_custom_pipeline and cognify call; yields PipelineRunInfo |
| Lightweight run_pipeline([...], data=...) | from cognee.pipelines import run_pipeline (after import cognee, the attribute cognee.pipelines.run_pipeline is the orchestrator) | Quick chains of task() specs with no permissions, locks, run rows, or migrations; returns the last step's outputs |

`cognee.run_custom_pipeline(tasks, data=None, dataset="main_dataset",
user=None, incremental_loading=False, data_per_batch=20,
run_in_background=False, pipeline_name="custom_pipeline", data_cache=False,
...) returns {dataset_id: PipelineRunInfo}` (the started run when
run_in_background=True). With data=None it runs over the dataset's
existing documents (Data rows).

Write a task

from cognee.modules.pipelines import Task
from cognee.modules.pipelines.models import PipelineContext
from cognee.modules.pipelines.tasks.task import task_summary
from cognee.pipelines import Drop


@task_summary("Tagged {n} chunk(s)")
async def tag_chunks(chunks: list, ctx: PipelineContext = None, label: str = "x"):
    for chunk in chunks:
        chunk.metadata["label"] = label
    return chunks  # or yield per item; return/yield Drop to discard


tag = Task(tag_chunks, label="reviewed", batch_size=10, needs_llm=False)
  • A task is an async def, a generator, an async generator, or a plain

def. Extra Task(fn, args, *kwargs) arguments are passed after the
pipeline data.

  • needs_llm=False on tasks that never call an LLM lets an LLM-free

pipeline skip the LLM connection check.

  • ctx (injected by the parameter name ctx) carries user,

data_item, dataset, pipeline_run_id, pipeline_name, and extras.

  • task.with_config(batch_size=..., **kwargs) returns a modified copy.

How data flows

  • Each document runs the whole chain on its own with run_custom_pipeline

or the orchestrator, and the first task receives it as a **one-element
list** ([data_item]), not the bare item. The lightweight run_pipeline
passes data to the first task unchanged.

  • data_per_batch (default 20) is how many documents run at the same

time. It is a concurrency limit, not a batch size.

  • batch_size belongs to the consumer. A task's batch_size decides how

the previous task's generator output is grouped before it is passed in.
Generator tasks always hand over lists; a coroutine or function hands over
its single return value.

  • Streaming: each upstream result goes down the chain immediately, so a

downstream task can run many times per document.

  • enriches=True: if the task returns None, its input is passed on

unchanged (coroutines and functions only, not generators).

  • Drop: returning or yielding it removes that item from the stream.
  • Every DataPoint passing through is stamped automatically with where it

came from (source_pipeline, source_task, source_user, …).

Store results

`add_data_points(data_points, custom_edges=None, embed_triplets=False,
graph_only=False) writes a list of DataPoint`s to the graph and indexes
their index_fields in the vector DB. It returns the same list, so it can
sit mid-chain. Give every node type identity_fields so repeated runs merge
instead of duplicating (see the cognee-custom-graph-models skill).

Work on the existing graph: memify

await cognee.memify(
    extraction_tasks=["extract_subgraph_chunks"],  # names or Task objects
    enrichment_tasks=[Task(my_enrichment, needs_llm=False)],
    dataset="people",
    node_name=["AI"],  # optional subgraph filter
)

With no data, memify passes the graph (or the node_type / node_name
subgraph) to the first task. Registered task names:
extract_subgraph, extract_subgraph_chunks, get_triplet_datapoints,
extract_user_sessions, cognify_session, extract_agent_trace_feedbacks,
cognify_agent_trace_feedback, apply_feedback_weights,
detect_entity_duplicates, merge_entity_duplicates, index_data_points.
improve() also forwards extraction_tasks / enrichment_tasks to memify,
but only inside its enrichment stage. With custom tasks that stage skips
the TRIPLET_EMBEDDING gate and the has-the-graph-changed check, so they
run on every improve (unless the stage is disabled, the lock is held, or
the fatal persist_session_qa stage errors and stops the run first).

Check status

status = await cognee.datasets.get_status([dataset_id], pipeline_names=["custom_pipeline"])

Without pipeline_names it reports only cognify_pipeline. It returns
{str(dataset_id): PipelineRunStatus} (`{str(dataset_id): {pipeline_name:
status}} for several pipeline_names): DATASET_PROCESSING_STARTED`,
_COMPLETED, or _ERRORED (_INITIATED exists only on legacy rows). The value run_custom_pipeline
returns per dataset is a PipelineRunInfo instead, whose class names the
outcome: PipelineRunCompleted, PipelineRunAlreadyCompleted,
PipelineRunErrored, and so on.

Pitfalls

  • Wrong run_pipeline. The one imported via

from cognee.pipelines import run_pipeline wants task()
specs called (extract(), not extract) and raises TypeError
otherwise; the orchestrator in cognee.modules.pipelines wants Task
objects and raises WrongTaskTypeError otherwise.

  • String task names only work in memify. run_custom_pipeline accepts

only Task objects despite its type hint.

  • Some callables are rejected by Task (`ValueError: Unsupported task

type): bound methods and functools.partial`s of a plain (non-generator)
sync function, and callable objects (instances with __call__). Generator
and async variants, plain functions, and lambdas work. When in doubt, wrap it
in a plain def / async def.

  • run_custom_pipeline does not run database migrations. On an existing

database, run await cognee.run_migrations() (or any remember() first).

  • Keep pipeline_name="custom_pipeline" unless you add your name to

WRITE_PIPELINE_NAMES in cognee/modules/improve/graph_changes.py.
Otherwise improve() does not notice your graph writes and may skip
enrichment as "already completed".

  • memify defaults. An omitted or empty task list is replaced by the

defaults: index_data_points enrichment, plus get_triplet_datapoints
extraction only when TRIPLET_EMBEDDING=true (off by default). memify
uses only the first dataset it resolves.

  • Nodes duplicate on every run when a DataPoint has no

identity_fields (or Dedup() fields). examples/guides/custom_data_models.py
and examples/guides/custom_tasks_and_pipelines.py have this bug; don't
copy it.

How it works

run_custom_pipeline → orchestrator run_pipeline (checks write
permission, takes the per-dataset lock, records a PipelineRun) →
run_tasks (a semaphore of data_per_batch, one chain per document) →
run_tasks_base (streams each task's output into the next, batching by the
consumer's batch_size, injecting ctx, stamping provenance).

  • Package overview and the runner semantics: cognee/modules/pipelines/__init__.py
  • Task, task(), TaskSpec, BoundTask, @task_summary: cognee/modules/pipelines/tasks/task.py
  • Orchestrator: cognee/modules/pipelines/operations/pipeline.py; execution:

run_tasks.py, run_tasks_base.py, run_tasks_data_item.py

  • Lightweight runner: cognee/modules/pipelines/operations/run_pipeline.py,

exported from cognee/pipelines/

  • Context: cognee/modules/pipelines/models/PipelineContext.py
  • run_custom_pipeline: cognee/modules/run_custom_pipeline/run_custom_pipeline.py
  • memify: cognee/modules/memify/memify.py,

cognee/memify_pipelines/memify_task_registry.py, memify_default_tasks.py

  • Storage: cognee/tasks/storage/add_data_points.py
  • Index of all shipped tasks: cognee/tasks/README.md

Examples:

  • examples/demos/custom_pipelines/custom_pipeline_single_object_example.py:

the best reference. It runs over added documents, does LLM extraction into
typed DataPoints, then recalls. Its models declare identity with Dedup()
(the Annotated alternative to metadata["identity_fields"]).

  • examples/demos/custom_pipelines/organizational_hierarchy/: low-level

run_tasks, no LLM, dedup via identity_fields, status polling.

  • examples/demos/custom_pipelines/custom_cognify_pipeline_example.py:

rebuilds add + cognify from the default task list.

  • examples/demos/custom_pipelines/memify_coding_agent_rule_extraction_example.py:

memify with a custom enrichment task.

Extending it

  • A new shipped task: put it in the cognee/tasks/ subpackage for its

stage, export it from that package's __init__.py, follow the template in
cognee/tasks/README.md, and add a unit test under cognee/tests/unit/tasks/.

  • A new memify task name: register it in

cognee/memify_pipelines/memify_task_registry.py.

  • A new write pipeline name: add it to WRITE_PIPELINE_NAMES.
  • Pipeline tests: cognee/tests/unit/modules/pipelines/ (runner semantics,

context, provenance, rollback) and cognee/tests/unit/pipelines/ (the
lightweight API).