JanusGraph并发Gremlin查询机制及随机游走性能问题排查
目标
明确JanusGraph处理并发Gremlin查询的运行机制:查询为串行执行还是并行执行,调度判定逻辑是怎样的?
核心需求为执行大量计算与图遍历操作,所有服务均部署在本地机器,虽已对gremlinpython脚本做并行化处理,但执行过程中存在明显性能瓶颈。
环境配置
- JanusGraph 0.6.1 完整版
- 本地图实例,使用默认
conf/remote.yaml配置文件
实现逻辑
每个工作线程均配置独立属性,所有线程均持有AnonymousTraversalSource实例;线程从起始顶点列表中取出元素后循环执行下述逻辑,直至起始顶点列表为空:
def job(vertex_id:int, g:AnonymousTraversalSource, length:int, nb_walks:int) -> str: random_walks = [] for _ in range(nb_walks): random_walk = g.V(vertex_id).repeat( __.local(__.both().sample(1)) ).times(length).path().next() random_walks.append(",".join([str(v.id) for v in random_walk])) return "\n".join(random_walks)
遍历源初始化逻辑如下:
connection = DriverRemoteConnection(<URL>, "g") g = traversal().with_remote(connection)
线程类定义逻辑如下:
class myThread(threading.Thread): def __init__(self, thread_id, g, length, nb_walks): threading.Thread.__init__(self) self.thread_id = thread_id self.thread_count = 0 self.gtraversal = g self.walk_length = length self.nb_walks = nb_walks def run(self): while True: start_ids_list_lock.acquire() try: start_id = start_ids_list.pop(0) start_ids_list_lock.release() except IndexError: start_ids_list_lock.release() break else: self.thread_count += 1 random_walk = job( vertex_id=start_id, g=self.gtraversal, length=self.walk_length, nb_walks=self.nb_walks ) random_walks_list_lock.acquire() random_walks_list.append(random_walk) random_walks_list_lock.release()
已尝试方案
已测试如下三种配置方案:
- 所有线程传入同一个
AnonymousTraversalSource对象 - 基于同一个
DriverRemoteConnection对象为不同线程实例化独立的AnonymousTraversalSource对象 - 为不同线程创建独立的
DriverRemoteConnection对象,并基于各自连接构建独立的AnonymousTraversalSource对象
三种方案的性能无明显差异,执行约500次随机游走耗时均在20-25秒区间。
待排查问题
当前DriverRemoteConnection或AnonymousTraversalSource对象的构建方式是否存在问题?
是否存在可行的性能优化方案?当前实现方式是否已达到性能上限?
1. JanusGraph并发查询调度逻辑
JanusGraph本身没有全局串行执行查询的限制,查询是否并行完全由上层Gremlin Server线程池配置、连接参数和存储后端能力决定:
- 默认配置下Gremlin Server使用固定大小工作线程池处理请求,线程数默认等于CPU核心数,只要并发请求数不超过线程池上限,多个查询会被并行调度,不存在全局锁。
- 你测试三种连接方式性能无差异的核心原因是:gremlinpython默认连接配置中,单连接同时只能处理1个在途请求,且默认连接池仅初始化1个连接,不管你在Python侧创建多少线程、多少个TraversalSource实例,所有请求最终都会排队通过这个连接发往服务端,服务端按请求接收顺序串行响应,这就是客户端多线程没有带来性能提升的根本原因,和Python侧TraversalSource、DriverRemoteConnection的构建方式无关。
2. 现有实现的额外性能损耗点
除了连接配置导致的请求串行问题,你的遍历写法本身存在大量冗余开销:
- 你在Python侧循环调用
.next(),每一次随机游走都对应一次独立RPC往返,网络IO开销占总耗时的60%以上。 - 随机游走逻辑完全可以下推到Gremlin服务端执行,不需要每走一步就和客户端做一次网络交互。
3. 可直接落地的优化方案
第一步:调整连接参数,启用并发请求能力
给DriverRemoteConnection传入自定义连接池配置,放开单连接并发限制、调整连接池大小:
conn_pool_conf = { "minConnectionPoolSize": 8, "maxConnectionPoolSize": 16, "maxInProcessPerConnection": 8, "maxSimultaneousUsagePerConnection": 8 } connection = DriverRemoteConnection(<URL>, "g", **conn_pool_conf) g = traversal().with_remote(connection)
同步修改服务端conf/gremlin-server.yaml配置,将threadPoolWorker参数设置为CPU核心数的1-2倍,保证服务端有足够线程处理并行请求。
第二步:将循环逻辑下推到服务端,减少RPC次数
把nb_walks次随机游走的循环逻辑放到Gremlin遍历中,一次请求返回单个顶点对应的所有游走路径,把N次RPC压缩为1次RPC:
def job(vertex_id:int, g:AnonymousTraversalSource, length:int, nb_walks:int) -> str: # 一次请求拉取当前顶点所有随机游走结果,无需客户端循环发请求 random_walks = g.V(vertex_id).repeat( __.local(__.both().sample(1)) ).times(length).path().id().fold().repeat( __.local( __.V(vertex_id).repeat(__.local(__.both().sample(1))).times(length).path().id() ) ).times(nb_walks-1).fold().next() return "\n".join([",".join([str(vid) for vid in walk]) for walk in random_walks])
如果调用JanusGraph内置的专用随机游走步骤(JanusGraphRandomWalkStep),遍历性能还能再提升30%以上,不需要自行编写repeat组合逻辑。
第三步:匹配存储后端配置
如果你使用默认内嵌的BerkeleyJE作为存储后端,它本身仅支持单写多读,并发写入性能上限极低,本地测试时可以开启cache.db-cache参数配置数据库级缓存,减少磁盘IO开销;如果需要更高并发能力,建议替换为HBase或Cassandra作为存储后端,并发性能会有数量级提升。
按上述方案调整后,500次随机游走的耗时通常可以降到2-3秒区间,远低于当前20-25秒的水平,现有实现远没有达到性能上限。
内容的提问来源于stack exchange,提问作者GregoirePelegrin

