Extraction pipeline

Run multiple extraction stages and merge their results, from Extractor classes reference. Import ExtractionPipeline from neo4j_agent_memory.extraction.

ExtractionPipeline

Run stages and merge their results. Import the direct MergeStrategy from neo4j_agent_memory.extraction.pipeline. Its values are distinct from the smaller settings enum.

def ExtractionPipeline(
    stages: list[ExtractionStage | EntityExtractor],
    merge_strategy: MergeStrategy=MergeStrategy.CONFIDENCE,
    stop_on_success: bool=False,
    min_entities_for_success: int=1,
    fallback_on_error: bool=True,
    *,
    ontology: 'OntologyDocument | None'=None,
    min_confidence: float | None=None,
): ...
Parameter Description

stages

List of extraction stages or extractors to run

merge_strategy

Strategy for merging results from stages

stop_on_success

Stop after first stage that returns entities

min_entities_for_success

Minimum entities to consider a stage successful

fallback_on_error

Continue to next stage if current stage errors

ontology

Optional OntologyDocument. When set, the merged result runs through validate_relations(ontology, mode="drop") as the final step.

min_confidence

Optional confidence floor for the merged result. Entities and relations below it are dropped, as are relations left without an endpoint.

Strategy Behavior

UNION

Keep unique entities across results.

INTERSECTION

Keep entities found by multiple results.

CONFIDENCE

Keep the highest-confidence entity for each matching name/type.

CASCADE

Keep earlier-stage entities and fill gaps from later results.

FIRST_SUCCESS

Stop after a stage reaches min_entities_for_success and return the first nonempty successful result.

stop_on_success=True also enables early stopping with the other merge strategies. FIRST_SUCCESS is not FIRST or LAST, and cannot be selected through the config.settings MergeStrategy enum.

from neo4j_agent_memory.extraction import ExtractionPipeline, GLiNER2Extractor, SpacyEntityExtractor
from neo4j_agent_memory.extraction.pipeline import MergeStrategy
from neo4j_agent_memory.ontology import get_template

podcast = get_template("podcast")
pipeline = ExtractionPipeline(
    stages=[SpacyEntityExtractor(), GLiNER2Extractor(ontology=podcast)],
    merge_strategy=MergeStrategy.CONFIDENCE,
    ontology=podcast,
)
result = await pipeline.extract("John discussed Acme in the podcast")

Final stage

After the stages are merged, the pipeline applies two optional steps to the merged result, for every merge strategy:

  1. With min_confidence set, entities and relations scoring below it are dropped, and so are relations whose source or target entity was dropped. This holds stages that take no threshold of their own (spaCy, the LLM extractor) to the same floor.

  2. With ontology set, relations the ontology does not permit are dropped with ExtractionResult.validate_relations(ontology, mode="drop"), and the number dropped is logged at debug level. Merging can only add relations the ontology does not permit: the GLiNER2.5 stage constrains endpoints while decoding, but the LLM stage does not. See Relation extractors.

Give the GLiNER2.5 stage and the pipeline the same ontology. Validating against a different document drops relations the stage decoded legitimately. The configuration factory resolves one ontology (POLE+O when nothing else is configured) and hands it to every stage and to the pipeline. ExtractorBuilder does the same with the ontology you name and validates nothing when you name none. A pipeline you construct yourself does no validation unless you pass ontology.

extract

Run extraction through all stages.

async def extract(
    text: str,
    *,
    entity_types: list[str] | None=None,
    extract_relations: bool=True,
    extract_preferences: bool=True,
) -> ExtractionResult: ...
Parameter Description

text

Text to extract from

entity_types

Optional list of entity types to extract

extract_relations

Whether to extract relations

extract_preferences

Whether to extract preferences

extract_with_details

Run extraction with detailed stage information.

async def extract_with_details(
    text: str,
    *,
    entity_types: list[str] | None=None,
    extract_relations: bool=True,
    extract_preferences: bool=True,
) -> PipelineResult: ...

extract_batch

Extract entities from multiple texts in parallel.

async def extract_batch(
    texts: list[str],
    *,
    batch_size: int=10,
    max_concurrency: int=5,
    entity_types: list[str] | None=None,
    extract_relations: bool=True,
    extract_preferences: bool=True,
    on_progress: ProgressCallback | None=None,
    fail_fast: bool=False,
) -> BatchExtractionResult: ...
Parameter Description

texts

List of texts to extract from

batch_size

Number of texts to process in each batch (for memory management)

max_concurrency

Maximum number of concurrent extractions

entity_types

Optional list of entity types to extract

extract_relations

Whether to extract relations

extract_preferences

Whether to extract preferences

on_progress

Optional callback called after each text is processed. Receives (completed_count, total_count).

fail_fast

If True, stop on first error. If False, continue and collect errors.

add_stage

Add a stage to the pipeline.

def add_stage(stage: ExtractionStage | EntityExtractor) -> None: ...

remove_stage

Remove a stage by name. Returns True if found and removed.

def remove_stage(name: str) -> bool: ...

stage_names lists configured stage names. extract_with_details returns PipelineResult with merged results and per-stage timing/error details; extract returns just ExtractionResult. BatchExtractionResult retains per-item results/errors.

StageResult

Fields and defaults:

Field Type Default Description

stage_name

str

required

result

ExtractionResult

required

success

bool

True

error

str | None

None

duration_ms

float

0.0

PipelineResult

Fields and defaults:

Field Type Default Description

final_result

ExtractionResult

required

stage_results

list[StageResult]

field(default_factory=list)

merge_strategy

MergeStrategy

MergeStrategy.CONFIDENCE

total_duration_ms

float

0.0

BatchItemResult

Fields and defaults:

Field Type Default Description

index

int

required

result

ExtractionResult

required

success

bool

True

error

str | None

None

duration_ms

float

0.0

BatchExtractionResult

Fields and defaults:

Field Type Default Description

results

list[BatchItemResult]

field(default_factory=list)

total_duration_ms

float

0.0

BatchExtractionResult exposes total_items, successful_items, failed_items, success_rate, total_entities, and total_relations properties plus get_extraction_results(), get_all_entities(), and get_errors() helpers.