Streaming extraction

Chunk long documents and combine extraction results, from Extractor classes reference. Import StreamingExtractor from neo4j_agent_memory.extraction.

StreamingExtractor

Chunk long documents and combine extraction results. Import from neo4j_agent_memory.extraction. The constructor defaults are DEFAULT_CHUNK_SIZE (4000) and DEFAULT_OVERLAP (200) in either chunking mode, so chunk_by_tokens=True without an explicit chunk_size gives 4000-token chunks. Pass chunk_size and overlap with chunk_by_tokens=True, or use create_streaming_extractor(extractor, chunk_by_tokens=True), which defaults to 1000 tokens with a 50-token overlap.

def StreamingExtractor(
    extractor: EntityExtractor,
    *,
    chunk_size: int=DEFAULT_CHUNK_SIZE,
    overlap: int=DEFAULT_OVERLAP,
    chunk_by_tokens: bool=False,
    split_on_sentences: bool=True,
): ...
Parameter Description

extractor

Base extractor to use for each chunk

chunk_size

Size of each chunk (chars or tokens)

overlap

Overlap between chunks (chars or tokens)

chunk_by_tokens

If True, chunk by tokens instead of characters

split_on_sentences

Try to split on sentence boundaries

chunk_document

Chunk a document according to settings.

def chunk_document(text: str) -> list[ChunkInfo]: ...
Parameter Description

text

Text to chunk

extract_streaming

Stream extraction results chunk by chunk.

async def extract_streaming(
    text: str,
    *,
    entity_types: list[str] | None=None,
    extract_relations: bool=True,
    on_chunk_complete: Callable[[StreamingChunkResult], None] | None=None,
) -> AsyncIterator[StreamingChunkResult]: ...
Parameter Description

text

Text to extract from

entity_types

Optional list of entity types to extract

extract_relations

Whether to extract relations

on_chunk_complete

Optional callback for each completed chunk

Yields one StreamingChunkResult for each chunk.

extract

Extract from a long document with automatic chunking.

async def extract(
    text: str,
    *,
    entity_types: list[str] | None=None,
    extract_relations: bool=True,
    deduplicate: bool=True,
    on_progress: Callable[[int, int], None] | None=None,
) -> StreamingExtractionResult: ...
Parameter Description

text

Text to extract from

entity_types

Optional list of entity types to extract

extract_relations

Whether to extract relations

deduplicate

Whether to deduplicate entities across chunks

on_progress

Optional progress callback (completed_chunks, total_chunks)

extract_to_result

Extract and return standard ExtractionResult.

async def extract_to_result(
    text: str,
    *,
    entity_types: list[str] | None=None,
    extract_relations: bool=True,
    deduplicate: bool=True,
    on_progress: Callable[[int, int], None] | None=None,
) -> ExtractionResult: ...
Parameter Description

text

Text to extract from

entity_types

Optional list of entity types to extract

extract_relations

Whether to extract relations

deduplicate

Whether to deduplicate entities

on_progress

Optional progress callback

With deduplicate=True, entities are merged across chunks by normalized name and type, and relations by source name, relation type and target name, keeping the highest-confidence copy of each.

Each chunk is a separate extraction call, so relations whose endpoints fall in different chunks are not found, and mention ids (ExtractedEntity.id, ExtractedRelation.source_id and target_id) are scoped to the chunk that produced them. GLiNER2Extractor windows any input longer than its max_words (384 words by default) itself, so a chunk longer than that is split again inside the extractor. Use StreamingExtractor when you want chunk-level results, progress reporting or control over the boundaries.

from neo4j_agent_memory.extraction import GLiNER2Extractor, StreamingExtractor

extractor = GLiNER2Extractor.for_schema("podcast")
streamer = StreamingExtractor(extractor, chunk_size=4000, overlap=200)
async for chunk_result in streamer.extract_streaming(long_document):
    print(chunk_result.chunk.index, chunk_result.entity_count)
# extract returns StreamingExtractionResult; extract_to_result returns ExtractionResult.
result = await streamer.extract_to_result(long_document, deduplicate=True)