This is the multi-page printable view of this section. .
HugeGraph Computing (OLAP)
- 1: HugeGraph-Vermeer Quick Start
- 2: HugeGraph-Computer Quick Start
- 3: HugeGraph-Computer Configuration Reference
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
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:
| Role | Configuration key | Default address | Purpose |
|---|---|---|---|
| master | http_peer | 0.0.0.0:6688 | REST API for host-side clients |
| master | grpc_peer | 0.0.0.0:6689 | Workers connect to master |
| worker | http_peer | 0.0.0.0:6788 | Worker HTTP service |
| worker | grpc_peer | 0.0.0.0:6789 | Worker 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
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:
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:
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:
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.
- 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:
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:
View logs / stop / remove:
- 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:
Create a custom bridge network (once):
Run master (use your absolute CONFIG_DIR and adjust IPs as needed):
Run worker:
View logs / stop / remove:
- Option 3: Build from Source
Build following the Vermeer README.
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:
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 responsetask.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=0indicates a successful query; inspecttask.statefor 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:
- 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:
- HugeGraph
Request example:
Replace request addresses, graph name, and credentials with your actual connection settings.
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.
- HDFS
Request example:
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:
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:
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:
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:
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:
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:
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:
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:
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:
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:
3.10 KOUT
Starting from a point, get the k-layer nodes of this point.
Request example:
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:
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:
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:
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:
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:
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:
🚧, 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.
| Service | Example address or port | Purpose |
|---|---|---|
| HugeGraph-Server | http://127.0.0.1:8080 | Read graph data and write algorithm results; configured by hugegraph.url. |
| etcd | http://127.0.0.1:2379 | BSP job coordination; configured by bsp.etcd_endpoints. |
| Computer master RPC | TCP 8190 | Workers connect to master; port in the distribution configuration. |
| Computer worker data transfer | OS-assigned locally; K8s Operator defaults to 8099 | Transfer vertices and messages between workers. Advertised addresses and ports must be reachable across hosts. |
| MinIO (K8s manifests) | HTTP 9000 | Input 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 service | Operator manifests use cert-manager.io/v1 Certificate, Issuer, and CA injection. Install a compatible version before deploying the Operator. |
| HDFS | Cluster-specific | Required 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.
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.
To build current master, clone and package the source:
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.
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:
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:
If the current read mode hides OLAP writes, an administrator can set it to ALL:
Query vertices to verify the result property:
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:
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.
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:
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.
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:
Inspect runtime information or failure logs:
After checking results, delete the CR to clean up associated resources and restore the default policy:
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, withlog4j2.xmlin the same directory. Mavenpackagecopies these into distributionconf/, runtime dependencies intolib/, and built-in algorithm JARs intoalgorithm/. Templates are maintained in source; there is no separate configuration generator. - Standalone and YARN: startup scripts read distribution
conf/computer.propertiesby default. Override it withbin/start-computer.sh -c <configuration-path>; master and workers need the same job parameters. - Kubernetes Operator: users supply
spec.computerConfin the CRD; the Operator writes a ConfigMap mounted ascomputer.properties. It supplies job ID, worker count, Pod addresses, and etcd when unspecified. Unspecified or zero transfer/RPC ports become8099/8190;transport.server_hostandrpc.server_hostbecome 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)remoteJarUrifor 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.
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 option | default value | description |
|---|---|---|
| hugegraph.url | http://127.0.0.1:8080 | The HugeGraph server URL to load data and write results back. |
| hugegraph.name | hugegraph | The 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.id | local_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_count | 1 | The number of workers for one graph algorithm job. In K8s, this option is set by the Operator. |
| job.partitions_count | 1 | The number of partitions for computing one graph algorithm job. |
| job.partitions_thread_nums | 4 | The number of threads for partition parallel compute. |
2. Algorithm Configuration
Algorithm-specific configuration for computation logic.
| config option | default value | description |
|---|---|---|
| algorithm.params_class | ComputerOptions.Null placeholder class | Required. The class used to pass algorithm parameters before the algorithm runs. |
| algorithm.result_class | ComputerOptions.Null placeholder class | The vertex value class used to store computation results. |
| algorithm.message_class | ComputerOptions.Null placeholder class | The 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 option | default value | description |
|---|---|---|
| input.source_type | hugegraph-server | The 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 option | default value | description |
|---|---|---|
| input.split_size | 1048576 (1 MB) | The input split size in bytes. |
| input.split_max_splits | 10000000 | The maximum number of input splits. |
| input.split_page_size | 500 | The page size for streamed load input split data. |
| input.split_fetch_timeout | 300 | The timeout in seconds to fetch input splits. |
3.3 Input Processing
| config option | default value | description |
|---|---|---|
| input.filter_class | org.apache.hugegraph.computer.core.input.filter.DefaultInputFilter | The class to create input-filter object. Input-filter is used to filter vertex edges according to user needs. |
| input.edge_direction | OUT | The 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_freq | MULTIPLE | The 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_vertex | 200 | The 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 option | default value | description |
|---|---|---|
| input.send_thread_nums | 4 | The 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 option | default value | description |
|---|---|---|
| snapshot.write | false | Whether to write snapshots of input vertex/edge partitions. |
| snapshot.load | false | Whether 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 option | default value | description |
|---|---|---|
| snapshot.minio_endpoint | "" (empty) | MinIO service endpoint (e.g., http://minio:9000). Required when using MinIO. |
| snapshot.minio_access_key | minioadmin | MinIO access key for authentication. |
| snapshot.minio_secret_key | minioadmin | MinIO 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):
Example: MinIO Snapshot (in K8s CRD computerConf):
5. Worker & Master Configuration
Configuration for worker and master computation logic.
5.1 Master Configuration
| config option | default value | description |
|---|---|---|
| master.computation_class | org.apache.hugegraph.computer.core.master.DefaultMasterComputation | Master-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 option | default value | description |
|---|---|---|
| worker.computation_class | org.apache.hugegraph.computer.core.config.Null | The class to create worker-computation object. Worker-computation is used to compute each vertex in each superstep. |
| worker.combiner_class | org.apache.hugegraph.computer.core.config.Null | Combiner can combine messages into one value for a vertex. For example, PageRank algorithm can combine messages of a vertex to a sum value. |
| worker.partitioner | org.apache.hugegraph.computer.core.graph.partition.HashPartitioner | The partitioner that decides which partition a vertex should be in, and which worker a partition should be in. |
5.3 Worker Combiners
| config option | default value | description |
|---|---|---|
| worker.vertex_properties_combiner_class | org.apache.hugegraph.computer.core.combiner.OverwritePropertiesCombiner | The combiner can combine several properties of the same vertex into one properties at input step. |
| worker.edge_properties_combiner_class | org.apache.hugegraph.computer.core.combiner.OverwritePropertiesCombiner | The combiner can combine several properties of the same edge into one properties at input step. |
5.4 Worker Buffers
| config option | default value | description |
|---|---|---|
| worker.received_buffers_bytes_limit | 104857600 (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_capacity | 52428800 (50 MB) | The initial size of write buffer that used to store vertex or message. |
| worker.write_buffer_threshold | 52428800 (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 option | default value | description |
|---|---|---|
| worker.data_dirs | [jobs] | The directories separated by ‘,’ that received vertices and messages can persist into. |
| worker.wait_sort_timeout | 600000 (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_timeout | 86400000 (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 option | default value | description |
|---|---|---|
| output.output_class | org.apache.hugegraph.computer.core.output.LogOutput | The class to output the computation result of each vertex. Called after iteration computation. |
| output.result_name | value | The value is assigned dynamically by #name() of instance created by WORKER_COMPUTATION_CLASS. |
| output.result_write_type | OLAP_COMMON | The result write-type to output to HugeGraph, allowed values: [OLAP_COMMON, OLAP_SECONDARY, OLAP_RANGE]. |
6.2 Output Behavior
| config option | default value | description |
|---|---|---|
| output.with_adjacent_edges | false | Whether to output the adjacent edges of the vertex. |
| output.with_vertex_properties | false | Whether to output the properties of the vertex. |
| output.with_edge_properties | false | Whether to output the properties of the edge. |
6.3 Batch Output
| config option | default value | description |
|---|---|---|
| output.batch_size | 500 | The batch size of output. |
| output.batch_threads | 1 | The number of threads used for batch output. |
| output.single_threads | 1 | The number of threads used for single output. |
6.4 HDFS Output
| config option | default value | description |
|---|---|---|
| output.hdfs_url | hdfs://127.0.0.1:9000 | The HDFS URL for output. |
| output.hdfs_user | hadoop | The HDFS user for output. |
| output.hdfs_path_prefix | /hugegraph-computer/results | The directory of HDFS output results. |
| output.hdfs_delimiter | , (comma) | The delimiter of HDFS output. |
| output.hdfs_merge_partitions | true | Whether to merge output files of multiple partitions. |
| output.hdfs_replication | 3 | The 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_enable | false | Whether 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.conf | Kerberos configuration file path. |
6.5 Retry & Timeout
| config option | default value | description |
|---|---|---|
| output.retry_times | 3 | The retry times when output fails. |
| output.retry_interval | 10 | The retry interval (in seconds) when output fails. |
| output.thread_pool_shutdown_timeout | 60 | The 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 option | default value | description |
|---|---|---|
| transport.server_host | 127.0.0.1 | Worker data listener and advertised address. It must be reachable by other workers across hosts; the Operator sets it to the Pod IP. |
| transport.server_port | 0 (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_threads | 4 | The number of transport threads for server. |
7.2 Client Configuration
| config option | default value | description |
|---|---|---|
| transport.client_threads | 4 | The number of transport threads for client. |
| transport.client_connect_timeout | 3000 | The timeout (in ms) of client connect to server. |
7.3 Protocol Configuration
| config option | default value | description |
|---|---|---|
| transport.provider_class | org.apache.hugegraph.computer.core.network.netty.NettyTransportProvider | The transport provider, currently only supports Netty. |
| transport.io_mode | AUTO | The network IO mode, allowed values: [NIO, EPOLL, AUTO]. AUTO means selecting the appropriate mode automatically. |
| transport.tcp_keep_alive | true | Whether to enable TCP keep-alive. |
| transport.transport_epoll_lt | false | Whether to enable EPOLL level-trigger (only effective when io_mode=EPOLL). |
7.4 Buffer Configuration
| config option | default value | description |
|---|---|---|
| transport.send_buffer_size | 0 | The size of socket send-buffer in bytes. 0 means using system default value. |
| transport.receive_buffer_size | 0 | The size of socket receive-buffer in bytes. 0 means using system default value. |
| transport.write_buffer_high_mark | 67108864 (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_mark | 33554432 (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 option | default value | description |
|---|---|---|
| transport.max_pending_requests | 8 | The max number of client unreceived ACKs. It will trigger sending unavailable if the number of unreceived ACKs >= max_pending_requests. |
| transport.min_pending_requests | 6 | The minimum number of client unreceived ACKs. It will trigger sending available if the number of unreceived ACKs < min_pending_requests. |
| transport.min_ack_interval | 200 | The minimum interval (in ms) of server reply ACK. |
7.6 Timeouts
| config option | default value | description |
|---|---|---|
| transport.close_timeout | 10000 | The timeout (in ms) of close server or close client. |
| transport.sync_request_timeout | 10000 | The timeout (in ms) to wait for response after sending sync-request. |
| transport.finish_session_timeout | 0 | The timeout (in ms) to finish session. 0 means using (transport.sync_request_timeout × transport.max_pending_requests). |
| transport.write_socket_timeout | 3000 | The timeout (in ms) to write data to socket buffer. |
| transport.server_idle_timeout | 360000 (6 minutes) | The max timeout (in ms) of server idle. |
7.7 Heartbeat
| config option | default value | description |
|---|---|---|
| transport.heartbeat_interval | 20000 (20 seconds) | The minimum interval (in ms) between heartbeats on client side. |
| transport.max_timeout_heartbeat_count | 120 | The 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 option | default value | description |
|---|---|---|
| transport.max_syn_backlog | 511 | The capacity of SYN queue on server side. 0 means using system default value. |
| transport.recv_file_mode | true | Whether 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_retries | 3 | The 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.
| Option | Distribution value | Description |
|---|---|---|
| rpc.server_host | 127.0.0.1 (K8s: Pod IP) | Master RPC address, reachable by workers. |
| rpc.server_port | 8190 (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 option | default value | description |
|---|---|---|
| hgkv.max_file_size | 2147483648 (2 GB) | The max number of bytes in each HGKV file. |
| hgkv.max_data_block_size | 65536 (64 KB) | The max byte size of HGKV file data block. |
| hgkv.max_merge_files | 10 | The max number of files to merge at one time. |
| hgkv.temp_file_dir | /tmp/hgkv | This folder is used to store temporary files during the file merging process. |
8.2 Value File Configuration
| config option | default value | description |
|---|---|---|
| valuefile.max_segment_size | 1073741824 (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 option | default value | description |
|---|---|---|
| bsp.etcd_endpoints | http://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_step | 10 (packaged: 2) | The max super step of the algorithm. |
| bsp.register_timeout | 300000 (packaged: 100000) | The max timeout (in ms) to wait for master and workers to register. |
| bsp.wait_workers_timeout | 86400000 (24 hours) | The max timeout (in ms) to wait for workers BSP event. |
| bsp.wait_master_timeout | 86400000 (24 hours) | The max timeout (in ms) to wait for master BSP event. |
| bsp.log_interval | 30000 (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 option | default value | description |
|---|---|---|
| allocator.max_vertices_per_thread | 10000 | Maximum number of vertices per thread processed in each memory allocator. |
| sort.thread_nums | 4 | The 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.
| Option | Managed by | Description |
|---|---|---|
| bsp.etcd_endpoints | K8s Operator | Uses INTERNAL_ETCD_URL only when absent from computerConf. |
| transport.server_host | K8s Operator | Overrides with Pod IP; workers must reach each other. |
| transport.server_port | K8s Operator | Defaults unspecified or zero values to 8099, not a random port. |
| job.id | K8s Operator | Set from CRD job ID. |
| job.workers_count | K8s Operator | Set from CRD workerInstances. |
| rpc.server_host | K8s Operator | Overrides with master Pod IP. |
| rpc.server_port | K8s Operator | Defaults unspecified or zero values to 8190. |
| rpc.remote_url | Computer startup | Removed 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 option | default value | description |
|---|---|---|
| k8s.auto_destroy_pod | true | Delete the job CR after completion or failure; CR deletion triggers cleanup of associated computing resources. |
| k8s.close_reconciler_timeout | 120 | The max timeout (in ms) to close reconciler. |
| k8s.internal_etcd_url | Code: http://127.0.0.1:2379; manifest: http://hugegraph-computer-operator-etcd.hugegraph-computer-operator-system:2379 | etcd URL used by Operator jobs; the supplied manifest uses the etcd Service address. |
| k8s.internal_minio_url | Code: http://127.0.0.1:9000; manifest: http://hugegraph-computer-operator-minio.hugegraph-computer-operator-system:9000 | MinIO address supplied to jobs by the Operator; needed only for MinIO snapshots. |
| k8s.max_reconcile_retry | 3 | The max retry times of reconcile. |
| k8s.probe_backlog | 50 | The maximum backlog for serving health probes. |
| k8s.probe_port | 9892 | The port that the controller binds to for serving health probes. |
| k8s.ready_check_internal | 1000 | The time interval (ms) of check ready. |
| k8s.ready_timeout | 30000 | The max timeout (in ms) of check ready. |
| k8s.reconciler_count | Code: Runtime.getRuntime().availableProcessors(); manifest: 6 | Maximum reconciler thread count. |
| k8s.resync_period | 600000 | The minimum frequency at which watched resources are reconciled. |
| k8s.timezone | Asia/Shanghai | The timezone of computer job and operator. |
| k8s.watch_namespace | hugegraph-computer-operator-system | Namespace watched for custom resources; also used by the supplied manifest. Use * to watch all namespaces. |
HugeGraph-Computer CRD
| spec | default value | description | required |
|---|---|---|---|
| algorithmName | The name of algorithm. | true | |
| jobId | The job id. | true | |
| image | Job container image, which must include the Computer runtime. Bundle algorithm JARs or download them with remoteJarUri. | true | |
| computerConf | The map of computer config options. | true | |
| workerInstances | The number of worker instances, it will override the ‘job.workers_count’ option. | true | |
| pullPolicy | Kubernetes selects Always for latest, otherwise IfNotPresent, when unset | Explicit values: Always, Never, or IfNotPresent. The CRD has no default; see image pull policy. | false |
| pullSecrets | The pull-secrets of Image, detail please refer to: https://kubernetes.io/docs/concepts/containers/images/#specifying-imagepullsecrets-on-a-pod | false | |
| masterCpu | The 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-cpu | false | |
| workerCpu | The 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-cpu | false | |
| masterMemory | The 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-memory | false | |
| workerMemory | The 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-memory | false | |
| log4jXml | The content of log4j.xml for computer job. | false | |
| jarFile | Path to the algorithm JAR inside the image. | false | |
| remoteJarUri | HTTP(S) algorithm JAR URL, downloaded and loaded by the startup script. | false | |
| jvmOptions | The java startup parameters of computer job. | false | |
| envVars | please refer to: https://kubernetes.io/docs/tasks/inject-data-application/define-interdependent-environment-variables/ | false | |
| envFrom | please refer to: https://kubernetes.io/docs/tasks/inject-data-application/define-environment-variable-container/ | false | |
| masterCommand | bin/start-computer.sh | The 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 |
| workerCommand | bin/start-computer.sh | The 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 |
| volumes | Please refer to: https://kubernetes.io/docs/concepts/storage/volumes/ | false | |
| volumeMounts | Please refer to: https://kubernetes.io/docs/concepts/storage/volumes/ | false | |
| secretPaths | The map of k8s-secret name and mount path. | false | |
| configMapPaths | The map of k8s-configmap name and mount path. | false | |
| podTemplateSpec | Please refer to: https://kubernetes.io/docs/reference/kubernetes-api/workload-resources/pod-template-v1/#PodTemplateSpec | false | |
| securityContext | Please refer to: https://kubernetes.io/docs/tasks/configure-pod-container/security-context/ | false |
KubeDriver Config Options
| config option | default value | description |
|---|---|---|
| k8s.build_image_bash_path | The path of command used to build image. | |
| k8s.enable_internal_algorithm | true | Whether enable internal algorithm. |
| k8s.framework_image_url | hugegraph/hugegraph-computer:latest | The image url of computer framework. |
| k8s.image_repository_password | The password for login image repository. | |
| k8s.image_repository_registry | The address for login image repository. | |
| k8s.image_repository_url | hugegraph/hugegraph-computer | The url of image repository. |
| k8s.image_repository_username | The 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_url | hugegraph/hugegraph-computer:latest | The image url of internal algorithm. |
| k8s.jar_file_dir | /cache/jars/ | The directory where the algorithm jar will be uploaded. |
| k8s.kube_config | ~/.kube/config | The path of k8s config file. |
| k8s.log4j_xml_path | The log4j.xml path for computer job. | |
| k8s.namespace | hugegraph-computer-operator-system | Namespace of the HugeGraph-Computer system. |
| k8s.pull_secret_names | [] | The names of pull-secret for pulling image. |