Process documents in batch
|
Available on NAMS: No. The procedure on this page needs the Bolt backend. With the hosted NAMS backend its operations are unavailable: depending on the call, the client raises |
How to efficiently process large document collections to build your context graph.
Prerequisites
-
Install
pip install 'neo4j-agent-memory[gliner2]==0.7.0'for the GLiNER2.5 examples and allow the configured model (fastino/gliner2.5-base-v1by default) to download before an offline run. -
Configure a connected client with
MemorySettingsandasync with MemoryClient(settings) as client:as in Store, retrieve, and summarize messages. Storage snippets below run inside that block. -
Use a test database and synthetic documents first. Choose positive batch sizes and concurrency limits that fit available memory.
The first sections use documents: list[str]. The financial and catalog variations define their own inputs. GLiNER2.5 decodes relations only when its ontology declares relationship types, and the business template used on this page declares none. Persist relations only when the chosen pipeline actually returns them.
1. Extract and inspect a batch
Account for batch and chunk extraction results covers ExtractionPipeline.extract_batch in full, with a maintained runnable program: progress tracking, per-item result access, and the distinction between a failed item and an empty successful result. Read that guide first if you have not already.
The rest of this page assumes an extractor (a GLiNER2Extractor), a pipeline (an ExtractionPipeline wrapping it), a documents: list[str], and the batch result from extracting them. The examples call pipeline.extract_batch(), which returns a BatchExtractionResult with a success flag per item. GLiNER2Extractor.extract_batch() returns a plain list[ExtractionResult] instead; see Tune batch size.
from neo4j_agent_memory.extraction import ExtractionPipeline, GLiNER2Extractor
extractor = GLiNER2Extractor.for_schema("business")
pipeline = ExtractionPipeline(stages=[extractor], fallback_on_error=False)
# List of documents to process
documents = [
"Customer John ordered Nike Air Max shoes from our Manhattan store.",
"Jane returned the Adidas Ultraboost due to sizing issues.",
"Order #12345 shipped via FedEx to Brooklyn, NY.",
# ... hundreds more
]
# Batch extraction
result = await pipeline.extract_batch(
texts=documents,
batch_size=10,
max_concurrency=5,
)
print(f"Processed: {result.successful_items}/{result.total_items}")
print(f"Total entities: {result.total_entities}")
print(f"Total relations: {result.total_relations}")
2. Choose chunked extraction for long documents
Account for batch and chunk extraction results also covers StreamingExtractor and chunk sizing. This page only adds the storage step: writing each chunk’s entities as they arrive, instead of waiting for the combined result. The alternative extract(…, deduplicate=True) path is described after the example.
from neo4j_agent_memory.extraction import StreamingExtractor
from pathlib import Path
streamer = StreamingExtractor(
extractor=extractor,
chunk_size=4000, # Characters per chunk
overlap=200, # Overlap to avoid splitting entities
)
long_document = Path("annual_report.txt").read_text(encoding="utf-8")
# Stream results as they're extracted, storing each chunk's entities immediately
async for chunk_result in streamer.extract_streaming(long_document):
if not chunk_result.success:
print(f"Chunk {chunk_result.chunk.index} failed: {chunk_result.error}")
continue
print(f"Chunk {chunk_result.chunk.index}: {chunk_result.entity_count} entities")
for entity in chunk_result.result.entities:
await client.long_term.add_entity(
name=entity.name,
entity_type=entity.type,
)
To store the cross-chunk deduplicated result instead of storing per chunk, call await streamer.extract(long_document, deduplicate=True) and store from the returned result.entities once extraction completes.
GLiNER2.5 also windows any chunk longer than its max_words limit (384 words by default), and a relation whose endpoints fall in different windows is not decoded. When relations matter, keep each chunk under that limit.
3. Store successful batch results
The Bolt write merges entities by exact name and type. As of SDK 0.6.0,
add_entity re-reads the matched node when that MERGE matches an existing
entity, so the returned Entity.id already carries the persisted value;
provenance and relationship writes can use it directly, with no separate
readback required. A batch import can still run many writers against the
same name/type pair concurrently, so define this application helper once for
the remaining examples — it double-checks that exactly one node is stored
before using its ID, which is cheap insurance against a concurrent write
race, not a required workaround:
from uuid import UUID
async def store_entity(client, **kwargs):
entity, dedup_result = await client.long_term.add_entity(**kwargs)
rows = await client.query.cypher(
"MATCH (e:Entity {name: $name, type: $type}) RETURN e.id AS id",
{"name": entity.name, "type": entity.type},
)
if len(rows) != 1:
raise ValueError("Expected exactly one persisted entity for name and type")
return entity.model_copy(update={"id": UUID(rows[0]["id"])}), dedup_result
This extra read guards against concurrent writers racing on the same name/type pair during a large batch import; it is not required to work around SDK ID handling on 0.7.0, and it is not an atomic upsert contract either way. It assumes library UUID entity IDs and no concurrent deletion. Inspect ambiguous matches instead of choosing a semantic nearest neighbor.
Basic storage
# Use the connected Bolt client from the prerequisites.
# Extract batch
result = await pipeline.extract_batch(texts=documents)
# Store all entities
for item in result.results:
if not item.success:
continue
for entity in item.result.entities:
await store_entity(
client,
name=entity.name,
entity_type=entity.type,
attributes={
"confidence": entity.confidence,
"source_index": item.index,
},
)
With deduplication
stored_count = 0
merged_count = 0
for item in result.results:
if not item.success:
continue
for entity in item.result.entities:
stored, dedup_result = await store_entity(
client,
name=entity.name,
entity_type=entity.type,
deduplicate=True,
)
if dedup_result.action == "none":
stored_count += 1
elif dedup_result.action == "merged":
merged_count += 1
print(f"No deduplication match: {stored_count}")
print(f"Merged with existing: {merged_count}")
add_entity accepts the per-call deduplicate boolean, not a
deduplication configuration keyword. Its deduplication actions are none,
merged, and flagged; none is not an insertion-count guarantee. With the
default Bolt resolver, the decision uses the ResolutionConfig bands
(auto_merge_threshold=0.90, review_threshold=0.85), the same ones message
ingestion uses. For thresholds see
Tune entity resolution; for reviewing
flagged pairs see Merge a reviewed entity pair.
With provenance tracking
# Register the extractor for provenance
await client.long_term.register_extractor(
name=extractor.name, # "gliner2"
version=extractor.version, # installed gliner2 package version
config={"model": extractor.model_id, "schema": "business"},
)
# Store with provenance
for item in result.results:
doc_index = item.index
if not item.success:
continue
# Store document as a message/source
message = await client.short_term.add_message(
session_id="batch-import",
role="system",
content=documents[doc_index],
metadata={"doc_index": doc_index},
extraction_mode="skip", # This batch already performed extraction.
)
for entity in item.result.entities:
stored, _ = await store_entity(
client,
name=entity.name,
entity_type=entity.type,
)
# Link to source document
linked = await client.long_term.link_entity_to_message(
entity=stored,
message_id=message.id,
confidence=entity.confidence,
start_pos=entity.start_pos,
end_pos=entity.end_pos,
)
if not linked:
raise RuntimeError("Source provenance link was not persisted")
# Link to extractor
await client.long_term.link_entity_to_extractor(
entity=stored,
extractor_name=extractor.name,
confidence=entity.confidence,
)
To let the client extract instead of the batch above, store the documents with
client.short_term.add_messages_batch(session_id, messages, extract_entities=True),
where each message is a {"role": …, "content": …} dict. It runs the
client’s configured extractor, resolves the extracted entities against stored
entities by default (resolution.resolve_on_ingest=True), and links them with
MENTIONS rather than EXTRACTED_FROM. Messages without a timestamp are
stored in input order.
Optional: Process financial documents
Example for processing financial documents:
from neo4j_agent_memory.extraction import ExtractionPipeline, GLiNER2Extractor
# Configure for financial domain
pipeline = ExtractionPipeline(
stages=[GLiNER2Extractor.for_schema("business")],
fallback_on_error=False,
)
async def process_financial_documents(documents: list[dict]):
"""Process financial documents and build context graph."""
# Extract texts
texts = [doc["content"] for doc in documents]
# Batch extraction
result = await pipeline.extract_batch(
texts=texts,
batch_size=10,
on_progress=lambda d, t: print(f"Extracting: {d}/{t}"),
)
# Store in context graph
for item in result.results:
i = item.index
if not item.success:
print(f"Failed: {documents[i]['id']} - {item.error}")
continue
doc = documents[i]
# Store document metadata as entity
doc_entity, _ = await store_entity(
client,
name=doc.get("title", f"Document {doc['id']}"),
entity_type="DOCUMENT",
attributes={
"doc_id": doc["id"],
"doc_type": doc.get("type"),
"date": doc.get("date"),
},
)
# Keep exact names mapped to the IDs returned by storage.
stored_by_name = {}
# Store extracted entities with relationships
for entity in item.result.entities:
stored, _ = await store_entity(
client,
name=entity.name,
entity_type=entity.type,
)
stored_by_name.setdefault(entity.name, []).append(stored.id)
# Link entity to document
await client.long_term.add_relationship(
source=stored.id,
target=doc_entity.id,
relationship_type="MENTIONED_IN",
confidence=entity.confidence,
)
# Store relationships if extracted
for relation in item.result.relations:
# Do not use semantic nearest-neighbor search to resolve write IDs.
source_ids = stored_by_name.get(relation.source, [])
target_ids = stored_by_name.get(relation.target, [])
if len(source_ids) == len(target_ids) == 1:
await client.long_term.add_relationship(
source=source_ids[0],
target=target_ids[0],
relationship_type=relation.relation_type,
confidence=relation.confidence,
)
else:
print(f"Review ambiguous or missing relation endpoints: {relation}")
return result
# Usage
documents = [
{
"id": "10K-2024",
"title": "Annual Report 2024",
"type": "SEC Filing",
"date": "2024-03-15",
"content": "Apple Inc. reported record revenue...",
},
{
"id": "earnings-Q4",
"title": "Q4 Earnings Call",
"type": "Transcript",
"date": "2024-01-25",
"content": "CEO Tim Cook discussed...",
},
# ... more documents
]
result = await process_financial_documents(documents)
Optional: Import an ecommerce catalog
Example for a one-time product catalog import. sku is an attribute, not a uniqueness constraint. Two distinct SKUs with the same product name and type can match the same entity, and re-running does not update every attribute. Define an explicit SKU-to-entity identity and update policy before using this as a recurring import.
async def import_product_catalog(products: list[dict]):
"""Import product catalog into context graph."""
# Prepare product descriptions for extraction
texts = []
for product in products:
text = f"""
Product: {product['name']}
Brand: {product.get('brand', 'Unknown')}
Category: {product.get('category', 'Unknown')}
Description: {product.get('description', '')}
"""
texts.append(text)
# Extract additional entities from descriptions
pipeline = ExtractionPipeline(
stages=[GLiNER2Extractor.for_schema("business")],
fallback_on_error=False,
)
result = await pipeline.extract_batch(
texts=texts,
batch_size=20,
)
for i, product in enumerate(products):
# Store product as primary entity
product_entity, _ = await store_entity(
client,
name=product["name"],
entity_type="PRODUCT",
description=product.get("description"),
attributes={
"sku": product["sku"],
"price": product.get("price"),
"rating": product.get("rating"),
},
deduplicate=False, # One-time import only; this does not enforce SKU uniqueness.
)
# Store brand as entity
if product.get("brand"):
brand_entity, _ = await store_entity(
client,
name=product["brand"],
entity_type="BRAND",
)
await client.long_term.add_relationship(
source=product_entity.id,
target=brand_entity.id,
relationship_type="MADE_BY",
)
# Store category as entity
if product.get("category"):
category_entity, _ = await store_entity(
client,
name=product["category"],
entity_type="CATEGORY",
)
await client.long_term.add_relationship(
source=product_entity.id,
target=category_entity.id,
relationship_type="IN_CATEGORY",
)
# Link extracted entities to product. The business schema returns the
# product as OBJECT:PRODUCT and the brand as ORGANIZATION:COMPANY, so skip
# them by name rather than by the PRODUCT/BRAND types stored above.
already_stored = {product["name"].casefold(), (product.get("brand") or "").casefold()}
if result.results[i].success:
for entity in result.results[i].result.entities:
if entity.name.casefold() in already_stored:
continue
linked, _ = await store_entity(
client,
name=entity.name,
entity_type=entity.type,
)
await client.long_term.add_relationship(
source=product_entity.id,
target=linked.id,
relationship_type="RELATED_TO",
)
print(f"Imported {len(products)} products")
# Usage
products = [
{
"sku": "NKE-AM90-001",
"name": "Nike Air Max 90",
"brand": "Nike",
"category": "Running Shoes",
"price": 129.99,
"description": "Classic running shoe with Air cushioning...",
},
# ... more products
]
await import_product_catalog(products)
4. Tune throughput and handle failures
Tune batch size
# Smaller batches for memory-constrained environments
result = await pipeline.extract_batch(
texts=documents,
batch_size=5, # Smaller batches
max_concurrency=2,
)
# Larger batches for high-memory environments
result = await pipeline.extract_batch(
texts=documents,
batch_size=50, # Larger batches
max_concurrency=10,
)
ExtractionPipeline.extract_batch() still sends each text through the stages
separately, up to max_concurrency at once; batch_size only groups that
work. To decode several texts
in one forward pass on a GPU, call GLiNER2Extractor.extract_batch() directly.
It returns one ExtractionResult per input, in input order, with no
max_concurrency, fail_fast, or per-item success flag: an error in a slice
raises out of the call.
gpu_extractor = GLiNER2Extractor.for_schema("business", device="cuda")
results = await gpu_extractor.extract_batch(documents, batch_size=32)
for index, extraction in enumerate(results):
print(f"Doc {index}: {extraction.entity_count} entities")
Choose an extractor for your data
Measure extraction quality and throughput on labeled samples before choosing
spaCy, GLiNER2.5, or an LLM stage. The settings above are starting points, not
hardware-specific performance guarantees. A GLiNER2.5 stage returns
relationships only when its ontology declares relationship types: the poleo,
podcast, and news templates do, and business does not. Use one of those
templates, attach relationships with DomainSchema.to_ontology(relationships=[…]),
or add an LLM stage when your goal requires them.
Skip failures and continue
result = await pipeline.extract_batch(
texts=documents,
fail_fast=False, # Continue on errors (default)
)
# Check for failures
failed = [item for item in result.results if not item.success]
print(f"Failed: {len(failed)}/{result.total_items}")
for item in failed:
print(f" Index {item.index}: {item.error}")
Optional: Prepare larger imports
1. Validate before processing
# Filter out invalid documents
valid_docs = [
doc for doc in documents
if doc.get("content") and len(doc["content"]) > 10
]
print(f"Processing {len(valid_docs)} of {len(documents)} documents")
result = await pipeline.extract_batch(texts=[d["content"] for d in valid_docs])
2. Log progress and errors
import logging
logging.basicConfig(level=logging.INFO)
logger = logging.getLogger(__name__)
def on_progress(completed, total):
logger.info(f"Progress: {completed}/{total}")
result = await pipeline.extract_batch(
texts=documents,
on_progress=on_progress,
)
for item in result.results:
if not item.success:
logger.error(f"Failed doc {item.index}: {item.error}")
3. Process in manageable chunks
# For very large document sets, process in chunks
CHUNK_SIZE = 1000
for i in range(0, len(all_documents), CHUNK_SIZE):
chunk = all_documents[i:i + CHUNK_SIZE]
print(f"Processing chunk {i//CHUNK_SIZE + 1}")
result = await pipeline.extract_batch(texts=chunk)
# Apply the successful-item storage loop from step 3 before the next chunk.
print(result.successful_items, result.get_errors())
4. Use appropriate deduplication
# For known-clean data, skip deduplication
await store_entity(
client,
name=entity.name,
entity_type=entity.type,
deduplicate=False, # Skips similarity deduplication; does not make imports idempotent.
)
# For user-generated content, enable deduplication
await store_entity(
client,
name=entity.name,
entity_type=entity.type,
deduplicate=True,
)
5. Verify the import
Check result.failed_items and result.get_errors() before declaring the
batch complete. For the provenance example, query the session you imported:
MATCH (c:Conversation {session_id: 'batch-import'})-[:HAS_MESSAGE]->(m:Message)
OPTIONAL MATCH (e:Entity)-[:EXTRACTED_FROM]->(m)
RETURN m.id, m.content, collect(DISTINCT e.name) AS entities
ORDER BY m.content
Expect one source message for each successfully processed input, including
inputs with zero extracted entities. These links use EXTRACTED_FROM; storing
a source message with extraction skipped does not create MENTIONS links. Verify a labeled sample of entity names
and relation endpoints. A successful extraction call or nonzero entity count
does not establish extraction quality, persistence, or retry safety. Record
source IDs and completed batches outside this example before adding retries.
See also
-
Account for batch and chunk extraction results — the
extract_batch/StreamingExtractorAPI and its result bookkeeping in depth; this page covers ingesting and storing the documents that feed it.