Add lazily-evaluated pipeline DSL and local interpreter - #618
Merged
Merged
Conversation
jonnylaw
marked this pull request as draft
August 24, 2026 15:15
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):
|
jonnylaw
marked this pull request as ready for review
August 26, 2026 09:20
williedoran-neo4j
approved these changes
Aug 26, 2026
williedoran-neo4j
left a comment
Contributor
There was a problem hiding this comment.
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
This file contains hidden or bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Sign up for free
to join this conversation on GitHub.
Already have an account?
Sign in to comment
Add this suggestion to a batch that can be applied as a single commit.This suggestion is invalid because no changes were made to the code.Suggestions cannot be applied while the pull request is closed.Suggestions cannot be applied while viewing a subset of changes.Only one suggestion per line can be applied in a batch.Add this suggestion to a batch that can be applied as a single commit.Applying suggestions on deleted lines is not supported.You must change the existing code in this line in order to create a valid suggestion.Outdated suggestions cannot be applied.This suggestion has been applied or marked resolved.Suggestions cannot be applied from pending reviews.Suggestions cannot be applied on multi-line comments.Suggestions cannot be applied while the pull request is queued to merge.Suggestion cannot be applied right now. Please check back later.
Description
Adds a lazily-evaluated dataflow pipeline DSL (
neo4j_graphrag.pipeline) plus a local interpreter to run it.A
Pipelineis pure data — building one produces an operator graph and executes nothing. Evaluation is deferred to anInterpreter, 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.What's included
pipeline.py— thePipeline/ResultPipelinebuilder API.operators.py— the operator graph:Map,FlatMap,Filter,Take,TakeWhile,Skip,Grouped,ReduceOp,MapAsyncChunked, theTry*error-capturing variants, the*Okvariants that operate inside aResult, plusTee/Partitionfor fan-out.interpreter.py—Interpreterprotocol andLocalInterpreter, 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— theSourceandSinkboundary protocols.Type of 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?
tests/unit/pipeline/test_pipeline.py(~870 lines) covers the builder, each operator, laziness/short-circuiting, error propagation through theTry*and*Okvariants, and the tee/partition fan-out paths.Checklist
docs/entry yetSimpleKGBuilderbuilt on this DSL is the planned follow-up PRNote on review order
This is the first of two changes. The follow-up rewrites
SimpleKGBuilderon top of this DSL and adds a worked example; it's parked onpipeline-dsl-exampleuntil 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.