Skip to content

This is the multi-page printable view of this section. .

Return to the regular view of this page.

HugeGraph Computing (OLAP)

The HugeGraph-Computer repository contains two graph computing systems with different deployment and runtime models. Start with Go Vermeer for general graph algorithms; use Computer when you need distributed Java BSP/Pregel computation. Both can connect to HugeGraph, but their configuration and job entry points are not interchangeable.

Default entry: Go Vermeer

flowchart TB
    Master["Master"] --> Workers["Workers"]
    Master -->|PD| PD["PD"]
    Workers -->|Scan| Store["Store"]
    Workers -.->|REST| Server["Server"]

With load.type=hugegraph, the Vermeer master queries PD for partition metadata and workers scan HStore Store partitions directly. Workers write results through the Server REST API only when output.type=hugegraph; other input sources or output settings do not require those connections.

Java Computer

flowchart TB
    Config["Job configuration"] --> Master["Master"] --> Workers["Workers"]
    Master -. BSP .-> Etcd["etcd"]
    Workers <-->|REST| Server["Server"]
    Workers <-->|HDFS| HDFS["HDFS (optional)"]

Job configuration can be submitted through the Kubernetes Operator or YARN. Workers read graph data through the HugeGraph Server REST API and can write results back according to the output configuration. With HDFS input or output, workers access HDFS directly. The master uses etcd to coordinate BSP jobs.

1 - HugeGraph-Vermeer Quick Start

Vermeer high-performance in-memory graph computing: start once, execute repeatedly, with 15+ OLAP algorithms, seconds-to-minutes execution, deployment, loading, PageRank, and community detection.

1. Overview of Vermeer

1.1 Architecture

Vermeer is a high-performance, memory-first graph computing framework written in Go: start once and execute repeatedly. It supports fast execution of 15+ OLAP algorithms, often in seconds to minutes; actual time depends on graph size, algorithm parameters, and resources. A single master currently schedules multiple workers.

The master handles communication, forwarding, and aggregation, with modest computation and resource usage. Workers store graph data and execute tasks, consuming most memory and CPU. gRPC handles internal communication; REST provides external APIs.

At startup, built-in defaults are loaded first, followed by the [default] section of config/<env>.ini under the working directory, then explicit command-line overrides. For example, --env=master reads config/master.ini. The Docker image copies repository configuration to /go/bin/config/ and works from /go/bin/. A host bind mount hides the bundled files, so it must contain every ini file used. The configuration reader does not read environment variables.

Default ports:

RoleConfiguration keyDefault addressPurpose
masterhttp_peer0.0.0.0:6688REST API for host-side clients
mastergrpc_peer0.0.0.0:6689Workers connect to master
workerhttp_peer0.0.0.0:6788Worker HTTP service
workergrpc_peer0.0.0.0:6789Worker gRPC listener, also advertised to master for node communication

The Docker examples publish only master HTTP on host loopback 127.0.0.1:6688:6688; master and worker gRPC ports stay within the container network.

flowchart LR
  Client["curl / Python client"] -->|"HTTP :6688"| Master["master"]
  Worker["worker"] <-->|"Bidirectional gRPC: master :6689, worker :6789"| Master
  Master -->|"gRPC partition lookup"| PD["HugeGraph PD"]
  Worker -->|"gRPC partition scan"| Store["HugeGraph Store"]

1.2 Running Vermeer

Warning

In production, enable Server authentication and authorization, an IP allowlist, and minimum permissions; retain audit-*.log. Server Auth does not protect independent Vermeer, PD, or Store APIs. Restrict these HTTP/gRPC ports to trusted networks and callers, and configure access controls at the Vermeer external entry point.

master.ini defaults to auth=none, disabling authentication for ordinary and administrative APIs. Local examples publish loopback only. Before remote access, enable auth=token or restrict callers through a protected network or gateway.

Both Docker options need a host configuration directory. From the Vermeer repository root, copy the supplied master.ini and worker.ini templates. The mount hides image configuration under /go/bin/config, so never mount an empty directory or your entire home directory there:

CONFIG_DIR="$HOME/vermeer-config"
mkdir -p "$CONFIG_DIR"
cp config/master.ini config/worker.ini "$CONFIG_DIR/"

Keep master HTTP/gRPC listeners at 0.0.0.0:6688 and 0.0.0.0:6689. The Compose example assigns 172.20.0.10 to master and 172.20.0.11 to worker; set these options in the copied worker.ini:

[default]
http_peer=0.0.0.0:6788
grpc_peer=172.20.0.11:6789
master_peer=172.20.0.10:6689
run_mode=worker
worker_group=$

This single-worker example uses literal worker_group=$, the general group used when no named group is bound. The repository template defaults to worker_group=default; if retaining that named group, bind it to the task space or graph before submitting tasks. For default space $DEFAULT, use:

curl --fail --show-error -X POST \
  'http://localhost:6688/admin/workers/alloc/default/%24DEFAULT'

An errcode=0 response confirms binding; with token authentication, also supply the authorization header. master_peer must reach master gRPC from the worker network. grpc_peer is both the listener and advertised worker address, so do not advertise 0.0.0.0, which other containers cannot reach. Update these IPs when changing the subnet or addresses; each additional worker needs a unique, mutually reachable grpc_peer.

  1. Option 1: Docker Compose (Recommended)

Run from the Vermeer repository root. Modify existing docker-compose.yaml or use the following example. The repository file currently mounts all of ~/ as configuration and does not publish master HTTP; replace the mount and add the port mapping before starting:

services:
  vermeer-master:
    image: hugegraph/vermeer
    container_name: vermeer-master
    ports:
      - "127.0.0.1:6688:6688"
    volumes:
      - /home/user/vermeer-config:/go/bin/config:ro
    command: --env=master
    networks:
      vermeer_network:
        ipv4_address: 172.20.0.10 # Static master address

  vermeer-worker:
    image: hugegraph/vermeer
    container_name: vermeer-worker
    volumes:
      - /home/user/vermeer-config:/go/bin/config:ro
    command: --env=worker
    networks:
      vermeer_network:
        ipv4_address: 172.20.0.11 # Static worker address

networks:
  vermeer_network:
    driver: bridge
    ipam:
      config:
        - subnet: 172.20.0.0/24 # Change the subnet as needed

Replace /home/user/vermeer-config with your actual absolute configuration directory. Changing the subnet or static IPs of vermeer_network also requires updating grpc_peer and master_peer in worker.ini. Do not use the default grpc_peer=0.0.0.0:6789 as an advertised container address.

Build and start from the project directory:

# Build the image from the Vermeer root directory
docker build -t hugegraph/vermeer .

# Start from the Vermeer root directory
docker-compose up -d
# Or use the newer CLI:
# docker compose up -d

View logs / stop / remove:

docker-compose logs -f
docker-compose down
  1. Option 2: Separate docker run Commands (Manual Network and Static IPs)

Set CONFIG_DIR to the prepared configuration directory, with grpc_peer and master_peer matching the static container addresses. Replace the example CONFIG_DIR with your actual absolute path.

Build the image:

docker build -t hugegraph/vermeer .

Create a custom bridge network (once):

docker network create --driver bridge \
  --subnet 172.20.0.0/24 \
  vermeer_network

Run master (use your absolute CONFIG_DIR and adjust IPs as needed):

CONFIG_DIR=/home/user/vermeer-config

docker run -d \
  --name vermeer-master \
  --network vermeer_network --ip 172.20.0.10 \
  -p 127.0.0.1:6688:6688 \
  -v ${CONFIG_DIR}:/go/bin/config:ro \
  hugegraph/vermeer \
  --env=master

Run worker:

docker run -d \
  --name vermeer-worker \
  --network vermeer_network --ip 172.20.0.11 \
  -v ${CONFIG_DIR}:/go/bin/config:ro \
  hugegraph/vermeer \
  --env=worker

View logs / stop / remove:

docker logs -f vermeer-master
docker logs -f vermeer-worker

docker stop vermeer-master vermeer-worker
docker rm vermeer-master vermeer-worker

# Delete the custom network if needed
docker network rm vermeer_network
  1. Option 3: Build from Source

Build following the Vermeer README.

go build

Start from the Vermeer root directory with ./vermeer --env=master and ./vermeer --env=worker01. In worker01.ini, set grpc_peer to an address bindable locally and reachable by master and other workers; point master_peer to master gRPC.

After starting master, check its HTTP port from the host:

curl --fail --show-error http://localhost:6688/graphs

Expect HTTP 200 and JSON errcode=0.

2. Task Creation REST API

2.1 Introduction

Submit a load task, wait for loading to finish, then submit a compute task. A loaded graph can support repeated computations and is not deleted after completion. Asynchronous APIs return creation results and task information, including ID, before completion; synchronous APIs wait for success or failure, so client and proxy HTTP timeouts must be sufficiently long. Query task states: loaded for successful loading, complete for successful computation, error for failure, and canceled for cancellation; other states remain waiting or running. Graphs loading or in an error state cannot be computed. Deletion requires a deletable graph state and no current usage.

Available URLs:

  • Asynchronous: POST http://master_ip:port/tasks/create; read the ID from response task.id.
  • Synchronous: POST http://master_ip:port/tasks/create/sync; returns after the task ends.
  • Query task: GET http://master_ip:port/task/{task_id}. errcode=0 indicates a successful query; inspect task.state for task status. Set a client polling deadline for asynchronous jobs; reaching it stops client waiting without canceling the server task.

2.2 Loading Graph Data

The examples list common load parameters. Each loader reads its own keys; unused keys do not change behavior.

Vermeer provides three loading methods:

  1. Local files

load.vertex_files and load.edge_files map worker address host portions to file paths read by those workers. With containers, mount data into the worker and use container paths. In the Compose example, mount a host data directory at worker /data (such as - /host/data:/data:ro) and use 172.20.0.11 from grpc_peer as the mapping key. Other deployments use the host portion of their advertised worker addresses.

Obtain a dataset such as Twitter-2010; the first twitter-2010.txt.gz file is sufficient.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "load",
 "graph": "testdb",
 "params": {
  "load.parallel": "50",
  "load.type": "local",
  "load.vertex_files": "{\"172.20.0.11\":\"/data/twitter-2010.v_[0,99]\"}",
  "load.edge_files": "{\"172.20.0.11\":\"/data/twitter-2010.e_[0,99]\"}",
  "load.use_out_degree": "1",
  "load.use_outedge": "1"
 }
}
  1. HugeGraph

Request example:

Replace request addresses, graph name, and credentials with your actual connection settings.

POST http://localhost:6688/tasks/create
{
  "task_type": "load",
  "graph": "testdb",
  "params": {
    "load.parallel": "50",
    "load.type": "hugegraph",
    "load.hg_pd_peers": "[\"<pd-address-reachable-from-vermeer>:8686\"]",
    "load.hugegraph_name": "DEFAULT/hugegraph2/g",
    "load.hugegraph_username": "admin",
    "load.hugegraph_password": "<your-password-here>",
    "load.use_out_degree": "1",
    "load.use_outedge": "1"
  }
}

Vermeer master connects to PD through load.hg_pd_peers to look up partitions; workers read data from the returned Store addresses. These services must be reachable from the corresponding Vermeer hosts or containers. Inside Docker, 127.0.0.1 identifies the container itself, not the host or another container.

  1. HDFS

Request example:

POST http://localhost:6688/tasks/create
{
  "task_type": "load",
  "graph": "testdb",
  "params": {
    "load.parallel": "50",
    "load.type": "hdfs",
    "load.hdfs_namenode": "name_node1:9000",
    "load.hdfs_conf_path": "/path/to/conf",
    "load.krb_realm": "EXAMPLE.COM",
    "load.krb_name": "user@EXAMPLE.COM",
    "load.krb_keytab_path": "/path/to/keytab",
    "load.krb_conf_path": "/path/to/krb5.conf",
    "load.hdfs_use_krb": "1",
    "load.vertex_files": "/data/graph/vertices",
    "load.edge_files": "/data/graph/edges",
    "load.use_out_degree": "1",
    "load.use_outedge": "1"
  }
}

2.3 Outputting Computation Results

Current result writers support local, hdfs, and hugegraph through output.type; none disables output. Source retains an afs constant, but current master registers no AFS loader or writer. Set output.need_statistics=1 to put statistics into task information; supported statistics operators depend on each algorithm implementation.

The examples list common computation and output parameters; supported algorithm parameters depend on current Vermeer implementations.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "pagerank",
 "compute.parallel": "10",
 "compute.max_step": "10",
 "output.type": "local",
 "output.parallel": "1",
 "output.file_path": "result/pagerank"
  }
}

output.type=local writes to the executing worker’s local filesystem. For containers, mount the output directory if the host needs to read results.

3. Supported Algorithms

3.1 PageRank

The PageRank algorithm, also known as the web ranking algorithm, is a technique used by search engines to calculate the relevance and importance of web pages (nodes) based on their mutual hyperlinks.

  • If a web page is linked to by many other web pages, it indicates that the web page is relatively important, and its PageRank value will be relatively high.
  • If a web page with a high PageRank value links to other web pages, the PageRank value of the linked web pages will also increase accordingly.

The PageRank algorithm is suitable for scenarios such as web page ranking and identifying key figures in social networks.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "pagerank",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/pagerank",
 "compute.max_step":"10"
 }
}

3.2 WCC (Weakly Connected Components)

The weakly connected components algorithm calculates all connected subgraphs in an undirected graph and outputs the weakly connected subgraph ID to which each vertex belongs, indicating the connectivity between points and distinguishing different connected communities.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "wcc",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/wcc",
 "compute.max_step":"10"
 }
}

3.3 LPA (Label Propagation Algorithm)

The label propagation algorithm is a graph clustering algorithm commonly used in social networks to discover potential communities.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "lpa",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/lpa",
 "compute.max_step":"10"
 }
}

3.4 Degree Centrality

The degree centrality algorithm calculates the degree centrality value of each node in the graph, supporting both undirected and directed graphs. Degree centrality is an important indicator of node importance; the more edges a node has with other nodes, the higher its degree centrality value, and the more important the node is in the graph. In an undirected graph, degree centrality is calculated based on edge information to count the number of times a node appears, resulting in the degree centrality value of the node. In a directed graph, it is based on the direction of the edges, filtering based on input or output-edge information to count the number of times a node appears, resulting in the in-degree or out-degree value of the node. It indicates the importance of each point, with more important points having higher degrees.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "degree",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/degree",
 "degree.direction":"both"
 }
}

3.5 Closeness Centrality

Closeness centrality is used to calculate the inverse of the shortest distance from a node to all other reachable nodes, accumulating and normalizing the value. Closeness centrality can be used to measure the time it takes for information to be transmitted from the node to other nodes. The larger the closeness centrality of a node, the closer its position in the graph is to the center, suitable for scenarios such as identifying key nodes in social networks.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "closeness_centrality",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/closeness_centrality",
 "closeness_centrality.sample_rate":"0.01"
 }
}

3.6 Betweenness Centrality

The betweenness centrality algorithm determines the value of a node as a “bridge” node; the larger the value, the more likely it is to be a necessary path between two points in the graph. Typical examples include mutual followers in social networks. It is suitable for measuring the degree of aggregation around a node in a community.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "betweenness_centrality",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/betweenness_centrality",
 "betweenness_centrality.sample_rate":"0.01"
 }
}

3.7 Triangle Count

The triangle count algorithm calculates the number of triangles passing through each vertex, suitable for calculating the relationships between users and whether the associations form triangles. The more triangles, the higher the degree of association between nodes in the graph, and the tighter the organizational relationship. In social networks, triangles indicate cohesive communities, and identifying triangles helps understand clustering and interconnections among individuals or groups in the network. In financial or transaction networks, the presence of triangles may indicate suspicious or fraudulent activities, and triangle counting can help identify transaction patterns that may require further investigation.

The output result is the Triangle Count corresponding to each vertex, i.e., the number of triangles the vertex is part of.

Note: This algorithm is for undirected graphs and ignores edge directions.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "triangle_count",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/triangle_count"
 }
}

3.8 K-Core

The K-Core algorithm marks all vertices with a degree of K, suitable for graph pruning and finding the core part of the graph.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "kcore",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/kcore",
 "kcore.degree_k":"5"
 }
}

3.9 SSSP (Single Source Shortest Path)

The single source the shortest path algorithm calculates the shortest distance from one point to all other points.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "sssp",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/degree",
 "sssp.source":"tom"
 }
}

3.10 KOUT

Starting from a point, get the k-layer nodes of this point.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "kout",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/kout",
 "kout.source":"tom",
 "compute.max_step":"6"
 }
}

3.11 Louvain

The Louvain algorithm is a community detection algorithm based on modularity. The basic idea is that nodes in the network try to traverse all neighbor community labels and choose the community label that maximizes the modularity increment. After maximizing modularity, each community is regarded as a new node, and the process is repeated until the modularity no longer increases.

The distributed Louvain algorithm implemented on Vermeer is affected by factors such as node order and parallel computation. Due to the random traversal order of the Louvain algorithm, community compression also has a certain randomness, leading to different results in multiple executions. However, the overall trend will not change significantly.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "louvain",
 "compute.parallel":"10",
 "compute.max_step":"1000",
 "louvain.threshold":"0.0000001",
 "louvain.resolution":"1.0",
 "louvain.step":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/louvain"
  }
 }

3.12 Jaccard Similarity Coefficient

The Jaccard index, also known as the Jaccard similarity coefficient, is used to compare the similarity and diversity between finite sample sets. The larger the Jaccard coefficient value, the higher the similarity of the samples. It is used to calculate the Jaccard similarity coefficient between a given source point and all other points in the graph.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "jaccard",
 "compute.parallel":"10",
 "compute.max_step":"2",
 "jaccard.source":"123",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/jaccard"
 }
}

3.13 Personalized PageRank

The goal of personalized PageRank is to calculate the relevance of all nodes relative to user u. Starting from the node corresponding to user u, at each node, there is a probability of 1-d to stop walking and start again from u, or a probability of d to continue walking, randomly selecting a node from the nodes pointed to by the current node to walk down. It is used to calculate the personalized PageRank score starting from a given starting point, suitable for scenarios such as social recommendations.

Since the calculation requires using out-degree, load.use_out_degree needs to be set to 1 when reading the graph.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "ppr",
 "compute.parallel":"100",
 "compute.max_step":"10",
 "ppr.source":"123",
 "ppr.damping":"0.85",
 "ppr.diff_threshold":"0.00001",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/ppr"
 }
}

3.14 Global Kout

Calculate the k-degree neighbors of all nodes in the graph (excluding themselves and 1~k-1 degree neighbors). Due to the severe memory expansion of the global kout algorithm, k is currently limited to 1 and 2. Additionally, the global kout algorithm supports filtering functions (parameters such as “compute.filter”:“risk_level==1”), and the filtering condition is judged when calculating the k-degree. The final result set includes those that meet the filtering condition. The algorithm’s final output is the number of neighbors that meet the condition.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "kout_all",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"10",
 "output.file_path":"result/kout",
 "compute.max_step":"2",
 "compute.filter":"risk_level==1"
 }
}

3.15 Clustering Coefficient

The clustering coefficient represents the coefficient of the clustering degree of nodes in a graph. In real networks, especially in specific networks, nodes tend to establish a tightly organized relationship due to relatively high-density connection points. The clustering coefficient algorithm (Cluster Coefficient) is used to calculate the clustering degree of nodes in the graph. This algorithm is for local clustering coefficients. The local clustering coefficient can measure the clustering degree around each node in the graph.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "clustering_coefficient",
 "compute.parallel":"100",
 "compute.max_step":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/cc"
 }
}

3.16 SCC (Strongly Connected Components)

In the mathematical theory of directed graphs, if every vertex of a graph can be reached from any other point in the graph, the graph is said to be strongly connected. The parts of any directed graph that can achieve strong connectivity are called strongly connected components. It indicates the connectivity between points and distinguishes different connected communities.

Request example:

POST http://localhost:6688/tasks/create
{
 "task_type": "compute",
 "graph": "testdb",
 "params": {
 "compute.algorithm": "scc",
 "compute.parallel":"10",
 "output.type":"local",
 "output.parallel":"1",
 "output.file_path":"result/scc",
 "compute.max_step":"200"
 }
}

🚧, further updates and improvements will be made at any time. Suggestions and feedback are welcome.

2 - HugeGraph-Computer Quick Start

1. Component Overview

HugeGraph-Computer is a Java distributed graph computing framework based on BSP (Bulk Synchronous Parallel), with algorithms running in iterative supersteps. Kubernetes Operator or YARN can schedule jobs; master and worker processes can also run on one machine for small trial jobs.

Computer and Vermeer, in the same repository, are separate implementations: Computer uses a Java/BSP runtime for distributed computing, while Vermeer is a Go in-memory graph computing platform with a master-worker architecture. They share HugeGraph data sources, but their deployment and job configurations are not interchangeable.

Computer reads graph data from HugeGraph or HDFS and writes results to either system. The runtime can spill some data to disk; whether a job completes still depends on input size, resources, and configuration. Disk spilling does not remove resource limits.

2. Prerequisites and Connections

Building and running require JDK 11 or later; source builds also require Maven 3.5 or later. The PageRank example needs a running HugeGraph-Server with graph data, plus etcd reachable by the master and workers.

ServiceExample address or portPurpose
HugeGraph-Serverhttp://127.0.0.1:8080Read graph data and write algorithm results; configured by hugegraph.url.
etcdhttp://127.0.0.1:2379BSP job coordination; configured by bsp.etcd_endpoints.
Computer master RPCTCP 8190Workers connect to master; port in the distribution configuration.
Computer worker data transferOS-assigned locally; K8s Operator defaults to 8099Transfer vertices and messages between workers. Advertised addresses and ports must be reachable across hosts.
MinIO (K8s manifests)HTTP 9000Input partition snapshots. Computer 1.7.0 manifests incorrectly map Service port 9000 to MinIO Console port 9090; change targetPort to 9000 before enabling snapshots. Defaults are snapshot.write=false and snapshot.load=false, so MinIO access is unnecessary when snapshots are disabled.
cert-manager (K8s)In-cluster serviceOperator manifests use cert-manager.io/v1 Certificate, Issuer, and CA injection. Install a compatible version before deploying the Operator.
HDFSCluster-specificRequired only when input or output uses HDFS.

For Kubernetes jobs, hugegraph.url must be reachable from all computing Pods; do not use a localhost address available only on your computer. If HugeGraph authentication is enabled, configure the username and password and supply matching credentials for REST queries.

Warning

In production, enable Server authentication and authorization, retain Server audit logs (normally audit-*.log), give Computer a dedicated account with only the read/write permissions required by its jobs, and configure a source IP allowlist for Server. hugegraph.username and hugegraph.password are credentials for connecting to Server, not a replacement for Server authentication.

See the Computer configuration reference for more options.

3. Obtain Source and Build a Distribution

The Apache HugeGraph download directory provides versioned Computer source archives, without separate precompiled Computer binaries. This example uses VERSION=1.7.0, whose source archive and build output names include -incubating-. For other versions, use the actual filenames in the download directory.

VERSION=1.7.0 # Replace with the release to download
ARCHIVE="apache-hugegraph-computer-incubating-${VERSION}-src.tar.gz"
DOWNLOAD_BASE="https://downloads.apache.org/hugegraph/${VERSION}"
curl -fLO "${DOWNLOAD_BASE}/${ARCHIVE}"
curl -fLO "${DOWNLOAD_BASE}/${ARCHIVE}.sha512"
curl -fLO "${DOWNLOAD_BASE}/${ARCHIVE}.asc"
curl -fLO https://downloads.apache.org/hugegraph/KEYS
shasum -a 512 -c "${ARCHIVE}.sha512"
gpg --import KEYS
gpg --verify "${ARCHIVE}.asc" "${ARCHIVE}"
tar -xzf "${ARCHIVE}"
cd "apache-hugegraph-computer-incubating-${VERSION}-src/computer"
mvn clean package -DskipTests
tar -xzf "target/apache-hugegraph-computer-incubating-${VERSION}.tar.gz"
cd "apache-hugegraph-computer-incubating-${VERSION}"

To build current master, clone and package the source:

git clone https://github.com/apache/hugegraph-computer.git
cd hugegraph-computer/computer
mvn clean package -DskipTests

The computer/computer-dist module assembles the distribution, including bin/start-computer.sh, runtime dependencies in lib/, built-in algorithms in algorithm/builtin-algorithm.jar, and default conf/computer.properties and conf/log4j2.xml. The source configuration is computer/computer-dist/src/assembly/static/conf/computer.properties. Release 1.7.0 build output includes -incubating-; current master omits it and takes its version from the POM. Use archive names matching the source you built.

4. Run PageRank Locally

Edit conf/computer.properties in the distribution directory with your HugeGraph URL, graph name, and credentials; point bsp.etcd_endpoints to reachable etcd. The default configuration selects built-in PageRankParams. Master and workers must use the same configuration and job.id; concurrent jobs need distinct job IDs.

Start processes from the distribution directory in two terminals. The script reads conf/computer.properties by default; use -c to select another file.

# Terminal 1: start master
bin/start-computer.sh -d local -r master
# Terminal 2: start worker
bin/start-computer.sh -d local -r worker

Master waits for the configured workers to register before running the job. Check both terminal outputs; the default logging configuration also writes master and worker logs under the distribution directory’s logs/. Process startup alone does not confirm completion: verify that input, superstep computation, and output all finish normally in the master log.

PageRank writes results to HugeGraph under page_rank. First ensure every target vertex label allows this property. Computer creates a DOUBLE, OLAP_COMMON property key but does not add it to vertex labels. If absent, create it:

curl --fail --request POST \
  --header 'Content-Type: application/json' \
  --data '{"name":"page_rank","data_type":"DOUBLE","cardinality":"SINGLE","write_type":"OLAP_COMMON"}' \
  'http://127.0.0.1:8080/graphspaces/DEFAULT/graphs/hugegraph/schema/propertykeys'

If the property key already exists, verify its type is DOUBLE and write_type is OLAP_COMMON. List vertex labels, then add nullable page_rank to every label being computed; replace person with the actual label:

curl --fail --compressed \
  'http://127.0.0.1:8080/graphspaces/DEFAULT/graphs/hugegraph/schema/vertexlabels'

curl --fail --request PUT \
  --header 'Content-Type: application/json' \
  --data '{"name":"person","properties":["page_rank"],"nullable_keys":["page_rank"]}' \
  'http://127.0.0.1:8080/graphspaces/DEFAULT/graphs/hugegraph/schema/vertexlabels/person?action=append'

If the current read mode hides OLAP writes, an administrator can set it to ALL:

curl --fail --request PUT \
  --header 'Content-Type: application/json' \
  --data '"ALL"' \
  'http://127.0.0.1:8080/graphspaces/DEFAULT/graphs/hugegraph/graph_read_mode'

Query vertices to verify the result property:

curl --fail --compressed \
  'http://127.0.0.1:8080/graphspaces/DEFAULT/graphs/hugegraph/graph/vertices?limit=3'

When authentication is required, add --user "$HG_USER:$HG_PASSWORD" to curl. See the graph read-mode REST API for the endpoint and permissions.

5. Run PageRank on Kubernetes

Ensure computing Pods can reach HugeGraph-Server, then install a cluster-compatible cert-manager using its official installation guide. The Computer Operator manifests deploy the Operator, etcd, and MinIO. Use the same Computer release for the CRD and Operator manifests; this example uses 1.7.0:

VERSION=1.7.0 # Use the same release for CRD and Operator manifests
kubectl apply -f "https://raw.githubusercontent.com/apache/hugegraph-computer/${VERSION}/computer/computer-k8s-operator/manifest/hugegraph-computer-crd.v1.yaml"
kubectl apply -f "https://raw.githubusercontent.com/apache/hugegraph-computer/${VERSION}/computer/computer-k8s-operator/manifest/hugegraph-computer-operator.yaml"
kubectl rollout status deployment/hugegraph-computer-operator-controller-manager \
  -n hugegraph-computer-operator-system --timeout=120s
Optional: enable MinIO snapshots (1.7.0)

The 1.7.0 MinIO Service maps S3 API port 9000 to Console port 9090. Fix this mapping before enabling snapshots and set snapshot.minio_endpoint to http://hugegraph-computer-operator-minio.hugegraph-computer-operator-system.svc:9000. Skip this when snapshots remain disabled.

kubectl patch service hugegraph-computer-operator-minio \
  -n hugegraph-computer-operator-system \
  --type=json \
  -p='[{"op":"replace","path":"/spec/ports/0/targetPort","value":9000}]'

Replace the HugeGraph URL with a service address reachable by computing Pods, then submit a HugeGraphComputerJob. This example uses the official Docker Hub image hugegraph/hugegraph-computer:latest and its built-in PageRank JAR. Custom algorithm JARs can be bundled in the image and selected with jarFile, or downloaded from HTTP(S) with remoteJarUri. The partition count must be at least the worker count.

Save the following YAML as pagerank-job.yaml, then apply it:

apiVersion: operator.hugegraph.apache.org/v1
kind: HugeGraphComputerJob
metadata:
  namespace: hugegraph-computer-operator-system
  name: pagerank-sample
spec:
  jobId: pagerank-sample
  algorithmName: page_rank
  image: hugegraph/hugegraph-computer:latest # Official Docker Hub image
  jarFile: /hugegraph/hugegraph-computer/algorithm/builtin-algorithm.jar
  pullPolicy: IfNotPresent
  workerInstances: 1
  masterCpu: 500m
  workerCpu: 500m
  masterMemory: 1Gi
  workerMemory: 2Gi
  computerConf:
    job.partitions_count: "1"
    algorithm.params_class: org.apache.hugegraph.computer.algorithm.centrality.pagerank.PageRankParams
    hugegraph.url: http://hugegraph-server:8080
    hugegraph.name: hugegraph
kubectl apply -f pagerank-job.yaml
kubectl get hcjob pagerank-sample -n hugegraph-computer-operator-system --watch

SUCCEEDED means the job completed. By default, the Operator deletes the CR and computing resources after completion. Expand the following instructions when you need to retain state or investigate a failed job.

Optional: retain job state, inspect logs, or clean up resources

The Operator defaults to AUTO_DESTROY_POD=true, so short jobs may lose their CR and Pods before inspection. Disable this before submitting the job, wait for the replacement Controller Pod to become ready, then submit. Change this setting in the Java Operator’s controller container:

kubectl set env deployment/hugegraph-computer-operator-controller-manager \
  -n hugegraph-computer-operator-system \
  --containers=controller AUTO_DESTROY_POD=false
kubectl rollout status deployment/hugegraph-computer-operator-controller-manager \
  -n hugegraph-computer-operator-system --timeout=120s

Inspect runtime information or failure logs:

kubectl logs --follow <master-pod-name> -n hugegraph-computer-operator-system
kubectl logs --follow <worker-pod-name> -n hugegraph-computer-operator-system

After checking results, delete the CR to clean up associated resources and restore the default policy:

kubectl delete hcjob pagerank-sample -n hugegraph-computer-operator-system
kubectl set env deployment/hugegraph-computer-operator-controller-manager \
  -n hugegraph-computer-operator-system \
  --containers=controller AUTO_DESTROY_POD=true

After PageRank writes to HugeGraph, set the read mode and query as in the previous section. For HDFS output, results are under output.hdfs_path_prefix/<job.id>/; filenames and partition layout depend on job configuration.

See the CRD configuration reference for all fields.

6. Built-in Algorithms and Development

Current built-in algorithms include:

  • Centrality: PageRank, Betweenness Centrality, Closeness Centrality, Degree Centrality.
  • Communities and structure: Clustering Coefficient, K-core, LPA, Triangle Count, WCC.
  • Paths and sampling: ring detection, filtered ring detection, single-source shortest path, Random Walk.

Implementations are in computer/computer-algorithm. Custom algorithms must follow the Computer API and be packaged as loadable JARs; see the Computer README for modules and development entry points.

3 - HugeGraph-Computer Configuration Reference

Computer Config Options

The defaults in the tables come from ComputerOptions.java in computer-api; explicit conf/computer.properties overrides are shown as “code default (distribution: actual value)”. Common options such as rpc.* come from HugeGraph Commons; distribution values are listed separately below.

Configuration Sources

  • Source template: computer/computer-dist/src/assembly/static/conf/computer.properties, with log4j2.xml in the same directory. Maven package copies these into distribution conf/, runtime dependencies into lib/, and built-in algorithm JARs into algorithm/. Templates are maintained in source; there is no separate configuration generator.
  • Standalone and YARN: startup scripts read distribution conf/computer.properties by default. Override it with bin/start-computer.sh -c <configuration-path>; master and workers need the same job parameters.
  • Kubernetes Operator: users supply spec.computerConf in the CRD; the Operator writes a ConfigMap mounted as computer.properties. It supplies job ID, worker count, Pod addresses, and etcd when unspecified. Unspecified or zero transfer/RPC ports become 8099/8190; transport.server_host and rpc.server_host become Pod IPs. Job startup scripts substitute only ${POD_IP}, ${HOSTNAME}, ${POD_NAME}, and ${POD_NAMESPACE}.
  • Kubernetes job images must contain the Computer runtime. Bundle algorithm JARs and select them with jarFile, or use an HTTP(S) remoteJarUri for startup download. See the Computer Quick Start and CRD table below.

Apache downloads provide versioned Computer source archives without separate precompiled binaries. Release 1.7.0 source and distribution names include -incubating-; run mvn clean package -DskipTests in the tagged computer/ project to build them. Current master omits that marker and takes its version from the POM. Source archives cannot be started as built distributions.

Warning

Empty HugeGraph credentials and example MinIO keys below are for local demonstrations. In production, enable Server authentication and authorization, retain Server audit-*.log, use a Computer account with minimum required Server permissions, and configure a Server source IP allowlist. Do not use example MinIO keys in production.


1. Basic Configuration

Core job settings for HugeGraph-Computer.

config optiondefault valuedescription
hugegraph.urlhttp://127.0.0.1:8080The HugeGraph server URL to load data and write results back.
hugegraph.namehugegraphThe graph name to load data and write results back.
hugegraph.username"" (empty)The username for HugeGraph authentication (leave empty if authentication is disabled).
hugegraph.password"" (empty)The password for HugeGraph authentication (leave empty if authentication is disabled).
job.idlocal_0001 (packaged: local_001)The job identifier on YARN cluster or K8s cluster.
job.namespace"" (empty)Optional etcd job-key prefix for namespace isolation; distinct from Kubernetes metadata.namespace and not populated automatically by the Operator.
job.workers_count1The number of workers for one graph algorithm job. In K8s, this option is set by the Operator.
job.partitions_count1The number of partitions for computing one graph algorithm job.
job.partitions_thread_nums4The number of threads for partition parallel compute.

2. Algorithm Configuration

Algorithm-specific configuration for computation logic.

config optiondefault valuedescription
algorithm.params_classComputerOptions.Null placeholder classRequired. The class used to pass algorithm parameters before the algorithm runs.
algorithm.result_classComputerOptions.Null placeholder classThe vertex value class used to store computation results.
algorithm.message_classComputerOptions.Null placeholder classThe message class passed while computing a vertex.

3. Input Configuration

Configuration for loading input data from HugeGraph or other sources.

3.1 Input Source

config optiondefault valuedescription
input.source_typehugegraph-serverThe source type to load input data, allowed values: [‘hugegraph-server’, ‘hugegraph-loader’]. The ‘hugegraph-loader’ means use hugegraph-loader to load data from HDFS or file. If using ‘hugegraph-loader’, please configure ‘input.loader_struct_path’ and ‘input.loader_schema_path’.
input.loader_struct_path"" (empty)The structure path for Loader input. It takes effect only when input.source_type=hugegraph-loader.
input.loader_schema_path"" (empty)The schema path for Loader input. It takes effect only when input.source_type=hugegraph-loader.

3.2 Input Splits

config optiondefault valuedescription
input.split_size1048576 (1 MB)The input split size in bytes.
input.split_max_splits10000000The maximum number of input splits.
input.split_page_size500The page size for streamed load input split data.
input.split_fetch_timeout300The timeout in seconds to fetch input splits.

3.3 Input Processing

config optiondefault valuedescription
input.filter_classorg.apache.hugegraph.computer.core.input.filter.DefaultInputFilterThe class to create input-filter object. Input-filter is used to filter vertex edges according to user needs.
input.edge_directionOUTThe direction of edges to load, allowed values: [OUT, IN, BOTH]. When the value is BOTH, edges in both OUT and IN directions will be loaded.
input.edge_freqMULTIPLEThe frequency of edges that can exist between a pair of vertices, allowed values: [SINGLE, SINGLE_PER_LABEL, MULTIPLE]. SINGLE means only one edge can exist between a pair of vertices (identified by sourceId + targetId); SINGLE_PER_LABEL means each edge label can have one edge between a pair of vertices (identified by sourceId + edgeLabel + targetId); MULTIPLE means many edges can exist between a pair of vertices (identified by sourceId + edgeLabel + sortValues + targetId).
input.max_edges_in_one_vertex200The maximum number of adjacent edges allowed to be attached to a vertex. The adjacent edges will be stored and transferred together as a batch unit.

3.4 Input Performance

config optiondefault valuedescription
input.send_thread_nums4The number of threads for parallel sending of vertices or edges.

4. Snapshot & Storage Configuration

HugeGraph-Computer supports snapshot functionality to save vertex/edge partitions to local storage or MinIO object storage, enabling checkpoint recovery or accelerating repeated computations.

4.1 Basic Snapshot Configuration

config optiondefault valuedescription
snapshot.writefalseWhether to write snapshots of input vertex/edge partitions.
snapshot.loadfalseWhether to load from snapshots of vertex/edge partitions.
snapshot.name"" (empty)User-defined snapshot name to distinguish different snapshots.

4.2 MinIO Integration (Optional)

MinIO can be used as a distributed object storage backend for snapshots in K8s deployments.

config optiondefault valuedescription
snapshot.minio_endpoint"" (empty)MinIO service endpoint (e.g., http://minio:9000). Required when using MinIO.
snapshot.minio_access_keyminioadminMinIO access key for authentication.
snapshot.minio_secret_keyminioadminMinIO secret key for authentication.
snapshot.minio_bucket_name"" (empty)MinIO bucket name for storing snapshot data.

Usage Scenarios:

  • Checkpoint Recovery: Resume from snapshots after job failures, avoiding data reloading
  • Repeated Computations: Load data from snapshots when running the same algorithm multiple times
  • A/B Testing: Save multiple snapshot versions of the same dataset to test different algorithm parameters

Example: Local Snapshot (in computer.properties):

snapshot.write=true
snapshot.name=pagerank-snapshot-20260201

Example: MinIO Snapshot (in K8s CRD computerConf):

computerConf:
  snapshot.write: "true"
  snapshot.name: "pagerank-snapshot-v1"
  snapshot.minio_endpoint: "http://minio:9000"
  snapshot.minio_access_key: "my-access-key"
  snapshot.minio_secret_key: "my-secret-key"
  snapshot.minio_bucket_name: "hugegraph-snapshots"

5. Worker & Master Configuration

Configuration for worker and master computation logic.

5.1 Master Configuration

config optiondefault valuedescription
master.computation_classorg.apache.hugegraph.computer.core.master.DefaultMasterComputationMaster-computation is computation that can determine whether to continue to the next superstep. It runs at the end of each superstep on the master.

5.2 Worker Computation

config optiondefault valuedescription
worker.computation_classorg.apache.hugegraph.computer.core.config.NullThe class to create worker-computation object. Worker-computation is used to compute each vertex in each superstep.
worker.combiner_classorg.apache.hugegraph.computer.core.config.NullCombiner can combine messages into one value for a vertex. For example, PageRank algorithm can combine messages of a vertex to a sum value.
worker.partitionerorg.apache.hugegraph.computer.core.graph.partition.HashPartitionerThe partitioner that decides which partition a vertex should be in, and which worker a partition should be in.

5.3 Worker Combiners

config optiondefault valuedescription
worker.vertex_properties_combiner_classorg.apache.hugegraph.computer.core.combiner.OverwritePropertiesCombinerThe combiner can combine several properties of the same vertex into one properties at input step.
worker.edge_properties_combiner_classorg.apache.hugegraph.computer.core.combiner.OverwritePropertiesCombinerThe combiner can combine several properties of the same edge into one properties at input step.

5.4 Worker Buffers

config optiondefault valuedescription
worker.received_buffers_bytes_limit104857600 (100 MB)The limit bytes of buffers of received data. The total size of all buffers can’t exceed this limit. If received buffers reach this limit, they will be merged into a file (spill to disk).
worker.write_buffer_capacity52428800 (50 MB)The initial size of write buffer that used to store vertex or message.
worker.write_buffer_threshold52428800 (50 MB)The threshold of write buffer. Exceeding it will trigger sorting. The write buffer is used to store vertex or message.

5.5 Worker Data & Timeouts

config optiondefault valuedescription
worker.data_dirs[jobs]The directories separated by ‘,’ that received vertices and messages can persist into.
worker.wait_sort_timeout600000 (10 minutes)The max timeout (in ms) for message-handler to wait for sort-thread to sort one batch of buffers.
worker.wait_finish_messages_timeout86400000 (24 hours)The max timeout (in ms) for message-handler to wait for finish-message of all workers.

6. I/O & Output Configuration

Configuration for output computation results.

6.1 Output Class & Result

config optiondefault valuedescription
output.output_classorg.apache.hugegraph.computer.core.output.LogOutputThe class to output the computation result of each vertex. Called after iteration computation.
output.result_namevalueThe value is assigned dynamically by #name() of instance created by WORKER_COMPUTATION_CLASS.
output.result_write_typeOLAP_COMMONThe result write-type to output to HugeGraph, allowed values: [OLAP_COMMON, OLAP_SECONDARY, OLAP_RANGE].

6.2 Output Behavior

config optiondefault valuedescription
output.with_adjacent_edgesfalseWhether to output the adjacent edges of the vertex.
output.with_vertex_propertiesfalseWhether to output the properties of the vertex.
output.with_edge_propertiesfalseWhether to output the properties of the edge.

6.3 Batch Output

config optiondefault valuedescription
output.batch_size500The batch size of output.
output.batch_threads1The number of threads used for batch output.
output.single_threads1The number of threads used for single output.

6.4 HDFS Output

config optiondefault valuedescription
output.hdfs_urlhdfs://127.0.0.1:9000The HDFS URL for output.
output.hdfs_userhadoopThe HDFS user for output.
output.hdfs_path_prefix/hugegraph-computer/resultsThe directory of HDFS output results.
output.hdfs_delimiter, (comma)The delimiter of HDFS output.
output.hdfs_merge_partitionstrueWhether to merge output files of multiple partitions.
output.hdfs_replication3The replication number of HDFS.
output.hdfs_core_site_path"" (empty)The HDFS core site path.
output.hdfs_site_path"" (empty)The HDFS site path.
output.hdfs_kerberos_enablefalseWhether Kerberos authentication is enabled for HDFS.
output.hdfs_kerberos_principal"" (empty)The HDFS principal for Kerberos authentication.
output.hdfs_kerberos_keytab"" (empty)The HDFS keytab file for Kerberos authentication.
output.hdfs_krb5_conf/etc/krb5.confKerberos configuration file path.

6.5 Retry & Timeout

config optiondefault valuedescription
output.retry_times3The retry times when output fails.
output.retry_interval10The retry interval (in seconds) when output fails.
output.thread_pool_shutdown_timeout60The timeout (in seconds) of output thread pool shutdown.

7. Network & Transport Configuration

Configuration for network communication between workers and master.

7.1 Server Configuration

config optiondefault valuedescription
transport.server_host127.0.0.1Worker data listener and advertised address. It must be reachable by other workers across hosts; the Operator sets it to the Pod IP.
transport.server_port0 (distribution: 0; K8s Operator: 8099)Worker data port. Locally, 0 lets the OS assign a port; the Operator uses fixed port 8099 by default for Pod communication.
transport.server_threads4The number of transport threads for server.

7.2 Client Configuration

config optiondefault valuedescription
transport.client_threads4The number of transport threads for client.
transport.client_connect_timeout3000The timeout (in ms) of client connect to server.

7.3 Protocol Configuration

config optiondefault valuedescription
transport.provider_classorg.apache.hugegraph.computer.core.network.netty.NettyTransportProviderThe transport provider, currently only supports Netty.
transport.io_modeAUTOThe network IO mode, allowed values: [NIO, EPOLL, AUTO]. AUTO means selecting the appropriate mode automatically.
transport.tcp_keep_alivetrueWhether to enable TCP keep-alive.
transport.transport_epoll_ltfalseWhether to enable EPOLL level-trigger (only effective when io_mode=EPOLL).

7.4 Buffer Configuration

config optiondefault valuedescription
transport.send_buffer_size0The size of socket send-buffer in bytes. 0 means using system default value.
transport.receive_buffer_size0The size of socket receive-buffer in bytes. 0 means using system default value.
transport.write_buffer_high_mark67108864 (64 MB)The high water mark for write buffer in bytes. It will trigger sending unavailable if the number of queued bytes > write_buffer_high_mark.
transport.write_buffer_low_mark33554432 (32 MB)The low water mark for write buffer in bytes. It will trigger sending available if the number of queued bytes < write_buffer_low_mark.

7.5 Flow Control

config optiondefault valuedescription
transport.max_pending_requests8The max number of client unreceived ACKs. It will trigger sending unavailable if the number of unreceived ACKs >= max_pending_requests.
transport.min_pending_requests6The minimum number of client unreceived ACKs. It will trigger sending available if the number of unreceived ACKs < min_pending_requests.
transport.min_ack_interval200The minimum interval (in ms) of server reply ACK.

7.6 Timeouts

config optiondefault valuedescription
transport.close_timeout10000The timeout (in ms) of close server or close client.
transport.sync_request_timeout10000The timeout (in ms) to wait for response after sending sync-request.
transport.finish_session_timeout0The timeout (in ms) to finish session. 0 means using (transport.sync_request_timeout × transport.max_pending_requests).
transport.write_socket_timeout3000The timeout (in ms) to write data to socket buffer.
transport.server_idle_timeout360000 (6 minutes)The max timeout (in ms) of server idle.

7.7 Heartbeat

config optiondefault valuedescription
transport.heartbeat_interval20000 (20 seconds)The minimum interval (in ms) between heartbeats on client side.
transport.max_timeout_heartbeat_count120The maximum times of timeout heartbeat on client side. If the number of timeouts waiting for heartbeat response continuously > max_timeout_heartbeat_count, the channel will be closed from client side.

7.8 Advanced Network Settings

config optiondefault valuedescription
transport.max_syn_backlog511The capacity of SYN queue on server side. 0 means using system default value.
transport.recv_file_modetrueWhether to enable receive buffer-file mode. It will receive buffer and write to file from socket using zero-copy if enabled. Note: Requires OS support for zero-copy (e.g., Linux sendfile/splice).
transport.network_retries3The number of retry attempts for network communication if network is unstable.

7.9 Master RPC Configuration

These options come from HugeGraph Commons RPC configuration; the distribution template explicitly sets host and port. Kubernetes Operator advertises the Pod IP and defaults an unspecified or zero port to 8190.

OptionDistribution valueDescription
rpc.server_host127.0.0.1 (K8s: Pod IP)Master RPC address, reachable by workers.
rpc.server_port8190 (K8s default: 8190)Master RPC listener port; allow worker connections in security groups and network policies.

8. Storage & Persistence Configuration

Configuration for HGKV (HugeGraph Key-Value) storage engine and value files.

8.1 HGKV Configuration

config optiondefault valuedescription
hgkv.max_file_size2147483648 (2 GB)The max number of bytes in each HGKV file.
hgkv.max_data_block_size65536 (64 KB)The max byte size of HGKV file data block.
hgkv.max_merge_files10The max number of files to merge at one time.
hgkv.temp_file_dir/tmp/hgkvThis folder is used to store temporary files during the file merging process.

8.2 Value File Configuration

config optiondefault valuedescription
valuefile.max_segment_size1073741824 (1 GB)The max number of bytes in each segment of value-file.

9. BSP & Coordination Configuration

Configuration for Bulk Synchronous Parallel (BSP) protocol and etcd coordination.

config optiondefault valuedescription
bsp.etcd_endpointshttp://localhost:2379 (distribution: http://127.0.0.1:2379)Comma-separated etcd client endpoints. The Operator uses INTERNAL_ETCD_URL only when computerConf does not specify this key.
bsp.max_super_step10 (packaged: 2)The max super step of the algorithm.
bsp.register_timeout300000 (packaged: 100000)The max timeout (in ms) to wait for master and workers to register.
bsp.wait_workers_timeout86400000 (24 hours)The max timeout (in ms) to wait for workers BSP event.
bsp.wait_master_timeout86400000 (24 hours)The max timeout (in ms) to wait for master BSP event.
bsp.log_interval30000 (30 seconds)The log interval (in ms) to print the log while waiting for BSP event.

10. Performance Tuning Configuration

Configuration for performance optimization.

config optiondefault valuedescription
allocator.max_vertices_per_thread10000Maximum number of vertices per thread processed in each memory allocator.
sort.thread_nums4The number of threads performing internal sorting.

11. System Administration Configuration

Kubernetes Operator fills or overrides these options. Avoid overriding them unless customizing the Operator or network. job.namespace is optional user configuration; see the basic configuration table.

OptionManaged byDescription
bsp.etcd_endpointsK8s OperatorUses INTERNAL_ETCD_URL only when absent from computerConf.
transport.server_hostK8s OperatorOverrides with Pod IP; workers must reach each other.
transport.server_portK8s OperatorDefaults unspecified or zero values to 8099, not a random port.
job.idK8s OperatorSet from CRD job ID.
job.workers_countK8s OperatorSet from CRD workerInstances.
rpc.server_hostK8s OperatorOverrides with master Pod IP.
rpc.server_portK8s OperatorDefaults unspecified or zero values to 8190.
rpc.remote_urlComputer startupRemoved when configuration is read; do not set as job configuration.

Why these values must match:

  • BSP/RPC: Must match deployed etcd/RPC services for coordination.
  • Job settings: Must match the CRD to provide the correct worker count.
  • Transport: Workers need mutually reachable Pod IPs and ports; the K8s default is 8099.

K8s Operator Config Options

NOTE: Option needs to be converted through environment variable settings, e.g. k8s.internal_etcd_url => INTERNAL_ETCD_URL

config optiondefault valuedescription
k8s.auto_destroy_podtrueDelete the job CR after completion or failure; CR deletion triggers cleanup of associated computing resources.
k8s.close_reconciler_timeout120The max timeout (in ms) to close reconciler.
k8s.internal_etcd_urlCode: http://127.0.0.1:2379; manifest: http://hugegraph-computer-operator-etcd.hugegraph-computer-operator-system:2379etcd URL used by Operator jobs; the supplied manifest uses the etcd Service address.
k8s.internal_minio_urlCode: http://127.0.0.1:9000; manifest: http://hugegraph-computer-operator-minio.hugegraph-computer-operator-system:9000MinIO address supplied to jobs by the Operator; needed only for MinIO snapshots.
k8s.max_reconcile_retry3The max retry times of reconcile.
k8s.probe_backlog50The maximum backlog for serving health probes.
k8s.probe_port9892The port that the controller binds to for serving health probes.
k8s.ready_check_internal1000The time interval (ms) of check ready.
k8s.ready_timeout30000The max timeout (in ms) of check ready.
k8s.reconciler_countCode: Runtime.getRuntime().availableProcessors(); manifest: 6Maximum reconciler thread count.
k8s.resync_period600000The minimum frequency at which watched resources are reconciled.
k8s.timezoneAsia/ShanghaiThe timezone of computer job and operator.
k8s.watch_namespacehugegraph-computer-operator-systemNamespace watched for custom resources; also used by the supplied manifest. Use * to watch all namespaces.

HugeGraph-Computer CRD

1.7.0 CRD: https://github.com/apache/hugegraph-computer/blob/1.7.0/computer/computer-k8s-operator/manifest/hugegraph-computer-crd.v1.yaml

specdefault valuedescriptionrequired
algorithmNameThe name of algorithm.true
jobIdThe job id.true
imageJob container image, which must include the Computer runtime. Bundle algorithm JARs or download them with remoteJarUri.true
computerConfThe map of computer config options.true
workerInstancesThe number of worker instances, it will override the ‘job.workers_count’ option.true
pullPolicyKubernetes selects Always for latest, otherwise IfNotPresent, when unsetExplicit values: Always, Never, or IfNotPresent. The CRD has no default; see image pull policy.false
pullSecretsThe pull-secrets of Image, detail please refer to: https://kubernetes.io/docs/concepts/containers/images/#specifying-imagepullsecrets-on-a-podfalse
masterCpuThe cpu limit of master, the unit can be ’m’ or without unit detail please refer to: https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-cpufalse
workerCpuThe cpu limit of worker, the unit can be ’m’ or without unit detail please refer to: https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-cpufalse
masterMemoryThe memory limit of master, the unit can be one of Ei、Pi、Ti、Gi、Mi、Ki detail please refer to: https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-memoryfalse
workerMemoryThe memory limit of worker, the unit can be one of Ei、Pi、Ti、Gi、Mi、Ki detail please refer to: https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-memoryfalse
log4jXmlThe content of log4j.xml for computer job.false
jarFilePath to the algorithm JAR inside the image.false
remoteJarUriHTTP(S) algorithm JAR URL, downloaded and loaded by the startup script.false
jvmOptionsThe java startup parameters of computer job.false
envVarsplease refer to: https://kubernetes.io/docs/tasks/inject-data-application/define-interdependent-environment-variables/false
envFromplease refer to: https://kubernetes.io/docs/tasks/inject-data-application/define-environment-variable-container/false
masterCommandbin/start-computer.shThe run command of master, equivalent to ‘Entrypoint’ field of Docker.false
masterArgs["-r master", “-d k8s”]The run args of master, equivalent to ‘Cmd’ field of Docker.false
workerCommandbin/start-computer.shThe run command of worker, equivalent to ‘Entrypoint’ field of Docker.false
workerArgs["-r worker", “-d k8s”]The run args of worker, equivalent to ‘Cmd’ field of Docker.false
volumesPlease refer to: https://kubernetes.io/docs/concepts/storage/volumes/false
volumeMountsPlease refer to: https://kubernetes.io/docs/concepts/storage/volumes/false
secretPathsThe map of k8s-secret name and mount path.false
configMapPathsThe map of k8s-configmap name and mount path.false
podTemplateSpecPlease refer to: https://kubernetes.io/docs/reference/kubernetes-api/workload-resources/pod-template-v1/#PodTemplateSpecfalse
securityContextPlease refer to: https://kubernetes.io/docs/tasks/configure-pod-container/security-context/false

KubeDriver Config Options

config optiondefault valuedescription
k8s.build_image_bash_pathThe path of command used to build image.
k8s.enable_internal_algorithmtrueWhether enable internal algorithm.
k8s.framework_image_urlhugegraph/hugegraph-computer:latestThe image url of computer framework.
k8s.image_repository_passwordThe password for login image repository.
k8s.image_repository_registryThe address for login image repository.
k8s.image_repository_urlhugegraph/hugegraph-computerThe url of image repository.
k8s.image_repository_usernameThe username for login image repository.
k8s.internal_algorithm[pageRank]The name list of all internal algorithm. Note: Algorithm names use camelCase here (e.g., pageRank), but algorithm implementations return underscore_case (e.g., page_rank).
k8s.internal_algorithm_image_urlhugegraph/hugegraph-computer:latestThe image url of internal algorithm.
k8s.jar_file_dir/cache/jars/The directory where the algorithm jar will be uploaded.
k8s.kube_config~/.kube/configThe path of k8s config file.
k8s.log4j_xml_pathThe log4j.xml path for computer job.
k8s.namespacehugegraph-computer-operator-systemNamespace of the HugeGraph-Computer system.
k8s.pull_secret_names[]The names of pull-secret for pulling image.