Neo4j中避免死锁与锁竞争的大规模并行导入

Lobsters Hottest 论文

摘要

本文提出一种改进的Neo4j大规模并行导入方法,采用基于哈希的分区和k-1着色算法,以避免数据加载过程中的死锁和锁竞争。

<p><a href="https://lobste.rs/s/iyouff/massive_parallel_imports_neo4j_without">评论</a></p>
查看原文
查看缓存全文

缓存时间: 2026/09/24 21:05

# Neo4j中大规模并行导入:避免死锁与锁竞争 来源:https://medium.com/neo4j/massive-parallel-imports-in-neo4j-without-deadlock-and-lock-contention-2c003a48d49a ## 基于工作线程与K-1着色算法优化图导入分区 作者:Maxime GUERY (https://medium.com/@maxime-guery?source=post_page---byline--2c003a48d49a-----------------------------------------) 本文是对Eric MONK (https://www.linkedin.com/in/dericmonk)在其文章"混合与批处理:Neo4j中快速并行关系加载的技术(https://neo4j.com/blog/developer/mix-and-batch-relationship-load/)"中设计方法的扩展。文中阐述了其理论基础与核心思想。后续文章将展示实际应用并提供源代码仓库。 他的方法利用源节点与目标节点ID的末位数字构建不重叠分区,然后将生成的矩阵单元格沿循环对角线分组为顺序批次,这些批次可以并行加载而不会同时访问相同节点集。**该方法有两个主要局限性**: 1. 根据节点ID末位数字划分分区无法使分区数适应可用工作线程数; 2. 当源节点与目标节点并非不相交时(例如在Person节点间创建关系),对角线批处理不安全。 因此,我们提出以下改进方案以克服上述限制: - 使用哈希函数基于工作线程数计算分区; - 新增针对源节点与目标节点来自同一集合(即非不相交)的关系导入函数; - 采用K-1着色算法计算可并行导入的分区批次。 ## 计算分区 > 通常,数据批量加载时会为节点和关系分别准备表。此步骤为这些表计算新列(称为`export_part`)。与原始文章不同,我们可以使用Mathieu EPLENIER (https://fr.linkedin.com/in/mathieu-eplenier)提出的优雅方法:通过以下公式计算源表的分区。 为实现最优导入,`partition_count`应等于服务器分配给数据导入的工作线程数。这提供了足够的独立任务来填充工作线程池,而无需在源数据中硬编码分区边界。 ### 节点分区 ``` export_part = HASH(id) modulo partition_count ``` Snowflake中的SQL表达式为: ``` CAST( ABS(HASH(node_id)) % partition_count AS VARCHAR) AS export_part ``` 若`partition_count = 5`,节点表最多包含`5`个分区值(`0`至`4`)。 ### 关系分区 ``` source_partition = HASH(id_source) modulo partition_count target_partition = HASH(id_target) modulo partition_count export_part = source_partition + " - " + target_partition ``` Snowflake中的SQL表达式为: ``` CAST( ABS(HASH(id_source)) % partition_count AS VARCHAR)|| ' - ' ||CAST( ABS(HASH(id_target)) % partition_count AS VARCHAR) AS export_part ``` 若`partition_count = 5`,关系表最多包含`5²=25`个分区值(从`0-0`到`4-4`)。 此方法具有确定性,可确保分区不相交。**但需注意:这无法保证分区大小相同,因为实际大小取决于数据结构。** > 此时需回答关键问题:**哪些分区可以并行导入?** ## 节点导入 节点表的导入较为简单:每个工作线程加载一个分区。因此所有分区可同时导入。 ## 关系导入 关系表的导入则更复杂,因为这取决于源节点与目标节点是**不相交**还是**非不相交**。为此,我们需要构建一个大小等于分区数(本例为5)的方阵。**这将帮助我们确定哪些分区可并行导入,从而最大化工作线程利用率。** ### 不相交节点间的导入 在"混合与批处理:Neo4j中快速并行关系加载技术"中使用的***对角线方法***,**在源节点与目标节点不属于同一集合时效果良好**。计算分区后,我们可以创建一个矩阵,每个单元格包含一个分区。矩阵的每条对角线对应一批可并行加载的分区。 *每种颜色代表一个批次,其中的分区可并行导入且不会产生死锁与锁竞争。* 例如,D3分区可并行导入且不会产生死锁或锁竞争,如下图所示: **Python代码** ```python def diagonals(partition_count: int) -> list[list[str]]: """根据源-目标分区矩阵返回并行批次(循环对角线)。 源节点与目标节点集合不相交,每个集合划分为``partition_count``个分区。 其组合形成包含``partition_count ** 2``个源-目标分区对的方阵。 函数将该矩阵划分为``partition_count``条循环对角线。 每条循环对角线代表一个可并行处理的分区批次。 批次内,每个源分区和目标分区恰好出现一次,防止并发任务访问相同分区。 参数: partition_count: 每个节点集合的分区数 返回: 并行批次列表,每个批次代表一条循环对角线, 包含格式为``"source - target"``的分区对 """ return [ [ f"{(diagonal_index + target_partition) % partition_count}" f" - {target_partition}" for target_partition in range(partition_count) ] for diagonal_index in range(partition_count) ] ``` ### 非不相交节点间的导入 **但若需导入来自同一集合节点间的关系,则会产生死锁错误。**因此我们需要另一个函数(如下定义)通过轮询旋转处理这种情况。 *此处每种颜色同样代表一个可并行导入且不会产生死锁或锁竞争的分区批次。****唯一区别在于每批次的分区数更少,这意味着比对角线函数生成的批次处理步骤更少。*** 例如,B5分区可并行导入且不会产生死锁或锁竞争,如下图所示: **Python代码** ```python def relationship_batches(partition_count: int) -> list[list[str]]: """返回非相交分区矩阵的并行批次。 源节点与目标节点属于同一节点集,该集合划分为``partition_count``个分区。 其组合形成包含``partition_count ** 2``个源-目标分区对的方阵。 每批次包含不共享任何分区的对。因此同一批次内的所有对均可并行处理, 而不会并发访问相同节点分区。批次通过轮询旋转生成。 正向对与反向对被置于不同批次,因为它们访问相同分区。 参数: partition_count: 节点集中的分区数 返回: 并行批次列表,包含格式为``"source - target"``的分区对 """ # 自引用对使用不同分区,因此可一起处理 batches = [ [ f"{partition} - {partition}" for partition in range(partition_count) ] ] partitions = list(range(partition_count)) # 轮询配对需要偶数个值。对于奇数个分区, # None代表轮次中休息的分区 if partition_count % 2: partitions.append(None) # 固定一个分区,围绕其旋转所有其他分区。 # 固定分区作为锚点,防止旋转产生重复对。 fixed_partition = partitions[0] rotating_partitions = partitions[1:] # 固定一个值并旋转剩余值可生成所有可能的无序分区对。 for _ in range(len(partitions) - 1): current_partitions = [ fixed_partition, *rotating_partitions # 展开旋转分区形成当前轮次对 ] # 配对当前轮次中位置相反的值。 # 每个实际分区最多出现在一个对中。 pairs = [ (current_partitions[index], current_partitions[-1 - index]) for index in range(len(current_partitions) // 2) if current_partitions[index] is not None and current_partitions[-1 - index] is not None ] # 正向对不共享任何物理分区,因此可并行处理。 batches.append([ f"{source_partition} - {target_partition}" for source_partition, target_partition in pairs ]) # 反向对必须置于单独批次,因为它们使用与其对应正向对相同的物理分区。 batches.append([ f"{target_partition} - {source_partition}" for source_partition, target_partition in pairs ]) # 旋转除固定锚点外的所有分区。最后一个旋转分区移至下一轮次开头。 rotating_partitions = [ rotating_partitions[-1], *rotating_partitions[:-1], ] return batches ``` **因此,使用这些函数计算的分区进行批量导入关系,有助于避免死锁与锁竞争。** ## 基于K-1着色算法的通用函数 此前我们使用两个Python函数确定哪些分区可并行导入。它们旨在从矩阵(因为关系在两个节点间创建)计算最优分区批次。但若需从更复杂的图(如单纯复形(https://en.wikipedia.org/wiki/Simplicial_complex))计算分区批次呢?这是一种高级图结构,其中关系可连接x个节点。 单纯复形(来源Wikipedia) 这正是可以使用Neo4j图数据科学库(https://neo4j.com/docs/graph-data-science/current/)中K-1着色算法(https://neo4j.com/docs/graph-data-science/current/algorithms/k1coloring/)计算分区批次的场景。我们将每个源-目标分区对建模为节点,当两个节点的关系加载至少一个共同节点分区时则连接它们。K-1着色为相邻节点分配不同颜色,使每个颜色组形成无冲突的并行批次(前提是算法已收敛),而不同颜色组则顺序处理。 > 由于算法非确定性,你能获得可行解,但并非最优解(在我们场景中指速度层面)。为说明这点,以下Cypher查询创建前述分区矩阵,并应用K-1着色算法计算分区批次。 ### 清理查询 ``` // 删除关系 MATCH ()-[r:SHARES_PARTITION_NOT_DISJOINT|SHARES_PARTITION_DISJOINT]->() CALL (r) { DELETE r } IN TRANSACTIONS OF 1000 ROWS FINISH; // 删除节点 MATCH (n:Partition) CALL (n) { DELETE n } IN TRANSACTIONS OF 1000 ROWS FINISH; ``` ### 图创建 ``` // 设置网格参数 :param grid => 5; // 创建节点网格的查询 WITH $grid AS grid UNWIND range(0, grid - 1) AS sourcePartition UNWIND range(0, grid - 1) AS targetPartition CALL (grid, sourcePartition, targetPartition) { MERGE (p:Partition { grid: grid, id: toString(sourcePartition) + "-" + toString(targetPartition) }) SET p.sourcePartition = sourcePartition, p.targetPartition = targetPartition } FINISH; // 创建节点间关系的查询 MATCH (a:Partition {grid: $grid}) MATCH (b:Partition {grid: $grid}) WHERE a < b CALL (a, b) { // 当节点非不相交时(源节点与目标节点属于同一集合) WITH a, b WHERE a.sourcePartition = b.sourcePartition OR a.sourcePartition = b.targetPartition OR a.targetPartition = b.sourcePartition OR a.targetPartition = b.targetPartition MERGE (a)-[r:SHARES_PARTITION_NOT_DISJOINT]->(b) UNION // 当节点不相交时(源节点与目标节点属于不同集合) WITH a, b WHERE a.sourcePartition = b.sourcePartition OR a.targetPartition = b.targetPartition MERGE (a)-[r:SHARES_PARTITION_DISJOINT]->(b) } FINISH; ``` ### 运行K-1着色 ``` // GDS投影 CYPHER runtime=parallel // SHARES_PARTITION_NOT_DISJOINT 或 SHARES_PARTITION_DISJOINT MATCH (source)-[r:SHARES_PARTITION_NOT_DISJOINT]->(target) RETURN gds.graph.project( 'grid', source, target, {}, { undirectedRelationshipTypes: ['*'] } ) AS graph; // 运行K-1着色算法(流模式) CALL gds.k1coloring.stream('grid', {maxIterations: 100, concurrency: 4}) YIELD nodeId, color RETURN color, collect(gds.util.asNode(nodeId).id) AS partitions; // 删除内存图投影 CALL gds.graph.drop('grid', false) YIELD graphName RETURN graphName; ``` 节点**不相交**时生成的批次(表格与图视图): ``` ╒═════╤═══════════════════════════════════╕ │color│partitions │ ╞═════╪═══════════════════════════════════╡ │0 │["0-0", "1-1", "2-2", "3-3", "4-4"]│ ├─────┼───────────────────────────────────┤ │1 │["0-1", "1-0", "3-2", "2-3"] │ ├─────┼───────────────────────────────────┤ │2 │["0-2", "2-0", "3-1", "1-3"] │ ├─────┼───────────────────────────────────┤ │3 │["0-3", "3-0", "2-1", "1-2"] │ ├─────┼───────────────────────────────────┤ │4 │["0-4", "4-0"] │ ├─────┼───────────────────────────────────┤ │5 │["4-1", "1-4"] │ ├─────┼───────────────────────────────────┤ │6 │["4-2", "2-4"] │ ├─────┼───────────────────────────────────┤ │7 │["4-3", "3-4"] │ └─────┴───────────────────────────────────┘ ``` 节点**非不相交**时生成的批次(表格与图视图): ``` ╒═════╤═══════════════════════════════════╕ │color│partitions │ ╞═════╪═══════════════════════════════════╡ │0 │["0-0", "1-1", "2-2", "3-3", "4-4"]│ ├─────┼───────────────────────────────────┤ │1 │["0-1", "2-3"] │ ├─────┼───────────────────────────────────┤ │2 │["0-2", "1-3"] │ ├─────┼───────────────────────────────────┤ │3 │["0-3", "1-2"] │ ├─────┼───────────────────────────────────┤ │4 │["0-4", "2-1"] │ ├─────┼───────────────────────────────────┤ │5 │["1-0", "2-4"] │ ├─────┼───────────────────────────────────┤ │6 │["2-0", "1-4"] │ ├─────┼───────────────────────────────────┤ │7 │["3-0", "4-1"] │ ├─────┼───────────────────────────────────┤ │8 │["4-0", "3-1"] │ ├─────┼───────────────────────────────────┤ │9 │["3-2"] │ ├─────┼───────────────────────────────────┤ │10 │["4-2"] │ ├─────┼───────────────────────────────────┤ │11 │["3-4"] │ ├─────┼───────────────────────────────────┤ │12 │["4-3"] │ └─────┴───────────────────────────────────┘ ``` 如图所示,K-1着色提供的解决方案不如Python函数最优。**但在复杂并行化场景中,这提供了一种有效替代方案。** 每种颜色代表一个可并行导入的分区批次。例如,使用颜色2时,一个工作线程将导入分区为`0-2`的关系,而另一个工作线程将导入分区为`4-4`的关系。 ## 结论 本文提出了一种**通用且优化**的方法...

相似文章

大规模并行 Postgres 备份

Hacker News Top

PlanetScale 描述了它如何通过为每个分片启动 EC2 实例、从对象存储恢复之前的备份并重放 WAL,来对分片 Postgres 数据库执行大规模并行备份,从而实现超过 50 GB/s 的 PB 级备份速度。

并行折叠

Lobsters Hottest

探索使用幺半群进行易并行数据处理,表明霍纳规则和Boyer-Moore多数投票算法等传统串行算法可通过幺半群组合实现并行化。同时介绍了垂直幺半群组合,用于高效的嵌套分组聚合。

让Postgres队列实现可扩展性

Hacker News Top

一篇详细的技术博文,解释如何使用SKIP LOCKED和适当的事务隔离级别来扩展基于PostgreSQL的队列,实现每秒3万次工作流执行。