From 11fd446c979db48c3c2b2299e70326df815e4083 Mon Sep 17 00:00:00 2001 From: berrazuriz1 Date: Mon, 25 May 2026 16:17:52 -0400 Subject: [PATCH] Batch node creation to avoid oversized Bolt transactions --- .../graph_db_manager/neo4j_manager.py | 26 ++++++++++++------- 1 file changed, 17 insertions(+), 9 deletions(-) diff --git a/blarify/repositories/graph_db_manager/neo4j_manager.py b/blarify/repositories/graph_db_manager/neo4j_manager.py index c4d3dfe0..9aee1a31 100644 --- a/blarify/repositories/graph_db_manager/neo4j_manager.py +++ b/blarify/repositories/graph_db_manager/neo4j_manager.py @@ -85,19 +85,27 @@ def save_graph(self, nodes: List[Any], edges: List[Any]): self.create_edges(edges) def create_nodes(self, nodeList: List[Any]): - # Function to create nodes in the Neo4j database if self.repo_id is None: raise ValueError("repo_id is required for creating nodes. Cannot create nodes with entity-wide scope.") + batch_size: int = 5000 + total_nodes: int = len(nodeList) + total_batches: int = (total_nodes + batch_size - 1) // batch_size + logger.info(f"Creating {total_nodes} nodes in batches of {batch_size}") + with self.driver.session(database=self.database) as session: - session.execute_write( - self._create_nodes_txn, - nodeList, - 1000, - repoId=self.repo_id, - entityId=self.entity_id, - environment=self.environment.value, - ) + for i in range(0, total_nodes, batch_size): + batch: List[Any] = nodeList[i:i + batch_size] + batch_num: int = i // batch_size + 1 + logger.info(f"Processing nodes batch {batch_num}/{total_batches} ({i}/{total_nodes})") + session.execute_write( + self._create_nodes_txn, + batch, + 1000, + repoId=self.repo_id, + entityId=self.entity_id, + environment=self.environment.value, + ) def create_edges(self, edgesList: List[Any]): # Function to create edges between nodes in the Neo4j database