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 |
|---|---|
|
List of extraction stages or extractors to run |
|
Strategy for merging results from stages |
|
Stop after first stage that returns entities |
|
Minimum entities to consider a stage successful |
|
Continue to next stage if current stage errors |
|
Optional |
|
Optional confidence floor for the merged result. Entities and relations below it are dropped, as are relations left without an endpoint. |
| Strategy | Behavior |
|---|---|
|
Keep unique entities across results. |
|
Keep entities found by multiple results. |
|
Keep the highest-confidence entity for each matching name/type. |
|
Keep earlier-stage entities and fill gaps from later results. |
|
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:
-
With
min_confidenceset, 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. -
With
ontologyset, relations the ontology does not permit are dropped withExtractionResult.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 to extract from |
|
Optional list of entity types to extract |
|
Whether to extract relations |
|
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 |
|---|---|
|
List of texts to extract from |
|
Number of texts to process in each batch (for memory management) |
|
Maximum number of concurrent extractions |
|
Optional list of entity types to extract |
|
Whether to extract relations |
|
Whether to extract preferences |
|
Optional callback called after each text is processed. Receives (completed_count, total_count). |
|
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 |
|---|---|---|---|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
PipelineResult
Fields and defaults:
| Field | Type | Default | Description |
|---|---|---|---|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
BatchItemResult
Fields and defaults:
| Field | Type | Default | Description |
|---|---|---|---|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
|
BatchExtractionResult
Fields and defaults:
| Field | Type | Default | Description |
|---|---|---|---|
|
|
|
|
|
|
|
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.