跳转到主要内容

这是本节的多页打印视图。 .

返回本页常规视图.

HugeGraph 图计算(OLAP)

HugeGraph-Computer 仓库包含两套部署和运行方式不同的图计算系统。通用图算法任务默认建议从 Go Vermeer 开始;需要 Java BSP/Pregel 分布式计算模型时,请使用 Computer。两者都可接入 HugeGraph,但配置和作业入口不能互换。

默认入口:Go Vermeer

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

当 load.type=hugegraph 时,Vermeer Master 通过 PD 查询分区元数据,Workers 直接扫描 HStore Store 分区。仅当 output.type=hugegraph 时,Workers 才通过 Server REST API 写回结果;其他数据源或输出配置不需要对应连接。

Java Computer

flowchart TB
    Config["作业配置"] --> Master["Master"] --> Workers["Workers"]
    Master -. BSP .-> Etcd["etcd"]
    Workers <-->|REST| Server["Server"]
    Workers <-->|HDFS| HDFS["HDFS(可选)"]

作业配置可由 Kubernetes Operator 或 YARN 提交。Workers 通过 REST API 从 HugeGraph Server 读取图数据,并可按输出配置写回计算结果;选择 HDFS 输入或输出时,Workers 直接访问 HDFS。Master 使用 etcd 协调 BSP 作业。

1 - HugeGraph-Vermeer 快速上手

Vermeer 高性能内存图计算:一次启动、多次执行,支持 15+ OLAP 算法及秒到分钟级计算,涵盖部署、数据加载、PageRank 和社区发现。

一、Vermeer 概述

1.1 运行架构

Vermeer 是使用 Go 编写的高性能内存优先图计算框架,支持一次启动、多次执行,以及 15+ OLAP 图算法的极速计算,大部分算法可在秒到分钟级完成。实际耗时取决于图规模、算法参数和可用资源。当前由一个 master 调度,可连接多个 worker。

master 是负责通信、转发、汇总的节点,计算量和占用资源量较少。worker 是计算节点,用于存储图数据和运行计算任务,占用大量内存和 cpu。grpc 和 rest 模块分别负责内部通信和外部调用。

启动时,程序先设置内置默认值,再读取工作目录下 config/<env>.ini 的 [default] 节,最后由显式命令行参数覆盖对应配置;例如 --env=master 读取 config/master.ini。Docker 镜像把仓库中的配置复制到 /go/bin/config/,并以 /go/bin/ 为工作目录。挂载宿主机目录到 /go/bin/config 会遮住镜像自带文件,因此该目录必须包含实际使用的 ini 文件。配置读取器不从环境变量取值。

默认端口如下:

角色配置键默认地址用途
masterhttp_peer0.0.0.0:6688REST API;宿主机客户端访问此端口
mastergrpc_peer0.0.0.0:6689worker 连接 master 的 gRPC 端口
workerhttp_peer0.0.0.0:6788worker HTTP 服务
workergrpc_peer0.0.0.0:6789worker gRPC 监听地址,同时会通告给 master 供节点间通信

Docker 示例只把 master 的 HTTP 端口发布到宿主机回环地址 127.0.0.1:6688:6688;master 和 worker 的 gRPC 端口留在容器网络内部。

flowchart LR
  Client["curl / Python 客户端"] -->|"HTTP :6688"| Master["master"]
  Worker["worker"] <-->|"双向 gRPC:master :6689,worker :6789"| Master
  Master -->|"gRPC 查询分区"| PD["HugeGraph PD"]
  Worker -->|"gRPC 扫描分区"| Store["HugeGraph Store"]

1.2 运行方法

警告

生产环境必须启用 HugeGraph Server 认证与授权、IP 白名单和最小权限授权,并保留 audit-*.log 审计记录。Server Auth 不会保护 Vermeer、PD 和 Store 的独立接口;这些 HTTP、gRPC 端口必须限制在可信网络及调用方范围内,Vermeer 对外入口需配置访问控制。

Vermeer 的 master.ini 默认 auth=none,普通及管理 API 均未启用鉴权。本机快速上手只发布回环端口;远程访问前必须启用 auth=token,或通过受保护网络、网关将访问限制到可信调用方。

下面两种 Docker 启动方式都需要先准备一个宿主机配置目录。请在 Vermeer 仓库根目录执行,将项目提供的 master.ini 和 worker.ini 模板复制到该目录;挂载会覆盖镜像里的 /go/bin/config,所以不要把空目录或整个用户主目录挂进去:

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

master.ini 保留 HTTP/gRPC 监听地址 0.0.0.0:6688 和 0.0.0.0:6689。Compose 示例为网络分配固定地址 172.20.0.10(master)和 172.20.0.11(worker),因此将复制后的 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=$

此单 worker 示例将 worker_group 设为 $,表示未绑定命名组时使用的通用组;请在 ini 文件中原样写入美元符号。仓库模板默认是 worker_group=default,若保留命名组,需先将它绑定到任务所在的空间或图,再提交任务。例如默认空间为 $DEFAULT 时,可执行:

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

该接口响应 errcode=0 表示绑定成功;启用 token 鉴权时还需带上授权请求头。master_peer 必须指向 worker 所在网络能访问的 master gRPC 地址。grpc_peer 同时用于 worker 监听和向 master 通告自己的地址,不能设为不可从其他容器访问的 0.0.0.0。如果更改网络子网或地址,需同步修改此处的 IP;多个 worker 还需要各自唯一且可互相访问的 grpc_peer 地址。

  1. 方案一:Docker Compose(推荐)

在 Vermeer 仓库根目录运行。可以修改仓库已有的 docker-compose.yaml,也可以使用下面的配置。仓库文件当前将整个 ~/ 挂载到配置目录且没有发布 master HTTP 端口,运行前必须按示例更换挂载路径并增加端口映射:

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 # master 固定地址

  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 # worker 固定地址

networks:
  vermeer_network:
    driver: bridge
    ipam:
      config:
        - subnet: 172.20.0.0/24 # 按需更换子网

把 /home/user/vermeer-config 换成上面实际的配置目录绝对路径。若修改 vermeer_network 的子网或固定 IP,也要同步更新 worker.ini 中的 grpc_peer 和 master_peer。不要复用默认 worker.ini 的 grpc_peer=0.0.0.0:6789 作为容器间通告地址。

在项目目录构建镜像并启动:

# 构建镜像(在项目根 vermeer 目录)
docker build -t hugegraph/vermeer .

# 启动(在 vermeer 根目录)
docker-compose up -d
# 或使用新版 CLI:
# docker compose up -d

查看日志 / 停止 / 删除:

docker-compose logs -f
docker-compose down
  1. 方案二:通过 docker run 单独启动(手动创建网络并分配静态 IP)

将 CONFIG_DIR 设为上面准备的配置目录,并确保复制后的 worker.ini 已按固定容器地址设置 grpc_peer 和 master_peer。该示例中的 CONFIG_DIR 需替换为实际的绝对路径。

构建镜像:

docker build -t hugegraph/vermeer .

创建自定义 bridge 网络(一次性操作):

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

运行 master(调整 CONFIG_DIR 为您的绝对配置路径,可以根据实际情况调整IP):

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

运行 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

查看日志 / 停止 / 删除:

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

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

# 删除自定义网络(如果需要)
docker network rm vermeer_network
  1. 方案三:从源码构建

构建。具体请参照 Vermeer Readme。

go build

从 Vermeer 仓库根目录启动,例如 ./vermeer --env=master 和 ./vermeer --env=worker01。worker01.ini 中的 grpc_peer 应填写可由 master 和其他 worker 访问、且本机可以绑定的地址;master_peer 指向 master 的 gRPC 监听地址。

启动 master 后,在宿主机验证 HTTP 端口:

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

请求应返回 HTTP 200,JSON 响应中的 errcode 为 0。

二、任务创建类 rest api

2.1 简介

创建任务的流程是先提交 load 任务,等待图加载完成,再提交 compute 任务。同一张已加载的图可重复用于计算,不会因任务完成而自动删除。异步接口返回创建结果和任务信息(包括任务 ID),不代表任务已经完成;同步接口会一直等待任务进入成功或失败状态,调用方及代理的 HTTP 超时需足够长。任务状态可通过查询接口读取:加载成功为 loaded,计算成功为 complete,失败为 error,取消为 canceled,其余状态仍在等待或执行中。图在加载中或错误状态时不能用于计算;删除图要求图处于可删除状态且当前未被使用。

可以使用的 url 如下:

  • 异步接口:POST http://master_ip:port/tasks/create,从响应的 task.id 取得任务 ID。
  • 同步接口:POST http://master_ip:port/tasks/create/sync,等待该任务结束后返回。
  • 查询单个任务:GET http://master_ip:port/task/{task_id};响应中的 errcode 为 0 表示查询成功,再按 task.state 判断业务状态。异步任务应在客户端设置轮询截止时间;到时只停止客户端等待,不会自动取消服务端任务。

2.2 加载图数据

以下示例列出常用的加载参数;引擎会按具体加载器读取参数,未被加载器使用的键不会改变行为。

vermeer提供三种加载方式:

  1. 从本地加载

load.vertex_files 和 load.edge_files 是“worker 地址主机部分到文件路径”的映射;路径由对应 worker 进程读取。容器部署时,数据文件必须先挂载进 worker 容器,并填写容器内路径。在上面的 Compose 示例中,应将主机数据目录挂到 worker 的 /data(例如增加 - /host/data:/data:ro),并把映射键设为 grpc_peer 中的 172.20.0.11;其他部署则使用各自 worker 上报地址的主机部分。

可以预先获取数据集,例如 twitter-2010 数据集。获取方式:https://snap.stanford.edu/data/twitter-2010.html,第一个 twitter-2010.txt.gz 即可。

request 示例:

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 示例:

将请求中的地址、图名和凭据替换为实际连接信息。

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 使用 load.hg_pd_peers 连接 PD 并查询分区,worker 再连接 PD 返回的 Store 地址读取数据。因此这些服务地址必须能从对应的 Vermeer 容器或主机访问;在 Docker 容器里,127.0.0.1 指向容器自身,不能用来代替宿主机或另一容器的地址。

  1. 从hdfs加载

request 示例:

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 输出计算结果

当前有实际写入器的结果方式为 local、hdfs 和 hugegraph,通过 output.type 指定;none 表示不写出结果。源码保留 afs 类型常量,但当前主线没有注册 AFS 加载器或写入器。指定 output.need_statistics 为 1 时,统计结果会写入任务信息;统计算子的适用范围取决于对应算法实现。

以下示例列出常用的计算和结果输出参数;算法支持的参数以 Vermeer 当前实现为准。

request 示例:

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 的结果写在执行任务的 worker 本地文件系统中;容器部署时,如需从宿主机读取结果,请将 worker 的输出目录挂载出来。

三、支持的算法

3.1 PageRank

PageRank 算法又称网页排名算法,是一种由搜索引擎根据网页(节点)之间相互的超链接进行计算的技

术,用来体现网页(节点)的相关性和重要性。

  • 如果一个网页被很多其他网页链接到,说明这个网页比较重要,也就是其 PageRank 值会相对较高。
  • 如果一个 PageRank 值很高的网页链接到其他网页,那么被链接到的网页的 PageRank 值会相应地提高。

PageRank 算法适用于网页排序、社交网络重点人物发掘等场景。

request 示例:

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(弱连通分量)

弱连通分量,计算无向图中所有联通的子图,输出各顶点所属的弱联通子图 id,表明各个点之间的连通性,区分不同的连通社区。

request 示例:

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(标签传播)

标签传递算法,是一种图聚类算法,常用在社交网络中,用于发现潜在的社区。

request 示例:

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(度中心性)

度中心性算法,算法用于计算图中每个节点的度中心性值,支持无向图和有向图。度中心性是衡量节点重要性的重要指标,节点与其它节点的边越多,则节点的度中心性值越大,节点在图中的重要性也就越高。在无向图中,度中心性的计算是基于边信息统计节点出现次数,得出节点的度中心性的值,在有向图中则基于边的方向进行筛选,基于输入边或输出边信息统计节点出现次数,得到节点的入度值或出度值。它表明各个点的重要性,一般越重要的点度数越高。

request 示例:

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)用于计算一个节点到所有其他可达节点的最短距离的倒数,进行累积后归一化的值。紧密中心度可以用来衡量信息从该节点传输到其他节点的时间长短。节点的“Closeness Centrality”越大,其在所在图中的位置越靠近中心,适用于社交网络中关键节点发掘等场景。

request 示例:

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(中介中心性算法)

中介中心性算法(Betweeness Centrality)判断一个节点具有"桥梁"节点的值,值越大说明它作为图中两点间必经路径的可能性越大,典型的例子包括社交网络中的共同关注的人。适用于衡量社群围绕某个节点的聚集程度。

request 示例:

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(三角形计数)

三角形计数算法,用于计算通过每个顶点的三角形个数,适用于计算用户之间的关系,关联性是不是成三角形。三角形越多,代表图中节点关联程度越高,组织关系越严密。社交网络中的三角形表示存在有凝聚力的社区,识别三角形有助于理解网络中个人或群体的聚类和相互联系。在金融网络或交易网络中,三角形的存在可能表示存在可疑或欺诈活动,三角形计数可以帮助识别可能需要进一步调查的交易模式。

输出的结果为 每个顶点对应一个 Triangle Count,即为每个顶点所在三角形的个数。

注:该算法为无向图算法,忽略边的方向。

request 示例:

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

K-Core 算法,标记所有度数为 K 的顶点,适用于图的剪枝,查找图的核心部分。

request 示例:

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(单元最短路径)

单源最短路径算法,求一个点到其他所有点的最短距离。

request 示例:

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

以一个点为起点,获取这个点的 k 层的节点。

request 示例:

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

Louvain 算法是一种基于模块度的社区发现算法。其基本思想是网络中节点尝试遍历所有邻居的社区标签,并选择最大化模块度增量的社区标签。在最大化模块度之后,每个社区看成一个新的节点,重复直到模块度不再增大。

Vermeer 上实现的分布式 Louvain 算法受节点顺序、并行计算等因素影响,并且由于 Louvain 算法由于其遍历顺序的随机导致社区压缩也具有一定的随机性,导致重复多次执行可能存在不同的结果。但整体趋势不会有大的变化。

request 示例:

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 相似度系数

Jaccard index , 又称为 Jaccard 相似系数(Jaccard similarity coefficient)用于比较有限样本集之间的相似性与差异性。Jaccard 系数值越大,样本相似度越高。用于计算一个给定的源点,与图中其他所有点的 Jaccard 相似系数。

request 示例:

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

个性化的 pagerank 的目标是要计算所有节点相对于用户 u 的相关度。从用户 u 对应的节点开始游走,每到一个节点都以 1-d 的概率停止游走并从 u 重新开始,或者以 d 的概率继续游走,从当前节点指向的节点中按照均匀分布随机选择一个节点往下游走。用于给定一个起点,计算此起点开始游走的个性化 pagerank 得分。适用于社交推荐等场景。

由于计算需要使用出度,需要在读取图时需要设置 load.use_out_degree 为 1。

request 示例:

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 全图 Kout

计算图的所有节点的k度邻居(不包含自己以及1~k-1度的邻居),由于全图kout算法内存膨胀比较厉害,目前k限制在1和2,另外,全局kout算法支持过滤功能( 参数如:“compute.filter”:“risk_level==1”),在计算第k度的是时候进行过滤条件的判断,符合过滤条件的进入最终结果集,算法最终输出是符合条件的邻居个数。

request 示例:

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

集聚系数表示一个图中节点聚集程度的系数。在现实的网络中,尤其是在特定的网络中,由于相对高密度连接点的关系,节点总是趋向于建立一组严密的组织关系。集聚系数算法(Cluster Coefficient)用于计算图中节点的聚集程度。本算法为局部集聚系数。局部集聚系数可以测量图中每一个结点附近的集聚程度。

request 示例:

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(强连通分量)

在有向图的数学理论中,如果一个图的每一个顶点都可从该图其他任意一点到达,则称该图是强连通的。在任意有向图中能够实现强连通的部分我们称其为强连通分量。它表明各个点之间的连通性,区分不同的连通社区。

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"
 }
}

🚧, 后续随时更新完善,欢迎随时提出建议和意见。

2 - HugeGraph-Computer 快速上手

1. 组件说明

HugeGraph-Computer 是基于 BSP(Bulk Synchronous Parallel,批量同步并行)模型的 Java 分布式图计算框架,算法按超步迭代运行。它可由 Kubernetes Operator 或 YARN 调度,也可在单机上启动 master 和 worker 进程进行小规模试跑。

它与同一仓库中的 Vermeer 是两套实现:Computer 使用 Java/BSP 运行时,面向分布式计算;Vermeer 是 Go 实现的内存图计算平台,采用 master-worker 结构。两者共享 HugeGraph 数据源,但部署和任务配置不能互换。

Computer 支持从 HugeGraph 或 HDFS 读取图数据,并可将结果写回 HugeGraph 或 HDFS。运行时可把部分数据溢写到磁盘;能否完成任务仍取决于输入规模、资源和配置,不应把溢写能力理解为不受资源限制。

2. 前置条件与连接

构建和运行需要 JDK 11 或更高版本。源码构建还需要 Maven 3.5 或更高版本。PageRank 示例要求一个已启动且包含待计算图数据的 HugeGraph-Server,以及可供 master、worker 访问的 etcd。

服务示例地址或端口用途
HugeGraph-Serverhttp://127.0.0.1:8080读取图数据并写回算法结果;以 hugegraph.url 为准。
etcdhttp://127.0.0.1:2379BSP 作业协调;以 bsp.etcd_endpoints 为准。
Computer master RPCTCP 8190worker 连接 master;发行配置中的端口。
Computer worker 数据传输本地默认由系统分配;K8s Operator 默认 8099worker 之间传输顶点和消息。跨主机时需保证公告地址和端口可达。
MinIO(K8s 清单)HTTP 9000用于输入分区快照。Computer 1.7.0 清单把 Service 的 9000 错映射到 MinIO Console 的 9090;启用快照前需将 targetPort 改为 9000。默认 snapshot.write=false、snapshot.load=false,不启用快照时无需访问 MinIO。
cert-manager(K8s)集群内服务Operator 清单使用 cert-manager.io/v1 的 Certificate、Issuer 和 CA 注入功能;部署 Operator 前必须先安装兼容版本。
HDFS按集群配置仅当输入或输出配置为 HDFS 时需要。

Kubernetes 作业中的 hugegraph.url 必须是各计算 Pod 都能访问的地址,不能填仅在个人电脑上可用的 localhost。如果启用了 HugeGraph 认证,应在配置中填写用户名和密码,并为 REST 查询使用对应凭据。

警告

生产环境必须为 HugeGraph Server 开启认证与授权(见认证与授权说明)并保留 Server 审计日志(标准日志文件为 audit-*.log),为 Computer 使用只具备作业所需读写权限的专用账号,并为 Server 网络入口设置来源 IP 白名单。hugegraph.username 和 hugegraph.password 只是 Computer 连接 Server 的凭据,不能替代 Server 端认证。

更多配置项见Computer 配置参考。

3. 获取源码并构建发行包

Apache HugeGraph 下载目录按版本提供 Computer 源码包,不提供单独的预编译 Computer 二进制包。下面以 VERSION=1.7.0 为例;该发布版的源码包和构建产物文件名含 -incubating-,其他版本请按下载目录中的实际文件名调整。

VERSION=1.7.0 # 替换为要下载的发布版本
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}"

从当前 master 构建时可直接克隆并打包:

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

发行包由 computer/computer-dist 组装,包含 bin/start-computer.sh、运行依赖 lib/、内置算法 algorithm/builtin-algorithm.jar 及默认配置 conf/computer.properties、conf/log4j2.xml。源码默认配置位于 computer/computer-dist/src/assembly/static/conf/computer.properties。1.7.0 发布源码的构建产物带 -incubating-,当前 master 的发行名不带该标记,且版本号由源码 POM 决定;不要混用两种源码对应的 tar 包名。

4. 单机运行 PageRank

先编辑发行目录中的 conf/computer.properties,按实际环境修改 HugeGraph 地址、图名和认证信息,并确认 bsp.etcd_endpoints 指向可访问的 etcd。默认配置已选择内置 PageRankParams;master 和 worker 必须使用同一份配置及相同的 job.id。每个并行作业应使用不同的 job.id。

在两个终端中都从发行目录启动进程。启动脚本默认读取 conf/computer.properties,也可用 -c 指定配置文件。

# 终端一:启动 master
bin/start-computer.sh -d local -r master
# 终端二:启动 worker
bin/start-computer.sh -d local -r worker

master 会等待配置要求的 worker 注册后运行作业。查看两个终端输出;默认日志配置也会在当前发行目录的 logs/ 下写入 master 和 worker 日志。只有进程启动成功并不代表计算完成,应确认 master 日志中输入、超步计算和输出阶段均正常结束。

PageRank 参数类将结果写回 HugeGraph,属性名为 page_rank。运行前先检查图中每个目标顶点类型都允许此属性;Computer 会按 DOUBLE、OLAP_COMMON 创建属性键,但不会把它加入顶点类型。若属性键不存在,可先创建:

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'

如果该属性键已存在,确认其类型为 DOUBLE 且 write_type 为 OLAP_COMMON。先列出图中的顶点类型,再对每个要计算的类型添加可空的 page_rank 属性;把示例中的 person 替换为实际类型名:

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'

若图当前读模式不显示 OLAP 写入,可由管理员把读模式设为 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'

查询顶点以确认结果属性:

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

需要认证时,在 curl 命令中添加 --user "$HG_USER:$HG_PASSWORD"。读模式接口和权限说明见图读模式 REST API。

5. 在 Kubernetes 中运行 PageRank

先确保 HugeGraph-Server 对计算 Pod 可达,并按 cert-manager 官方安装文档安装与集群兼容的 cert-manager。Computer Operator 清单会部署 Operator、etcd 和 MinIO;CRD 与 Operator 清单应使用同一 Computer 发行版本。下面以 1.7.0 发布版为例:

VERSION=1.7.0 # CRD 和 Operator 清单使用同一发布版本
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
可选:启用 MinIO 快照(1.7.0)

1.7.0 清单中的 MinIO Service 把 S3 API 的 9000 端口映射到了 Console 的 9090;启用快照时先修正映射,并将 snapshot.minio_endpoint 设为 http://hugegraph-computer-operator-minio.hugegraph-computer-operator-system.svc:9000。默认不启用快照时可跳过。

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

替换 HugeGraph 地址为计算 Pod 可访问的服务地址后,提交 HugeGraphComputerJob。示例使用官方 Docker Hub 的 hugegraph/hugegraph-computer:latest,并引用镜像内置的 PageRank JAR;自定义算法 JAR 可预置在镜像内通过 jarFile 指定,也可用 remoteJarUri 从 HTTP(S) 地址下载。分区数需不小于 worker 数量。

将以下 YAML 保存为 pagerank-job.yaml,然后应用该文件:

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 # 官方 Docker Hub 镜像
  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 表示作业完成。Operator 默认会在作业完成后删除 CR 和计算资源;需要保留结果或排查失败作业时,可按需展开下方说明。

可选:保留作业状态、查看日志或清理资源

Operator 默认 AUTO_DESTROY_POD=true,短作业的 CR 和 Pod 可能在查看前被清理。若要保留它们,请在提交作业前关闭该选项,并等待 Controller 新 Pod 就绪后再提交;该设置必须改在 Java Operator 的 controller 容器中:

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

查看运行信息或失败日志:

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

查询完结果后删除 CR 以清理关联资源,并恢复默认策略:

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

PageRank 写回 HugeGraph 后按上一节设置读模式并查询。若改用 HDFS 输出,结果位于 output.hdfs_path_prefix/<job.id>/ 下,文件名和分区布局由作业配置决定。

完整 CRD 字段见Computer 配置参考中的 CRD 说明。

6. 内置算法与开发入口

当前源码中的内置算法包括:

  • 中心性:PageRank、Betweenness Centrality、Closeness Centrality、Degree Centrality。
  • 社区与结构:Clustering Coefficient、K-core、LPA、Triangle Count、WCC。
  • 路径与采样:环检测、带过滤的环检测、单源最短路径、Random Walk。

算法实现位于 computer/computer-algorithm。自定义算法需要遵循 Computer API 并打包为可加载的 JAR;模块划分和开发入口见Computer README。

3 - HugeGraph-Computer 配置参考

Computer 配置选项

表格中的默认值来自 computer-api 模块的 ComputerOptions.java;发行包的 conf/computer.properties 显式设置的值以“代码默认值(发行包:实际值)”标注。rpc.* 等通用配置由 HugeGraph Commons 提供,以下会单独标出发行包实际值。

配置文件从哪里来

  • 源码模板:computer/computer-dist/src/assembly/static/conf/computer.properties;日志配置模板为同目录的 log4j2.xml。Maven package 阶段将模板复制进发行目录的 conf/,并将运行依赖和内置算法 JAR 装入 lib/、algorithm/。模板由源码维护,没有单独的配置生成命令。
  • 单机、YARN:启动脚本默认读取发行目录的 conf/computer.properties,可使用 bin/start-computer.sh -c <配置文件路径> 覆盖;master 和 worker 应读到同一组作业参数。
  • Kubernetes Operator:用户在 CRD 的 spec.computerConf 中提供键值;Operator 将它写成 ConfigMap,挂载为容器内的 computer.properties。Operator 会补入作业 ID、worker 数量、Pod 地址和未指定时的 etcd 地址,并在未指定或设为 0 时将数据传输端口设为 8099、master RPC 端口设为 8190。transport.server_host 与 rpc.server_host 会被设置为 Pod IP。作业容器启动脚本只对 ${POD_IP}、${HOSTNAME}、${POD_NAME}、${POD_NAMESPACE} 做环境变量替换。
  • Kubernetes 作业镜像必须包含 Computer 运行时。算法 JAR 可以预置在镜像内并由 CRD 的 jarFile 指向,也可以通过 remoteJarUri 设置 HTTP(S) 地址,让启动脚本下载并加载。示例与字段见Computer 快速上手及下方 CRD 表。

Apache HugeGraph 下载目录提供各发布版本的 Computer 源码包,不提供单独的预编译二进制包。1.7.0 的发布源码包名含 -incubating-;在该 tag 的 computer/ Maven 聚合工程运行 mvn clean package -DskipTests 后,发行包也带 -incubating-。当前 master 的发行名不含该标记,版本号由 POM 决定。源码包不能作为已构建发行包直接启动。

警告

表格中的空 HugeGraph 凭据和 MinIO 示例密钥仅适用于本地或演示环境。生产环境必须为 HugeGraph Server 开启认证与授权(见认证与授权说明)并保留 Server 审计日志(标准日志文件为 audit-*.log),为 Computer 配置最小必要权限的 Server 账号,并为 Server 网络入口设置来源 IP 白名单;不要将示例 MinIO 密钥用于生产。


1. 基础配置

HugeGraph-Computer 核心作业设置。

配置项默认值说明
hugegraph.urlhttp://127.0.0.1:8080HugeGraph 服务器 URL,用于加载数据和写回结果。
hugegraph.namehugegraph图名称,用于加载数据和写回结果。
hugegraph.username"" (空)HugeGraph 认证用户名(如果未启用认证则留空)。
hugegraph.password"" (空)HugeGraph 认证密码(如果未启用认证则留空)。
job.idlocal_0001 (打包: local_001)YARN 集群或 K8s 集群上的作业标识符。
job.namespace"" (空)可选的 etcd 作业键名前缀,可用于隔离不同命名空间;它不是 Kubernetes metadata.namespace,Operator 不会自动填入。
job.workers_count1执行一个图算法作业的 Worker 数量。在 K8s 中由 Operator 设置。
job.partitions_count1执行一个图算法作业的分区数量。
job.partitions_thread_nums4分区并行计算的线程数量。

2. 算法配置

计算逻辑的算法特定配置。

配置项默认值说明
algorithm.params_classComputerOptions.Null 占位类必填。用于在算法运行前传递算法参数的类。
algorithm.result_classComputerOptions.Null 占位类顶点值的类,用于存储顶点的计算结果。
algorithm.message_classComputerOptions.Null 占位类计算顶点时传递的消息类。

3. 输入配置

从 HugeGraph 或其他数据源加载输入数据的配置。

3.1 输入源

配置项默认值说明
input.source_typehugegraph-server加载输入数据的源类型,允许值:[‘hugegraph-server’, ‘hugegraph-loader’]。‘hugegraph-loader’ 表示使用 hugegraph-loader 从 HDFS 或文件加载数据。如果使用 ‘hugegraph-loader’,请配置 ‘input.loader_struct_path’ 和 ‘input.loader_schema_path’。
input.loader_struct_path"" (空)Loader 输入的结构路径,仅在 input.source_type=hugegraph-loader 时生效。
input.loader_schema_path"" (空)Loader 输入的 Schema 路径,仅在 input.source_type=hugegraph-loader 时生效。

3.2 输入分片

配置项默认值说明
input.split_size1048576 (1 MB)输入分片大小(字节)。
input.split_max_splits10000000最大输入分片数量。
input.split_page_size500流式加载输入分片数据的页面大小。
input.split_fetch_timeout300获取输入分片的超时时间(秒)。

3.3 输入处理

配置项默认值说明
input.filter_classorg.apache.hugegraph.computer.core.input.filter.DefaultInputFilter创建输入过滤器对象的类。输入过滤器用于根据用户需求过滤顶点边。
input.edge_directionOUT要加载的边的方向,允许值:[OUT, IN, BOTH]。当值为 BOTH 时,将加载 OUT 和 IN 两个方向的边。
input.edge_freqMULTIPLE一对顶点之间可以存在的边的频率,允许值:[SINGLE, SINGLE_PER_LABEL, MULTIPLE]。SINGLE 表示一对顶点之间只能存在一条边(通过 sourceId + targetId 标识);SINGLE_PER_LABEL 表示每个边标签在一对顶点之间可以有一条边(通过 sourceId + edgeLabel + targetId 标识);MULTIPLE 表示一对顶点之间可以存在多条边(通过 sourceId + edgeLabel + sortValues + targetId 标识)。
input.max_edges_in_one_vertex200允许附加到一个顶点的最大邻接边数量。邻接边将作为一个批处理单元一起存储和传输。

3.4 输入性能

配置项默认值说明
input.send_thread_nums4并行发送顶点或边的线程数量。

4. 快照与存储配置

HugeGraph-Computer 支持快照功能,可将顶点/边分区保存到本地存储或 MinIO 对象存储,用于断点恢复或加速重复计算。

4.1 基础快照配置

配置项默认值说明
snapshot.writefalse是否写入输入顶点/边分区的快照。
snapshot.loadfalse是否从顶点/边分区的快照加载。
snapshot.name"" (空)用户自定义的快照名称,用于区分不同的快照。

4.2 MinIO 集成(可选)

MinIO 可用作 K8s 部署中快照的分布式对象存储后端。

配置项默认值说明
snapshot.minio_endpoint"" (空)MinIO 服务端点(例如 http://minio:9000)。使用 MinIO 时必填。
snapshot.minio_access_keyminioadminMinIO 认证访问密钥。
snapshot.minio_secret_keyminioadminMinIO 认证密钥。
snapshot.minio_bucket_name"" (空)用于存储快照数据的 MinIO 存储桶名称。

使用场景:

  • 断点恢复:作业失败后从快照恢复,避免重新加载数据
  • 重复计算:多次运行同一算法时从快照加载数据以加速启动
  • A/B 测试:保存同一数据集的多个快照版本,测试不同的算法参数

示例:本地快照(在 computer.properties 中):

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

示例:MinIO 快照(在 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 配置

Worker 和 Master 计算逻辑的配置。

5.1 Master 配置

配置项默认值说明
master.computation_classorg.apache.hugegraph.computer.core.master.DefaultMasterComputationMaster 计算是可以决定是否继续下一个超步的计算。它在每个超步结束时在 master 上运行。

5.2 Worker 计算

配置项默认值说明
worker.computation_classorg.apache.hugegraph.computer.core.config.Null创建 worker 计算对象的类。Worker 计算用于在每个超步中计算每个顶点。
worker.combiner_classorg.apache.hugegraph.computer.core.config.NullCombiner 可以将消息组合为一个顶点的一个值。例如,PageRank 算法可以将一个顶点的消息组合为一个求和值。
worker.partitionerorg.apache.hugegraph.computer.core.graph.partition.HashPartitioner分区器,决定顶点应该在哪个分区中,以及分区应该在哪个 worker 中。

5.3 Worker 组合器

配置项默认值说明
worker.vertex_properties_combiner_classorg.apache.hugegraph.computer.core.combiner.OverwritePropertiesCombiner组合器可以在输入步骤将同一顶点的多个属性组合为一个属性。
worker.edge_properties_combiner_classorg.apache.hugegraph.computer.core.combiner.OverwritePropertiesCombiner组合器可以在输入步骤将同一边的多个属性组合为一个属性。

5.4 Worker 缓冲区

配置项默认值说明
worker.received_buffers_bytes_limit104857600 (100 MB)接收数据缓冲区的限制字节数。所有缓冲区的总大小不能超过此限制。如果接收缓冲区达到此限制,它们将被合并到文件中(溢出到磁盘)。
worker.write_buffer_capacity52428800 (50 MB)用于存储顶点或消息的写缓冲区的初始大小。
worker.write_buffer_threshold52428800 (50 MB)写缓冲区的阈值。超过它将触发排序。写缓冲区用于存储顶点或消息。

5.5 Worker 数据与超时

配置项默认值说明
worker.data_dirs[jobs]用逗号分隔的目录,接收的顶点和消息可以持久化到其中。
worker.wait_sort_timeout600000 (10 分钟)消息处理程序等待排序线程对一批缓冲区进行排序的最大超时时间(毫秒)。
worker.wait_finish_messages_timeout86400000 (24 小时)消息处理程序等待所有 worker 完成消息的最大超时时间(毫秒)。

6. I/O 与输出配置

输出计算结果的配置。

6.1 输出类与结果

配置项默认值说明
output.output_classorg.apache.hugegraph.computer.core.output.LogOutput输出每个顶点计算结果的类。在迭代计算后调用。
output.result_namevalue该值由 WORKER_COMPUTATION_CLASS 创建的实例的 #name() 动态分配。
output.result_write_typeOLAP_COMMON输出到 HugeGraph 的结果写入类型,允许值:[OLAP_COMMON, OLAP_SECONDARY, OLAP_RANGE]。

6.2 输出行为

配置项默认值说明
output.with_adjacent_edgesfalse是否输出顶点的邻接边。
output.with_vertex_propertiesfalse是否输出顶点的属性。
output.with_edge_propertiesfalse是否输出边的属性。

6.3 批量输出

配置项默认值说明
output.batch_size500输出的批处理大小。
output.batch_threads1用于批量输出的线程数量。
output.single_threads1用于单个输出的线程数量。

6.4 HDFS 输出

配置项默认值说明
output.hdfs_urlhdfs://127.0.0.1:9000输出的 HDFS URL。
output.hdfs_userhadoop输出的 HDFS 用户。
output.hdfs_path_prefix/hugegraph-computer/resultsHDFS 输出结果的目录。
output.hdfs_delimiter, (逗号)HDFS 输出的分隔符。
output.hdfs_merge_partitionstrue是否合并多个分区的输出文件。
output.hdfs_replication3HDFS 的副本数。
output.hdfs_core_site_path"" (空)HDFS core site 路径。
output.hdfs_site_path"" (空)HDFS site 路径。
output.hdfs_kerberos_enablefalse是否为 HDFS 启用 Kerberos 认证。
output.hdfs_kerberos_principal"" (空)HDFS 的 Kerberos 认证 principal。
output.hdfs_kerberos_keytab"" (空)HDFS 的 Kerberos 认证 keytab 文件。
output.hdfs_krb5_conf/etc/krb5.confKerberos 配置文件路径。

6.5 重试与超时

配置项默认值说明
output.retry_times3输出失败时的重试次数。
output.retry_interval10输出失败时的重试间隔(秒)。
output.thread_pool_shutdown_timeout60输出线程池关闭的超时时间(秒)。

7. 网络与传输配置

Worker 和 Master 之间网络通信的配置。

7.1 服务器配置

配置项默认值说明
transport.server_host127.0.0.1worker 数据服务监听并向其他 worker 公告的地址。多主机场景必须使用其他节点可达的地址;Operator 会改为 Pod IP。
transport.server_port0(发行包:0;K8s Operator:8099)worker 数据传输监听端口。单机值 0 表示系统分配端口;Operator 会将 K8s 默认值改为固定端口 8099,以便 Pod 间通信。
transport.server_threads4服务器传输线程的数量。

7.2 客户端配置

配置项默认值说明
transport.client_threads4客户端传输线程的数量。
transport.client_connect_timeout3000客户端连接到服务器的超时时间(毫秒)。

7.3 协议配置

配置项默认值说明
transport.provider_classorg.apache.hugegraph.computer.core.network.netty.NettyTransportProvider传输提供程序,目前仅支持 Netty。
transport.io_modeAUTO网络 IO 模式,允许值:[NIO, EPOLL, AUTO]。AUTO 表示自动选择适当的模式。
transport.tcp_keep_alivetrue是否启用 TCP keep-alive。
transport.transport_epoll_ltfalse是否启用 EPOLL 水平触发(仅在 io_mode=EPOLL 时有效)。

7.4 缓冲区配置

配置项默认值说明
transport.send_buffer_size0Socket 发送缓冲区大小(字节)。0 表示使用系统默认值。
transport.receive_buffer_size0Socket 接收缓冲区大小(字节)。0 表示使用系统默认值。
transport.write_buffer_high_mark67108864 (64 MB)写缓冲区的高水位标记(字节)。如果排队字节数 > write_buffer_high_mark,将触发发送不可用。
transport.write_buffer_low_mark33554432 (32 MB)写缓冲区的低水位标记(字节)。如果排队字节数 < write_buffer_low_mark,将触发发送可用。

7.5 流量控制

配置项默认值说明
transport.max_pending_requests8客户端未接收 ACK 的最大数量。如果未接收 ACK 的数量 >= max_pending_requests,将触发发送不可用。
transport.min_pending_requests6客户端未接收 ACK 的最小数量。如果未接收 ACK 的数量 < min_pending_requests,将触发发送可用。
transport.min_ack_interval200服务器回复 ACK 的最小间隔(毫秒)。

7.6 超时配置

配置项默认值说明
transport.close_timeout10000关闭服务器或关闭客户端的超时时间(毫秒)。
transport.sync_request_timeout10000发送同步请求后等待响应的超时时间(毫秒)。
transport.finish_session_timeout0完成会话的超时时间(毫秒)。0 表示使用 (transport.sync_request_timeout × transport.max_pending_requests)。
transport.write_socket_timeout3000将数据写入 socket 缓冲区的超时时间(毫秒)。
transport.server_idle_timeout360000 (6 分钟)服务器空闲的最大超时时间(毫秒)。

7.7 心跳配置

配置项默认值说明
transport.heartbeat_interval20000 (20 秒)客户端心跳之间的最小间隔(毫秒)。
transport.max_timeout_heartbeat_count120客户端超时心跳的最大次数。如果连续等待心跳响应超时的次数 > max_timeout_heartbeat_count,通道将从客户端关闭。

7.8 高级网络设置

配置项默认值说明
transport.max_syn_backlog511服务器端 SYN 队列的容量。0 表示使用系统默认值。
transport.recv_file_modetrue是否启用接收缓冲文件模式。如果启用,将使用零拷贝从 socket 接收缓冲区并写入文件。注意:需要操作系统支持零拷贝(例如 Linux sendfile/splice)。
transport.network_retries3网络通信不稳定时的重试次数。

7.9 Master RPC 配置

以下两个键由 HugeGraph Commons 的 RPC 配置定义;computer.properties 发行模板显式设置了主机和端口。Kubernetes Operator 会公告 Pod IP,并在未指定有效端口时使用 8190。

配置项发行包值说明
rpc.server_host127.0.0.1(K8s:Pod IP)Master RPC 服务地址,worker 必须能访问此地址。
rpc.server_port8190(K8s 默认:8190)Master RPC 监听端口;集群安全组和网络策略须允许 worker 连接。

8. 存储与持久化配置

HGKV(HugeGraph Key-Value)存储引擎和值文件的配置。

8.1 HGKV 配置

配置项默认值说明
hgkv.max_file_size2147483648 (2 GB)每个 HGKV 文件的最大字节数。
hgkv.max_data_block_size65536 (64 KB)HGKV 文件数据块的最大字节大小。
hgkv.max_merge_files10一次合并的最大文件数。
hgkv.temp_file_dir/tmp/hgkv此文件夹用于在文件合并过程中存储临时文件。

8.2 值文件配置

配置项默认值说明
valuefile.max_segment_size1073741824 (1 GB)值文件每个段的最大字节数。

9. BSP 与协调配置

批量同步并行(BSP)协议和 etcd 协调的配置。

配置项默认值说明
bsp.etcd_endpointshttp://localhost:2379(发行包:http://127.0.0.1:2379)etcd 客户端端点;多个地址用逗号分隔。K8s Operator 在未由 computerConf 指定时使用 Operator 的 INTERNAL_ETCD_URL。
bsp.max_super_step10 (打包: 2)算法的最大超步数。
bsp.register_timeout300000 (打包: 100000)等待 master 和 worker 注册的最大超时时间(毫秒)。
bsp.wait_workers_timeout86400000 (24 小时)等待 worker BSP 事件的最大超时时间(毫秒)。
bsp.wait_master_timeout86400000 (24 小时)等待 master BSP 事件的最大超时时间(毫秒)。
bsp.log_interval30000 (30 秒)等待 BSP 事件时打印日志的日志间隔(毫秒)。

10. 性能调优配置

性能优化的配置。

配置项默认值说明
allocator.max_vertices_per_thread10000每个内存分配器中每个线程处理的最大顶点数。
sort.thread_nums4执行内部排序的线程数量。

11. 系统管理配置

以下配置在 Kubernetes 作业中由 Operator 自动补齐或覆盖。除非自定义了 Operator 或网络,不应手动覆盖这些值;job.namespace 是可选的用户配置,见基础配置表。

配置项管理者说明
bsp.etcd_endpointsK8s Operator仅在 computerConf 未设置时,取 Operator 的 INTERNAL_ETCD_URL。
transport.server_hostK8s Operator覆盖为 Pod IP;worker 需要互相访问。
transport.server_portK8s Operator未设置或为 0 时设置为 8099,不是随机端口。
job.idK8s Operator自动从 CRD 设置为作业 ID
job.workers_countK8s Operator自动从 CRD workerInstances 设置
rpc.server_hostK8s Operator覆盖为 master Pod IP。
rpc.server_portK8s Operator未设置或为 0 时设置为 8190。
rpc.remote_urlComputer 启动流程读取配置时会移除该项,不应作为作业配置设置。

为什么禁止修改:

  • BSP/RPC 配置:必须与实际部署的 etcd/RPC 服务匹配。手动覆盖会破坏协调。
  • 作业配置:必须与 K8s CRD 规范匹配。不匹配会导致 worker 数量错误。
  • 传输配置:必须使用 worker 互相可达的 Pod IP 和端口;K8s 的默认端口为 8099。

K8s Operator 配置选项

注意:选项需要通过环境变量设置进行转换,例如 k8s.internal_etcd_url => INTERNAL_ETCD_URL

配置项默认值说明
k8s.auto_destroy_podtrue作业完成或失败后是否删除作业 CR;CR 删除后 Operator 会清理关联计算资源。
k8s.close_reconciler_timeout120关闭 reconciler 的最大超时时间(毫秒)。
k8s.internal_etcd_url代码默认 http://127.0.0.1:2379;发行清单覆盖为 http://hugegraph-computer-operator-etcd.hugegraph-computer-operator-system:2379Operator 作业使用的 etcd URL。使用随附清单时实际值是 etcd Service 地址。
k8s.internal_minio_url代码默认 http://127.0.0.1:9000;发行清单覆盖为 http://hugegraph-computer-operator-minio.hugegraph-computer-operator-system:9000Operator 为作业补入的 MinIO 地址;仅启用 MinIO 快照时需要。
k8s.max_reconcile_retry3reconcile 的最大重试次数。
k8s.probe_backlog50服务健康探针的最大积压。
k8s.probe_port9892controller 绑定的用于服务健康探针的端口。
k8s.ready_check_internal1000检查就绪的时间间隔(毫秒)。
k8s.ready_timeout30000检查就绪的最大超时时间(毫秒)。
k8s.reconciler_count代码默认 Runtime.getRuntime().availableProcessors();发行清单覆盖为 6reconciler 线程的最大数量。
k8s.resync_period600000被监视资源进行 reconcile 的最小频率。
k8s.timezoneAsia/Shanghaicomputer 作业和 operator 的时区。
k8s.watch_namespacehugegraph-computer-operator-system监视自定义资源的命名空间。随附发行清单也设置为此值;使用 * 可监视所有命名空间。

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

字段默认值说明必填
algorithmName算法名称。true
jobId作业 ID。true
imageComputer 作业使用的容器镜像,必须包含 Computer 运行时。目标算法 JAR 可预置在镜像内,或通过 remoteJarUri 下载。true
computerConfcomputer 配置选项的映射。true
workerInstancesworker 实例数量,将覆盖 ‘job.workers_count’ 选项。true
pullPolicy未设置时由 Kubernetes 按镜像标签决定:latest 为 Always,其他标签为 IfNotPresent可显式设置 Always、Never 或 IfNotPresent;CRD 未提供默认值。详见 镜像拉取策略。false
pullSecrets镜像拉取密钥,详情请参考:https://kubernetes.io/docs/concepts/containers/images/#specifying-imagepullsecrets-on-a-podfalse
masterCpumaster 的 CPU 限制,单位可以是 ’m’ 或无单位,详情请参考:https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-cpufalse
workerCpuworker 的 CPU 限制,单位可以是 ’m’ 或无单位,详情请参考:https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-cpufalse
masterMemorymaster 的内存限制,单位可以是 Ei、Pi、Ti、Gi、Mi、Ki 之一,详情请参考:https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-memoryfalse
workerMemoryworker 的内存限制,单位可以是 Ei、Pi、Ti、Gi、Mi、Ki 之一,详情请参考:https://kubernetes.io/docs/concepts/configuration/manage-resources-containers/#meaning-of-memoryfalse
log4jXmlcomputer 作业的 log4j.xml 内容。false
jarFile容器镜像内算法 JAR 的路径。false
remoteJarUri算法 JAR 的 HTTP(S) 地址;启动脚本会下载该 JAR 并加载。false
jvmOptionscomputer 作业的 Java 启动参数。false
envVars请参考:https://kubernetes.io/docs/tasks/inject-data-application/define-interdependent-environment-variables/false
envFrom请参考:https://kubernetes.io/docs/tasks/inject-data-application/define-environment-variable-container/false
masterCommandbin/start-computer.shmaster 的运行命令,等同于 Docker 的 ‘Entrypoint’ 字段。false
masterArgs["-r master", “-d k8s”]master 的运行参数,等同于 Docker 的 ‘Cmd’ 字段。false
workerCommandbin/start-computer.shworker 的运行命令,等同于 Docker 的 ‘Entrypoint’ 字段。false
workerArgs["-r worker", “-d k8s”]worker 的运行参数,等同于 Docker 的 ‘Cmd’ 字段。false
volumes请参考:https://kubernetes.io/docs/concepts/storage/volumes/false
volumeMounts请参考:https://kubernetes.io/docs/concepts/storage/volumes/false
secretPathsk8s-secret 名称和挂载路径的映射。false
configMapPathsk8s-configmap 名称和挂载路径的映射。false
podTemplateSpec请参考:https://kubernetes.io/docs/reference/kubernetes-api/workload-resources/pod-template-v1/#PodTemplateSpecfalse
securityContext请参考:https://kubernetes.io/docs/tasks/configure-pod-container/security-context/false

KubeDriver 配置选项

配置项默认值说明
k8s.build_image_bash_path用于构建镜像的命令路径。
k8s.enable_internal_algorithmtrue是否启用内部算法。
k8s.framework_image_urlhugegraph/hugegraph-computer:latestcomputer 框架的镜像 URL。
k8s.image_repository_password登录镜像仓库的密码。
k8s.image_repository_registry登录镜像仓库的地址。
k8s.image_repository_urlhugegraph/hugegraph-computer镜像仓库的 URL。
k8s.image_repository_username登录镜像仓库的用户名。
k8s.internal_algorithm[pageRank]所有内部算法的名称列表。注意:算法名称在这里使用驼峰命名法(例如 pageRank),但算法实现返回下划线命名法(例如 page_rank)。
k8s.internal_algorithm_image_urlhugegraph/hugegraph-computer:latest内部算法的镜像 URL。
k8s.jar_file_dir/cache/jars/算法 jar 将上传到的目录。
k8s.kube_config~/.kube/configk8s 配置文件的路径。
k8s.log4j_xml_pathcomputer 作业的 log4j.xml 路径。
k8s.namespacehugegraph-computer-operator-systemhugegraph-computer 系统的命名空间。
k8s.pull_secret_names[]拉取镜像的 pull-secret 名称。