Skip to content

Add lazily-evaluated pipeline DSL and local interpreter - #618

Merged
jonnylaw merged 7 commits into
neo4j:mainfrom
jonnylaw:pipeline-dsl
Aug 26, 2026
Merged

jonnylaw merged 7 commits into
neo4j:mainfrom
jonnylaw:pipeline-dsl

Conversation

@jonnylaw

Copy link
Copy Markdown
Contributor

Description

Adds a lazily-evaluated dataflow pipeline DSL (neo4j_graphrag.pipeline) plus a local interpreter to run it.

A Pipeline is pure data — building one produces an operator graph and executes nothing. Evaluation is deferred to an Interpreter, triggered by .collect(), .to_sink(), iteration, or by passing an interpreter explicitly. This separation means the same pipeline definition can later be run by alternative interpreters (batched, distributed, instrumented) without touching the definition site.

from neo4j_graphrag.pipeline import Pipeline

results = (
    Pipeline.from_source(my_source)
    .map(transform)
    .flat_map(expand)
    .map_async_chunked(async_transform)
    .reduce(zero=identity, combine=merge)
    .collect()
)

What's included

  • pipeline.py — the Pipeline / ResultPipeline builder API.
  • operators.py — the operator graph: Map, FlatMap, Filter, Take, TakeWhile, Skip, Grouped, ReduceOp, MapAsyncChunked, the Try* error-capturing variants, the *Ok variants that operate inside a Result, plus Tee / Partition for fan-out.
  • interpreter.py — Interpreter protocol and LocalInterpreter, which evaluates a graph in-process with bounded-concurrency chunking for the async operators.
  • result.py — Result / Ok / Err, so a failing element becomes a value in the stream rather than an exception that tears down the whole run.
  • source.py / sink.py — the Source and Sink boundary protocols.

Type of Change

  • New feature
  • Bug fix
  • Breaking change
  • Documentation update
  • Project configuration change

Purely additive — a new subpackage with no changes to existing modules, so nothing currently importable changes behaviour.

Complexity

Complexity: High — new subsystem, ~1,500 lines of implementation and a fair amount of design surface in the operator set.

How Has This Been Tested?

  • Unit tests
  • E2E tests
  • Manual tests

tests/unit/pipeline/test_pipeline.py (~870 lines) covers the builder, each operator, laziness/short-circuiting, error propagation through the Try* and *Ok variants, and the tee/partition fan-out paths.

Checklist

  • Documentation has been updated — module and method docstrings are in place, but there's no docs/ entry yet
  • Unit tests have been updated
  • E2E tests have been updated — not applicable, no I/O in this PR
  • Examples have been updated — a SimpleKGBuilder built on this DSL is the planned follow-up PR
  • New files have copyright header
  • CLA (https://neo4j.com/developer/cla/) has been signed
  • CHANGELOG.md updated if appropriate

Note on review order

This is the first of two changes. The follow-up rewrites SimpleKGBuilder on top of this DSL and adds a worked example; it's parked on pipeline-dsl-example until this lands. Reviewing the DSL on its own first is the intent — happy to open the second as a draft now if seeing the consumer helps judge the API.

@jonnylaw
jonnylaw requested a review from a team as a code owner August 24, 2026 15:14
@jonnylaw
jonnylaw marked this pull request as draft August 24, 2026 15:15
Comment thread src/neo4j_graphrag/pipeline/pipeline.py
Comment thread src/neo4j_graphrag/pipeline/pipeline.py Outdated
Comment thread src/neo4j_graphrag/pipeline/operators.py Outdated
Comment thread src/neo4j_graphrag/pipeline/sink.py
Comment thread src/neo4j_graphrag/pipeline/interpreter.py
Comment thread tests/unit/pipeline/test_pipeline.py Outdated
Comment thread src/neo4j_graphrag/pipeline/pipeline.py
@williedoran-neo4j

williedoran-neo4j commented Aug 26, 2026 •

Copy link
Copy Markdown
Contributor

This is a great update to the library. Adding a lazily-evaluated, interpreter-backed pipeline gives us a real foundation for the future, not just a nicer API, but a definition/evaluation split that leaves the door open to scale this out to genuinely large volumes.

I dont know if this is the right place for this... please point me to where i can contribute these ideas @stellasia

A few things I'd love to see us start dreaming about from here (not blockers, just directions):

  • AsyncInterpreter: one event loop across the whole chain, so loop-bound clients (httpx.AsyncClient) span the run instead of per-chunk asyncio.run.
  • Pluggable executor under map_async_chunked, swap the chunk dispatch (event loop → process pool → actor cluster) without touching the operator graph. Most of a KG build is embarrassingly-parallel map, and those stages should scale by "add cores," not "add a cluster."
  • Keyed reduce / group-by. ReduceOp is an unkeyed fold today, but entity resolution is a keyed reduction. That's the operator gap that matters once resolution outgrows a single node, a measured threshold, not a default.

@jonnylaw
jonnylaw marked this pull request as ready for review August 26, 2026 09:20

@williedoran-neo4j williedoran-neo4j left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

my comments were resolved, and added a few nuggets to start thinking about where and how we can use these changes to scale out workflows

@jonnylaw
jonnylaw merged commit c3099b0 into neo4j:main Aug 26, 2026
11 checks passed
Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

2 participants