/

Graph Computing

/

How We Orchestrate Graph Computing the Database Way

How We Orchestrate Graph Computing the Database Way

Bin Huang

Graph Computing

About the Author: Bin Huang is the Tech Lead for the NebulaGraph Database kernel and the lead architect of NebulaGraph Analytics. His work focuses on database architecture, query processing and optimization, distributed systems, and large-scale graph computing. He leads the design and evolution of database kernel and graph analytics technologies for high-performance graph workloads.

Synopsis—What This Post Covers

This post examines how NebulaGraph Analytics approaches graph computing as a database capability rather than as a standalone programming framework. Using PageRank as a running example, it traces the design from graph-computing models and GQL procedures to vectorized execution and distributed state management.

  • From graph APIs and DSLs to GQL procedures: How ideas from Pregel, PowerGraph, frontier-oriented systems, and graph DSLs influenced the design of NebulaGraph Analytics.

  • A language contract for graph algorithms: How constructs such as MATCH, PER NODE, PER PATH, active sets, and aggregators expose optimization-relevant semantics while preserving procedural control.

  • A shared database substrate: How graph analytics reuses database capabilities such as parsing, validation, planning, typed expressions, vectors, and pipelines, while adding analytics-specific runtime extensions.

  • Hybrid execution: Why NebulaGraph Analytics interprets procedure-level control flow while executing data-intensive graph operations through vectorized worker pipelines.

  • Distributed state and communication: How state ownership, aggregation semantics, and visibility boundaries allow the runtime to derive synchronization and communication without requiring users to manage machines or messages explicitly.

  • Optimization opportunities: How preserving algorithm intent gives the engine room to optimize frontier handling, state synchronization, local aggregation, execution strategies, and physical data movement.

  • Design boundaries and trade-offs: Where synchronization, aggregation, vectorization, temporary state, and synchronous execution introduce costs or constraints.

The central idea is that graph computing can benefit from the same separation of responsibilities long used in database systems: the user expresses algorithmic intent and semantics, while the engine owns physical execution and distributed communication.

Introduction

Graph computing is usually introduced through an algorithm or a programming model. This article takes a graph algorithm as a starting point to explore a broader systems design question: what needs to change when graph computing is built as a database capability?

For NebulaGraph Analytics, the answer spans the language contract, the frontend shared with the database, vectorized execution and distributed state movement. This article follows these design choices and explains why we choose GQL procedures as the boundary between algorithm authors and the underlying distributed machinery.

At the algorithm level, PageRank is straightforward: each node sends a fraction of its score along outgoing edges, the destination nodes combine the contributions they receive, and the process repeats until the scores stabilize.

The distributed program behind that description is considerably more complex. Which machine owns a node's state? Which machines need a readable copy? When do contributions become visible? Who decides when an iteration has completed? Which parts of the work should be executed as graph traversal, which belong to expression evaluation, and which belong to communication?

In NebulaGraph Analytics, the core iteration of the same algorithm can be expressed as a short procedure fragment:

The fragment exposes the structure that matters to the execution engine: iterative control, propagation along the graph topology, per-node state transitions, a global convergence value, and the boundaries between different phases of computation.

The procedure is the end result of a much longer engineering process. We began by studying graph-computing APIs and execution models, then worked toward a GQL-based language contract, and finally connected that contract to the database execution framework and distributed analytics state. This article focuses on how those pieces fit together as a system.

The central idea in our design is to give graph computing the same separation of responsibilities that databases have long relied on: users describe intent in a structured form, and the engine owns physical execution. The challenge is that graph algorithms require iteration, mutable temporary state, dynamically changing active frontiers, and topology-dependent communication—concepts that ordinary query languages cannot fully express.

Graph Computing Followed More Than One Path

The history of graph computing is sometimes compressed into a progression from API-based programming models to graph-specific DSLs. That framing overlooks the fact that these systems were built to solve different problems in the first place.

Pregel introduced a vertex-centric bulk-synchronous model. A vertex executes a user function, consumes messages, changes its local state, emits new messages, and participates in a sequence of supersteps. The framework hides message delivery, synchronization, placement, and failure handling. “Think like a vertex” worked because it turned a distributed system into a local programming abstraction plus a global protocol.

The abstraction also limits what the engine can observe. A callback may read several fields, construct arbitrary messages, and update local state, but the framework does not necessarily know the algebra behind those operations. It may know only that bytes are moving between vertices, without knowing whether those bytes represent a minimum, a sum, an idempotent flag, or a value replacement. Only with that information can the engine determine whether computation can be combined, reordered, or omitted.

PowerGraph addressed a different problem: the skewed degree distributions common in real-world graphs. A high-degree vertex can become a hotspot for both data placement and communication. PowerGraph's gather-apply-scatter model exposed more of the computation's structure, while its master/mirror mechanism allowed the incident edges of a logical vertex to be distributed across multiple machines.

Frontier-oriented systems followed another path. Ligra made vertex subsets explicit and offered operations over active vertices and their edges. This fits algorithms such as BFS and SSSP, where each iteration touches only a changing subset of the graph. It also provides the basis for choosing between sparse and dense execution: the representation of the current frontier can influence whether the engine should push outward from active vertices or pull over a broader portion of the graph.

Graph DSLs explored how to expose more algorithmic structure to the compiler. Green-Marl provided graph-specific constructs that could be translated into parallel code. GraphIt went further by separating the algorithm language from a scheduling language, allowing traversal direction, parallelization strategy, and data layout to change without rewriting the algorithm. TigerGraph GSQL took an SQL-inspired route, combining graph patterns, procedural stages, and global or vertex-attached accumulators, with code generation for installed queries.

These systems are better understood as parallel answers to different questions:

  • What is the basic unit of computation: a vertex, an edge, a frontier, or a matched pattern?

  • Which state belongs to the algorithm, and where is it attached?

  • Which synchronization processes need to be visible to the user?

  • How much of the algorithm's semantics can the engine inspect?

  • Which physical execution decisions can still be changed after the algorithm is written?

Our design is influenced by all of these paths. Pregel showed us the value of explicit iteration boundaries. PowerGraph provided a practical way to reason about distributed vertex identity. Frontier systems showed that the active set should not be hidden inside a callback. Graph DSLs demonstrated that an algorithm can retain procedural structure while exposing enough semantics for compilation and optimization.

Figure 1. The design draws from several branches; it is not the final step of a single historical sequence.

From MapReduce to SQL—and from Graph APIs to GQL

The evolution of distributed data processing offers a useful analogy. MapReduce let developers provide map and reduce functions while its runtime handled partitioning, scheduling, failure recovery, and communication. It made distributed programs easier to write, but in multi-stage applications, much of the execution structure still remained hidden insideuser-written functions and job sequences.

Spark SQL reflects the same design direction: Spark applications can express structured computation through SQL or DataFrames, allowing Catalyst and the physical engine to optimize it rather than treating every operation as an opaque RDD transformation.

The common direction is toward higher-level language representations for distributed computation, giving the engine optimization opportunities that low-level programming primitives cannot provide. We  apply the same idea to graph analytics: moving from vertex/edge callbacks toward GQL procedures that preserve iteration while exposing graph patterns, state scopes, aggregation semantics, and phase boundaries. An engine can optimize only the intent that survives the programming interface.

Graph analytics is still harder to express than a relational dataflow because it requires repeated updates to temporary state over an irregular topology, while the active frontier keeps changing. We ultimately chose not to write graph algorithms in SQL. Instead, we designed a procedure language that preserves algorithmic control logic while giving graph computation and state changes explicit, typed semantics.

There is also a standards boundary. ISO/IEC 39075:2024 standardized GQL as a property graph database language in 2024. The graph-computing constructs discussed here—including PER NODE, PER PATH, active sets, and node-attached aggregators—are NebulaGraph extensions built on top of GQL's data language and schema language.

The Language Preserves Optimization-Relevant Intent

The most important property of the procedure is not that it resembles SQL. It is that the language makes several dimensions of the algorithm explicit at the same time.

WHILE and IF preserve phase control and termination logic. MATCH describes topology-oriented work. PER PATHdenotes propagation or computation over the path bindings produced by MATCH. PER NODE denotes a state transition semantically owned by a node. Active sets describe the changing set of candidate nodes for subsequent computation. Global and node-attached aggregators describe where state lives and how concurrent contributions are combined.

These constructs are not intended to turn every graph algorithm into a declarative query. Procedures remain ordered and stateful. Together, they form a semantic intermediate representation between two undesirable extremes:

  • native callbacks that are flexible but expose too little semantic information to the engine;

  • a fixed catalog of built-in algorithms that can be fast but offers little programmability.

A PER PATH block processes path bindings produced by MATCH and can contribute to mergeable state. Because many paths may target the same logical node, direct replacement of node state is not allowed there; contributions become visible only after the task completes. A PER NODE block is responsible for state transitions within a node's scope, so it can explicitly replace that node's state. Global contributions become visible after the match-compute phase completes.

Those visibility rules tell the engine more than where a loop executes. They also specify which state must remain stable while parallel tasks are running, which updates can be combined locally, and where synchronization boundaries are semantically required.

The aggregator type provides another layer of information. SumAgg, MinAgg, MaxAgg, logical aggregators, and container aggregators do not merely store values; each defines its own merge behavior. The runtime no longer has to treat every update as an arbitrary message handler. It can build local partial results, combine them in stages, and then route node-attached state according to graph ownership.

Active sets provide the execution engine with similar information. A BFS frontier represented as an active set is visible both to the algorithm and to the planner. It can restrict the starting points of a scan, define the nodes collected for the next phase, and provide a basis for future sparse-versus-dense execution choices. If the same frontier existed only inside user-managed arrays, the graph engine would have much less room to participate in those decisions.

In practice, language design and systems optimization are connected. A language abstraction is valuable when it hides accidental complexity while preserving the information the runtime needs.

A Database Substrate, Not an Identical Query Engine

Once graph computation is represented as a language program, much of the infrastructure already present in the NebulaGraph database can be reused.

The parser builds a common syntax tree. The validator resolves names, graph bindings, scopes, and types before distributed execution begins. The planner translates graph patterns into scan, expand, and filter operations. The expression engine evaluates typed predicates and arithmetic expressions. Vector and pipeline infrastructure moves data through execution in batches. Explainable plan structures let us inspect and reason about execution before it runs.

NebulaGraph Analytics then adds the pieces that ordinary graph queries do not require: procedure-level iteration, node-attached temporary state, active sets, task-boundary analysis, distributed in-memory graphs, and state synchronization.

More precisely, NebulaGraph Analytics shares the underlying database substrate with regular graph queries, but adds a set of extensions specific to graph analytics. Saying that they use “exactly the same engine” would obscure the real differences. A regular graph query and an iterative PageRank procedure have different state lifetimes, communication patterns, and consistency requirements. The right things to reuse are the stable abstractions—syntax, types, expressions, plans, vectors, and pipelines—while the analytics runtime remains specialized where the workload requires it.

The resulting architecture has two layers.

  • At the upper layer, a procedure interpreter orchestrates variable declarations, control flow, graph-computationphases, and convergence checks. This layer handles relatively few operations, but each operation is semantically substantial.

  • At the lower layer, each graph-computing phase is translated into worker-side pipelines: scan the graph or active set, expand topology, evaluate filters and expressions in batches, and apply state updates. This is where most of the data processing happens.

Figure 2. Control flow coordinates coarse algorithm phases; data-intensive graph work remains in vectorized worker pipelines.

Why the Runtime Is Hybrid Rather Than Fully Compiled

Earlier design work considered translating the graph DSL into C++ and compiling each algorithm. This is a viable path: Green-Marl, GraphIt, and other graph DSLs demonstrate the value of compilation, while structured query engines demonstrate the value of specializing hot paths.

It is not, however, the only path from a high-level language to efficient execution.

Procedure-level control flow is typically relatively small and dynamic. A convergence loop may depend on a global value produced by the previous distributed phase. Its branches are important at the algorithm level, but they are usually not the dominant source of per-edge execution cost.

The picture is different for topology scanning, expansion, filtering, expression evaluation, and state accumulation. These operations run over large numbers of nodes and edges and can benefit from batching, vectorized expressions, pipelining, and multicore parallelism.

The current implementation therefore uses a hybrid boundary: the outer procedure is interpreted, while the hot graph-computation paths run through vectorized pipelines. Local control flow inside a graph-compute block can still operate over batches by using masks to select the rows participating in a branch or loop, without degrading the entire task into row-at-a-time scalar execution.

This choice avoids coupling every procedure to generated native code artifacts and to a specific version of the distributed runtime. Improvements to plans or operators can also benefit existing procedures directly, without recompiling the procedure source into a new algorithm library.

This choice also comes with tradeoffs. Interpretation is not free. Highly divergent row-level logic can reduce vectorizationefficiency. Sparse frontiers and dense phases may require different scan strategies and state representations. Code generation may still be valuable in the future for particular expressions or stable hot paths. We do not claim that vectorization is universally superior; rather, the hybrid model concentrates optimization effort on the dominant computational workload while preserving a flexible language layer.

Communication Follows State Semantics

This design is most directly reflected in distributed execution. The procedure author does not manually send a node value to another machine. The engine derives communication from three facts:

  1. which node state a task reads;

  2. which node state it contributes to;

  3. when those contributions must become visible.

Consider the PageRank propagation again:


The task reads rank and out_degree from s, then contributes to next_rank on t. The destination update is a SumAgg, so independently produced contributions have a defined merge operation. The surrounding PER PATH phase establishes that updates accumulate while paths are processed in parallel and become externally visible only after the phase completes.

In a distributed in-memory graph, the same logical node may have copies on multiple Workers because its incident edges can be distributed across machines. One copy acts as the master and holds the authoritative state, while the others act as mirrors for local graph computation.

For a path-oriented task, the state flow is conceptually straightforward:

  1. synchronize the required readable state from masters to mirrors;

  2. execute local scans, expansions, and expression evaluation;

  3. merge contributions from mirrors back into the corresponding masters.

Global node identity and partition ownership determine where state belongs. The user procedure does not need to specifya host or construct a message.

Global aggregators follow a different path because their logical owner is the procedure, rather than a graph node. Workers produce local partial results, the Driver combines them, and the fully merged value is made available to the next phase that reads it. PageRank's convergence value, delta, is a clear example.

Figure 3. Node-attached state follows graph ownership; global state is reduced at the procedure level.

The master/mirror mechanism is important, but it is only one part of the design. Other partitioning or communication strategies could still be introduced in the future. What we want to preserve is the abstraction boundary: the language defines state ownership and visibility at a logical level, while the runtime decides how to realize them physically.

This is the same reason databases distinguish a logical join from a hash join or merge join. The user expresses the logical relationship and its semantics, while the execution engine retains the freedom to choose or change the physical strategy.

Optimization Space Comes from Preserved Meaning

The database analogy becomes concrete when we ask what the engine knows before execution.

The engine can know which graph pattern defines the task, which node the computation starts from, whether an active set restricts the task, which node-attached states are read, which states are accumulated and which are overwritten, which global values are produced, and where results become visible.

That information gives the execution engine several ways to optimize:

  • start a node task only from authoritative node copies;

  • restrict graph expansion using an active frontier;

  • avoid synchronizing node state that the next phase does not read;

  • combine updates locally before sending them across Workers;

  • choose batch size, concurrency, and temporary-state representations independently of the procedure;

  • change traversal or communication strategies without changing the algorithm source.

The language abstraction does not make all of these optimizations appear at once. What it provides is an optimization space that the runtime and optimizer can exploit more fully over time.

Type specialization is another important optimization in the execution engine. Scalar reductions and container aggregators can use their own specialized representations, update paths, and merge implementations, avoiding the need to force workloads with very different characteristics through the same generic execution path.

The procedure is still not purely declarative. Statement order matters. A loop contains explicit state transitions. Export and logging are observable effects. Node-state assignment cannot be freely reordered relative to aggregation.

We therefore cannot claim that a relational cost-based optimizer can automatically discover the optimal distributed graph algorithm. A more accurate statement is that, compared with opaque callbacks, structured graph-computation semantics give the execution engine a much larger optimization space.

Design Boundaries and Costs

Every abstraction hides some details by moving responsibility into the system, and those responsibilities come with costs.

Task visibility boundaries give parallel updates clear semantics, but they also introduce synchronization. A slow Worker can delay the next phase. Frequently read node state may need to be synchronized repeatedly. High-degree nodes can still accumulate large volumes of contributions even when their edges are distributed across multiple machines.

Aggregator types make distributed updates composable, but they are not free. Numeric reductions are compact, while large list, set, map, or top-k states can create significant memory and communication pressure. A program that is valid at the language level is not necessarily efficient.

Vectorized execution amortizes expression-evaluation and dispatch costs across a batch, but its effectiveness decreases when the active set is small or when different rows follow highly divergent control-flow paths. Adaptive sparse-versus-dense execution and representation selection therefore remain important directions.

Temporary algorithm state is intentionally kept separate from persistent graph properties. This simplifies iteration and allows specialized layouts for analytical computation, but it also means that graph computation does not directly update persistent graph properties. Persisting results is a separate phase with its own I/O and consistency concerns.

Finally, batched task boundaries are better suited to deterministic, phase-oriented reasoning. Some algorithms may benefit from asynchronous updates, but asynchronous execution changes state visibility and convergence semantics. If the language contract already defines phase boundaries, the runtime cannot treat asynchronous execution as a low-level optimization that preserves program behavior; it requires an explicit semantic contract of its own.

These constraints arise from the design choices themselves, and they also define the boundaries for future work.

What “the Database Way” Means

The complete design can be summarized as a chain:

Figure 4. The database-shaped path from algorithm intent to distributed execution and state movement.

The user describes PageRank in terms of iteration, topology, node state, and convergence conditions. The language rules prevent updates whose semantics would be ambiguous under parallel execution. The planner turns graph-computation phases into tasks while preserving state dependencies. Worker pipelines process graph data in batches. The distributed runtime synchronizes node-attached state according to the visibility guarantees defined by the procedure and merges global values.

The system we built connects language-level semantics, the database execution framework, and distributed analytics state. The key is to define clear responsibility boundaries between these layers while preserving state-ownership and visibility semantics across those boundaries.

Doing graph computing “the database way” does not mean putting SQL syntax in front of a graph computing framework. It means applying practices long established in database engineering: defining explicit semantics, validating programs before execution, separating logical plans from physical strategies, processing data in vectorized form, and letting the engine manage communication.

Physical execution strategies can continue to evolve. We can improve frontier handling, choose state representations adaptively, reduce communication, specialize hot tasks, or introduce selective code generation when measurements show that it is worthwhile. A procedure should not need to change merely because its physical implementation changes.

The long-term value of the language is not to freeze graph computing into a particular model, but to provide a stable semantic boundary beneath which the system can continue to evolve.

Conclusion

Building graph computing as a database capability is ultimately a question of where semantics and responsibility should live.

In NebulaGraph Analytics, the procedure describes the algorithm's control flow, graph patterns, state transitions, aggregation semantics, and visibility boundaries. The database infrastructure provides the machinery for parsing, validation, planning, expression evaluation, and vectorized execution. The distributed runtime then takes responsibility for state synchronization, aggregation, and physical data movement.

This separation is what makes the system extensible. A graph algorithm does not need to encode where data is stored, which worker should communicate with another worker, or how a node-attached update should be merged. At the same time, the engine is not forced to treat the algorithm as an opaque callback. The semantics exposed by the procedure provide an optimization space that the runtime can exploit more fully over time.

There are real trade-offs. Synchronization has a cost, aggregators consume memory and communication bandwidth, vectorized execution is not equally effective for every workload, and synchronous phase boundaries constrain which execution strategies are semantics-preserving. These are not incidental implementation details; they are part of the design space of a programmable distributed graph-computing system.

The goal is not to claim that one execution model is universally optimal. It is to establish a stable boundary between what the algorithm means and how the system executes it. With that boundary in place, NebulaGraph Analytics can continue to improve its execution strategies—from frontier processing and state representation to communication and selective specialization—without requiring algorithm authors to rewrite their procedures.

That is what we mean by building graph computing the database way.


The Developers Behind NebulaGraph Analytics

NebulaGraph Analytics is the result of the collective efforts of the developers who contributed to its design and implementation. We thank everyone who helped bring the system from concept to production.

Developers: Bin Huang, Jie.Wang, Kyle.Cao, Shihao Wu, Yee, Yichen Wang, Yuxuan.Wang


Selected Primary References

Go From Zero to Graph in Minutes

Spin Up Your NebulaGraph Cluster Instantly! 

✅ 14-day free trial

✅ No credit card required

✅ Cancel anytime

Go From Zero to Graph in Minutes

Spin Up Your NebulaGraph Cluster Instantly! 

✅ 14-day free trial ✅ No credit card required ✅ Cancel anytime

Go From Zero to Graph in Minutes

Spin Up Your NebulaGraph Cluster Instantly! 

✅ 14-day free trial
✅ No credit card required
✅ Cancel anytime