Account for batch and chunk extraction results

Runs client-side. The extractors on this page run in your Python process and do not call a memory backend, so they work next to either backend. Configuring them as a MemoryClient extraction pipeline applies only to Bolt: NAMS extracts entities server-side and ignores the client’s extraction settings with a UserWarning. See the backend capabilities reference.

To process multiple source texts without losing failed inputs, use the pipeline’s batch result wrappers and retain the input index for every result. For longer individual texts, select explicit chunk units and inspect per-chunk errors.

For the rationale, see Understanding the extraction pipeline.

1. Install the local extractor

Use Python 3.10 or newer and a POSIX shell. Prepare the example files and install the required extra:

Create a local folder and virtual environment for the examples on this page. The commands reuse an existing environment without changing its files:

mkdir -p ~/agent-memory-tutorials
cd ~/agent-memory-tutorials
if [ -e .venv ]; then
  printf '%s\n' 'Using the existing virtual environment.'
else
  python3 -m venv .venv
fi
source .venv/bin/activate

Expected: ~/agent-memory-tutorials is your working directory and its virtual environment is active. Install the published SDK with the command below.

Each complete code block labelled Save as names a file to create in this folder using your editor. Copy the entire block, including imports and the entry point. Expand each helper disclosure and use Copy code to copy its full source. Keep all files together so their imports resolve.

When continuing from another tutorial or guide, retain the existing environment, configuration, session files, and .tutorial-state/. Reuse unchanged helper files; compare an existing file before replacing it, and finish any pending cleanup or recovery before changing the code that owns its state.

python -m pip install 'neo4j-agent-memory[gliner2]==0.7.0'

The maintained program extraction_recipes.py supplies its own imports, source texts, and event-loop entry point. It downloads the GLiNER2.5 checkpoint fastino/gliner2.5-base-v1 (about 407 MB) on first inference; allow disk space and network access for the download. These commands do not connect to a database or require an LLM API key.

Save the complete program below in agent-memory-tutorials/. It has no local helper imports. The next section explains the selected function.

Complete extraction_recipes.py
Save as extraction_recipes.py
"""Complete local-model extraction recipes without a database or LLM API."""

import argparse
import asyncio

from neo4j_agent_memory.extraction.domain_schemas import DomainSchema, get_schema
from neo4j_agent_memory.extraction.gliner2_extractor import GLiNER2Extractor
from neo4j_agent_memory.extraction.pipeline import ExtractionPipeline, MergeStrategy
from neo4j_agent_memory.extraction.streaming import StreamingExtractor
from neo4j_agent_memory.ontology import RelationshipDef

TEXTS = [
    "Maya Chen works at Northstar Robotics in Denver.",
    "Ravi Shah joined Summit Research in Boulder.",
]


# tag::schema[]
def custom_extractor():
    schema = DomainSchema(
        name="support_catalog",
        entity_types={"customer": "A named customer", "product": "A named purchased item"},
    )
    return GLiNER2Extractor(
        ontology=schema,
        label_mapping={"customer": ("PERSON", None), "product": ("OBJECT", "PRODUCT")},
        threshold=0.5,
    )


# end::schema[]


# tag::extract[]
async def extract(selected, text):
    result = await selected.extract(text, extract_relations=False, extract_preferences=False)
    if not result.entities:
        raise RuntimeError("No candidate entities; inspect the input/schema/model threshold")
    for entity in result.entities:
        print(entity.name, entity.type, entity.subtype, entity.confidence)
    print(f"Verified: {result.entity_count} candidate entities returned; inspect their accuracy")
    return result


# end::extract[]


# tag::relations[]
def relation_extractor():
    # The business catalog declares labels only; attach typed relationships.
    ontology = get_schema("business").to_ontology(
        relationships=[
            RelationshipDef(
                type="EMPLOYED_BY",
                source="person",
                target="company",
                description="The person works for or has joined this company",
            ),
            RelationshipDef(
                type="LOCATED_IN",
                source="company",
                target="location",
                description="The company is based or operates in this place",
            ),
        ]
    )
    return GLiNER2Extractor.for_ontology(ontology, threshold=0.5)


async def relations(selected, text):
    result = await selected.extract(text, extract_preferences=False)
    if not result.relations:
        raise RuntimeError("No relations; check the declared relationships and thresholds")
    for relation in result.relations:
        print(relation.source, relation.relation_type, relation.target, relation.confidence)
    print(f"Verified: {len(result.relations)} candidate relations returned; inspect them")
    return result


# end::relations[]


# tag::batch[]
async def batch(selected, texts):
    # Propagate stage errors so the batch can distinguish a failed item from an empty extraction.
    pipeline = ExtractionPipeline(
        stages=[selected], merge_strategy=MergeStrategy.CONFIDENCE, fallback_on_error=False
    )
    result = await pipeline.extract_batch(
        texts,
        batch_size=2,
        max_concurrency=1,
        fail_fast=False,
        extract_relations=False,
        extract_preferences=False,
        on_progress=lambda done, total: print(f"Progress: {done}/{total}"),
    )
    assert result.total_items == len(texts)
    for item in result.results:
        if item.success:
            print(f"Input {item.index}: {item.result.entity_count} entities")
        else:
            print(f"Input {item.index} failed: {item.error}")
    if result.failed_items:
        raise RuntimeError(
            f"Retry failed source indexes after fixing errors: {result.get_errors()}"
        )
    print(f"Verified: all {result.total_items} inputs accounted for")
    return result


# end::batch[]


# tag::streaming[]
async def streaming(selected, text):
    streamer = StreamingExtractor(selected, chunk_size=400, overlap=40, chunk_by_tokens=False)
    result = await streamer.extract(text, extract_relations=False)
    errors = [
        (chunk.chunk.index, chunk.error) for chunk in result.chunk_results if not chunk.success
    ]
    if errors:
        raise RuntimeError(f"Chunk extraction failed: {errors}")
    assert result.chunk_results
    combined = result.to_extraction_result(source_text=text)
    print(
        f"Verified: {len(result.chunk_results)} chunks completed; {combined.entity_count} merged entities"
    )
    return result


# end::streaming[]


async def main():
    parser = argparse.ArgumentParser()
    parser.add_argument("command", choices=["extract", "schema", "relations", "batch", "streaming"])
    args = parser.parse_args()
    if args.command == "schema":
        await extract(custom_extractor(), "Maya Chen bought a Trail Starter shoe.")
    elif args.command == "relations":
        await relations(relation_extractor(), TEXTS[0])
    else:
        selected = GLiNER2Extractor.for_schema("business")
        if args.command == "extract":
            await extract(selected, TEXTS[0])
        elif args.command == "batch":
            await batch(selected, TEXTS)
        else:
            await streaming(selected, "\n".join(TEXTS * 12))


if __name__ == "__main__":
    asyncio.run(main())

2. Configure and run the extraction

The runner selects the business template with GLiNER2Extractor.for_schema("business"). ExtractionPipeline(…​, fallback_on_error=False) propagates an extractor error to the batch wrapper; fail_fast=False then collects the item error instead of abandoning the remaining inputs. The example starts with one worker and a batch size of two; tune after measuring your workload.

The following function is included from the complete program. Run the maintained file with the command below it.

async def batch(selected, texts):
    # Propagate stage errors so the batch can distinguish a failed item from an empty extraction.
    pipeline = ExtractionPipeline(
        stages=[selected], merge_strategy=MergeStrategy.CONFIDENCE, fallback_on_error=False
    )
    result = await pipeline.extract_batch(
        texts,
        batch_size=2,
        max_concurrency=1,
        fail_fast=False,
        extract_relations=False,
        extract_preferences=False,
        on_progress=lambda done, total: print(f"Progress: {done}/{total}"),
    )
    assert result.total_items == len(texts)
    for item in result.results:
        if item.success:
            print(f"Input {item.index}: {item.result.entity_count} entities")
        else:
            print(f"Input {item.index} failed: {item.error}")
    if result.failed_items:
        raise RuntimeError(
            f"Retry failed source indexes after fixing errors: {result.get_errors()}"
        )
    print(f"Verified: all {result.total_items} inputs accounted for")
    return result
python extraction_recipes.py batch

3. Verify the result

Expected: progress for both inputs, an Input …​ line for each result, and Verified: all 2 inputs accounted for. Failed items are printed with their source indexes and cause a nonzero exit so they can be retried after correcting the cause. An empty but successful extraction is still an accounted-for item; inspect it for quality.

The command raises on the documented failure conditions. Inspect candidate names, types, and confidence against source text before storing them. Counts alone do not measure extraction accuracy.

Process a long text in chunks

  1. Read the maintained streaming function:

    async def streaming(selected, text):
        streamer = StreamingExtractor(selected, chunk_size=400, overlap=40, chunk_by_tokens=False)
        result = await streamer.extract(text, extract_relations=False)
        errors = [
            (chunk.chunk.index, chunk.error) for chunk in result.chunk_results if not chunk.success
        ]
        if errors:
            raise RuntimeError(f"Chunk extraction failed: {errors}")
        assert result.chunk_results
        combined = result.to_extraction_result(source_text=text)
        print(
            f"Verified: {len(result.chunk_results)} chunks completed; {combined.entity_count} merged entities"
        )
        return result
  2. Run it against the program’s repeated fictional sample text:

    python extraction_recipes.py streaming
  3. Verify Verified: …​ chunks completed; …​ merged entities. Any chunk error raises with its chunk index; inspect candidate quality separately.

chunk_size=400 and overlap=40 in this example are characters. The default is 4,000 characters with 200 overlap. With chunk_by_tokens=True, this implementation uses approximate whitespace token counts, not the exact tokenizer of every model. Choose a size based on the selected model’s accepted input, boundary behavior, memory use, latency, and a representative corpus; there is no universal 2,000/10,000/100,000-token cutoff.

GLiNER2.5 also windows any single text longer than its max_words limit (default 384 words, counted as the model splits them, with each punctuation mark a word of its own) into windows of that size that overlap by chunk_overlap words (default 64). A relation whose endpoints fall in different windows is lost without a warning. When you need relations, keep each chunk below max_words and split on paragraph or turn boundaries so both endpoints stay in one chunk.

The API accepts the complete input string and computes its chunks. It is not a file-streaming reader that avoids loading the document text into memory. extract_streaming yields StreamingChunkResult objects with .result.entities; extract returns StreamingExtractionResult, and to_extraction_result converts it to the common result model. Calling extraction twice repeats model work.

Preserve source indexes and measure quality

The batch API shown here belongs to ExtractionPipeline. GLiNER2Extractor.extract_batch has a different signature: it takes batch_size and on_progress but no max_concurrency or fail_fast, decodes each slice of texts in one forward pass, and returns a list of ExtractionResult in input order, not the aggregate wrapper. An error in any slice propagates out of the call. Keep the two APIs distinct.

Before persistence, associate each successful BatchItemResult.index with its source document; log errors against that same input identifier. The document program demonstrates message provenance and safe relationship endpoint matching. Measure precision/recall on reviewed samples and inspect chunk boundaries; a 100% completed-item rate is not a claim of extraction accuracy.