diff --git a/arango_rdf/controller.py b/arango_rdf/controller.py index 55cf7aa..9f52f79 100644 --- a/arango_rdf/controller.py +++ b/arango_rdf/controller.py @@ -1,7 +1,7 @@ #!/usr/bin/env python3 from typing import Set -from arango.database import StandardDatabase +from arango.database import AsyncDatabase, StandardDatabase from rdflib import Graph from .abc import AbstractArangoRDFController @@ -27,6 +27,7 @@ class ArangoRDFController(AbstractArangoRDFController): def __init__(self) -> None: self.db: StandardDatabase + self.async_db: AsyncDatabase self.rdf_graph: Graph def identify_best_class( diff --git a/arango_rdf/main.py b/arango_rdf/main.py index 270d317..c16bb87 100644 --- a/arango_rdf/main.py +++ b/arango_rdf/main.py @@ -12,7 +12,7 @@ import farmhash from arango.collection import StandardCollection from arango.cursor import Cursor -from arango.database import StandardDatabase +from arango.database import AsyncDatabase, StandardDatabase from arango.graph import Graph as ADBGraph from isodate import Duration from rdflib import RDF, RDFS, XSD, BNode @@ -68,6 +68,13 @@ class ArangoRDF(AbstractArangoRDF): that using an underscore "_", results in these attributes being treated as ArangoDB system attributes. Using "$" is an alternative non-system prefix. :type rdf_attribute_prefix: str + :param insert_async: If True, will insert documents asynchronously. + Defaults to False. + :type insert_async: bool + :param enable_pgt_cache: If True, will enable the PGT term metadata cache to avoid + repeated computations. Defaults to False. Not always useful, especially when + terms are not repeated alot in the RDF graph. + :type enable_pgt_cache: bool :raise TypeError: On invalid parameter types """ @@ -77,6 +84,8 @@ def __init__( controller: ArangoRDFController = ArangoRDFController(), logging_lvl: Union[str, int] = logging.INFO, rdf_attribute_prefix: str = "_", + insert_async: bool = False, + enable_pgt_cache: bool = False, ): self.set_logging(logging_lvl) @@ -88,9 +97,13 @@ def __init__( msg = "**controller** parameter must inherit from ArangoRDFController" raise TypeError(msg) - self.__db = db - self.__cntrl = controller - self.__cntrl.db = db + self.db: StandardDatabase = db + self.async_db: AsyncDatabase = db.begin_async_execution(return_result=False) + self.insert_async = insert_async + + self.controller: ArangoRDFController = controller + self.controller.db = db + self.controller.async_db = self.async_db # Set the RDF attribute prefix self.__rdf_attribute_prefix = rdf_attribute_prefix @@ -104,9 +117,10 @@ def __init__( self.__rdf_lang_attr = f"{rdf_attribute_prefix}lang" self.__rdf_datatype_attr = f"{rdf_attribute_prefix}datatype" - # An RDF to ArangoDB variable used as a buffer - # to store the to-be-inserted ArangoDB documents (RDF-to-ArangoDB). - self.__adb_docs: ADBDocs + # Pre-compile regex patterns for container predicates + rdf_ns = str(RDF) + self.__container_pattern_n = re.compile(f"^{re.escape(rdf_ns)}_[0-9]+$") + self.__container_pattern_li = re.compile(f"^{re.escape(rdf_ns)}li$") # Work-in-progress feature to enhance the Terminology Box of an RDF Graph # when importing to ArangoDB. @@ -125,6 +139,10 @@ def __init__( # e.g ( "4502") self.adb_key_uri = URIRef("http://www.arangodb.com/key") + # Cache for PGT term metadata to avoid repeated computations + self.enable_pgt_cache = enable_pgt_cache + self.pgt_term_metadata_cache: Dict[str, RDFTermMeta] = {} + # RDF Graph for maintaining the ArangoDB Collections & Keys # of the RDF Resources self.__adb_col_statements = RDFGraph() @@ -163,14 +181,6 @@ def __init__( logger.info(f"Instantiated ArangoRDF with database '{db.name}'") - @property - def db(self) -> StandardDatabase: - return self.__db # pragma: no cover - - @property - def controller(self) -> ArangoRDFController: - return self.__cntrl # pragma: no cover - @property def rdf_attribute_prefix(self) -> str: return self.__rdf_attribute_prefix # pragma: no cover @@ -695,8 +705,10 @@ def rdf_to_arangodb_by_rpt( self.__rdf_graph = rdf_graph self.__adb_key_statements = self.extract_adb_key_statements(rdf_graph) + # Create the ArangoDB documents buffer for this transformation + adb_docs: ADBDocs = defaultdict(lambda: defaultdict(dict)) + # Reset the ArangoDB Config - self.__adb_docs = defaultdict(lambda: defaultdict(dict)) self.__contextualize_graph = contextualize_graph self.__use_hashed_literals_as_keys = use_hashed_literals_as_keys @@ -724,7 +736,16 @@ def rdf_to_arangodb_by_rpt( # NOTE: Graph Contextualization is an experimental work-in-progress contextualize_statement_func = empty_func if contextualize_graph: - contextualize_statement_func = self.__rpt_contextualize_statement + + def contextualize_statement_func( + s_meta: RDFTermMeta, + p_meta: RDFTermMeta, + o_meta: RDFTermMeta, + sg_str: str, + ) -> None: + return self.__rpt_contextualize_statement( + adb_docs, s_meta, p_meta, o_meta, sg_str + ) self.__rdf_graph = self.__load_meta_ontology(self.__rdf_graph) @@ -744,6 +765,7 @@ def rdf_to_arangodb_by_rpt( self.__reified_subject_map = {} if flatten_reified_triples: self.__flatten_reified_triples( + adb_docs, self.__rpt_process_subject_predicate_object, contextualize_statement_func, batch_size, @@ -758,10 +780,10 @@ def rdf_to_arangodb_by_rpt( p: URIRef # Predicate o: RDFTerm # Object - rdf_graph_size = len(self.__rdf_graph) - batch_size = batch_size or rdf_graph_size + total = len(self.__rdf_graph) + batch_size = batch_size or total bar_progress = get_bar_progress("(RDF → ADB): RPT", "#BF23C4") - bar_progress_task = bar_progress.add_task("", total=rdf_graph_size) + bar_progress_task = bar_progress.add_task("", total=total) spinner_progress = get_import_spinner_progress(" ") statements = ( @@ -772,18 +794,21 @@ def rdf_to_arangodb_by_rpt( with Live(Group(bar_progress, spinner_progress)): for i, (s, p, o, *sg) in enumerate(statements((None, None, None)), 1): - bar_progress.advance(bar_progress_task) - logger.debug(f"RPT: {s} {p} {o} {sg}") self.__rpt_process_subject_predicate_object( - s, p, o, sg, None, contextualize_statement_func + adb_docs, s, p, o, sg, None, contextualize_statement_func ) if i % batch_size == 0: - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + bar_progress.update(bar_progress_task, advance=batch_size) + self.__insert_adb_docs( + adb_docs, spinner_progress, **adb_import_kwargs + ) - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + last_advance = total % batch_size if batch_size > 0 else 0 + bar_progress.update(bar_progress_task, advance=last_advance) + self.__insert_adb_docs(adb_docs, spinner_progress, **adb_import_kwargs) return self.__rpt_create_adb_graph(name) @@ -930,35 +955,38 @@ def rdf_to_arangodb_by_pgt( m = "Cannot specify both **uri_map_collection_name** and **resource_collection_name**." # noqa: E501 raise ValueError(m) - if not self.__db.has_collection(uri_map_collection_name): - self.__db.create_collection(uri_map_collection_name) + if not self.db.has_collection(uri_map_collection_name): + self.db.create_collection(uri_map_collection_name) - self.__uri_map_collection = self.__db.collection(uri_map_collection_name) + self.__uri_map_collection = self.db.collection(uri_map_collection_name) self.__resource_collection = None if resource_collection_name: - if not self.__db.has_collection(resource_collection_name): - self.__db.create_collection(resource_collection_name) + if not self.db.has_collection(resource_collection_name): + self.db.create_collection(resource_collection_name) - self.__resource_collection = self.__db.collection(resource_collection_name) + self.__resource_collection = self.db.collection(resource_collection_name) self.__predicate_collection = None if predicate_collection_name: - if not self.__db.has_collection(predicate_collection_name): - self.__db.create_collection(predicate_collection_name, edge=True) + if not self.db.has_collection(predicate_collection_name): + self.db.create_collection(predicate_collection_name, edge=True) - self.__predicate_collection = self.__db.collection( - predicate_collection_name - ) + self.__predicate_collection = self.db.collection(predicate_collection_name) + + # Create the ArangoDB documents buffer for this transformation + adb_docs: ADBDocs = defaultdict(lambda: defaultdict(dict)) # Reset the ArangoDB Config - self.__adb_docs = defaultdict(lambda: defaultdict(dict)) self.__contextualize_graph = contextualize_graph # A unique set of instance variables to # convert RDF Lists into JSON Lists during the PGT Process self.__rdf_list_heads: RDFListHeads = defaultdict(lambda: defaultdict(dict)) self.__rdf_list_data: RDFListData = defaultdict(lambda: defaultdict(dict)) + self.__rdf_list_subjects: Set[RDFTerm] = set() + self.__rdf_collection_subjects: Set[RDFTerm] = set() + self.__rdf_container_subjects: Set[RDFTerm] = set() # The ArangoDB Collection name of all unidentified RDF Resources self.__UNKNOWN_RESOURCE = f"{name}_UnknownResource" @@ -975,7 +1003,16 @@ def rdf_to_arangodb_by_pgt( # NOTE: Graph Contextualization is an experimental work-in-progress contextualize_statement_func = empty_func if contextualize_graph: - contextualize_statement_func = self.__pgt_contextualize_statement + + def contextualize_statement_func( + s_meta: RDFTermMeta, + p_meta: RDFTermMeta, + o_meta: RDFTermMeta, + sg_str: str, + ) -> None: + return self.__pgt_contextualize_statement( + adb_docs, s_meta, p_meta, o_meta, sg_str + ) self.__rdf_graph = self.__load_meta_ontology(self.__rdf_graph) @@ -1029,35 +1066,22 @@ def rdf_to_arangodb_by_pgt( self.__reified_subject_map = {} if flatten_reified_triples: self.__flatten_reified_triples( + adb_docs, self.__pgt_process_subject_predicate_object, contextualize_statement_func, batch_size, adb_import_kwargs, ) - ########################### - # PGT: Literal Statements # - ########################### + ############################## + # PGT: Pre-compute RDF lists # + ############################## - self.__pgt_parse_literal_statements( - contextualize_statement_func, - batch_size, - adb_import_kwargs, - ) - - ############# - # PGT: Main # - ############# - - s: RDFTerm # Subject - p: URIRef # Predicate - o: RDFTerm # Object + self.__precompute_rdf_list_info() - rdf_graph_size = len(self.__rdf_graph) - batch_size = batch_size or rdf_graph_size - bar_progress = get_bar_progress("(RDF → ADB): PGT", "#08479E") - bar_progress_task = bar_progress.add_task("", total=rdf_graph_size) - spinner_progress = get_import_spinner_progress(" ") + ########################### + # PGT: Prepare Statements # + ########################### statements = ( self.__rdf_graph.quads @@ -1065,54 +1089,68 @@ def rdf_to_arangodb_by_pgt( else self.__rdf_graph.triples ) - with Live(Group(bar_progress, spinner_progress)): - for i, (s, p, o, *sg) in enumerate(statements((None, None, None)), 1): - bar_progress.advance(bar_progress_task) + literal_statements = defaultdict(list) + non_literal_statements = defaultdict(list) - logger.debug(f"PGT: {s} {p} {o} {sg}") + with get_spinner_progress("(RDF → ADB): PGT [Prepare Statements]") as rp: + rp.add_task("") - # Address the possibility of (s, p, o) being a part of the - # structure of an RDF Collection or an RDF Container. - # TODO: Move out of loop, into a pre-processing step - rdf_list_col = self.__pgt_statement_is_part_of_rdf_list(s, p) - if rdf_list_col: - key = self.rdf_id_to_adb_label(str(p)) - doc = self.__rdf_list_data[rdf_list_col][s] - self.__pgt_rdf_val_to_adb_val(doc, key, o) - continue + for s, p, o, *sg in statements((None, None, None)): + if isinstance(o, Literal) and s not in self.__rdf_list_subjects: + literal_statements[(s, p)].append((o, sg)) + else: + non_literal_statements[(s, p)].append((o, sg)) - self.__pgt_process_subject_predicate_object( - s, p, o, sg, None, contextualize_statement_func - ) + ########################### + # PGT: Literal Statements # + ########################### - if i % batch_size == 0: - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + self.__pgt_parse_literal_statements( + adb_docs, + literal_statements, + contextualize_statement_func, + batch_size, + adb_import_kwargs, + ) + + ############################### + # PGT: Non-Literal Statements # + ############################### - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + self.__pgt_parse_non_literal_statements( + adb_docs, + non_literal_statements, + contextualize_statement_func, + batch_size, + adb_import_kwargs, + ) ################## # PGT: RDF Lists # ################## - bar_progress = get_bar_progress("(RDF → ADB): PGT [RDF Lists]", "#EF7D00") + bar_progress = get_bar_progress("(RDF → ADB): PGT [Lists]", "#EF7D00") + spinner_progress = get_import_spinner_progress(" ") with Live(Group(bar_progress, spinner_progress)): - self.__pgt_process_rdf_lists(bar_progress) - self.__insert_adb_docs(spinner_progress) + self.__pgt_process_rdf_lists(adb_docs, bar_progress) + self.__insert_adb_docs(adb_docs, spinner_progress, **adb_import_kwargs) ########################### # PGT: Namespace Prefixes # ########################### if namespace_collection_name: - if not self.__db.has_collection(namespace_collection_name): - self.__db.create_collection(namespace_collection_name) + if not self.db.has_collection(namespace_collection_name): + self.db.create_collection(namespace_collection_name) docs = [ {"prefix": prefix, "uri": uri, "_key": self.hash(uri)} for prefix, uri in namespace_prefixes ] - result = self.__db.collection(namespace_collection_name).insert_many( + db = self.db if self.insert_async else self.async_db + + result = db.collection(namespace_collection_name).insert_many( docs, overwrite=True, raise_on_document_error=True ) @@ -1120,6 +1158,46 @@ def rdf_to_arangodb_by_pgt( return self.__pgt_create_adb_graph(name) + def __precompute_rdf_list_info(self) -> None: + """Pre-compute RDF list information for optimization. + + Collects all RDF list subjects categorized by type once + to avoid repeated computation in processing loops. + """ + # Pre-compute collection subjects (RDF.first, RDF.rest) + for s in self.__rdf_graph.subjects(RDF.first, None): + self.__rdf_collection_subjects.add(s) + for s in self.__rdf_graph.subjects(RDF.rest, None): + self.__rdf_collection_subjects.add(s) + + # Pre-compute container subjects (container predicates _1, li, etc.) + for s, p, _ in self.__rdf_graph: + if isinstance(s, BNode): + p_str = str(p) + container_pattern_n = self.__container_pattern_n.match(p_str) + container_pattern_li = self.__container_pattern_li.match(p_str) + if container_pattern_n or container_pattern_li: + self.__rdf_container_subjects.add(s) + + self.__rdf_list_subjects = ( + self.__rdf_collection_subjects | self.__rdf_container_subjects + ) + + def __is_rdf_list_statement(self, s: RDFTerm, p: URIRef) -> str: + """Returns the list type or empty string if not a list statement. + + :param s: The RDF Subject + :param p: The RDF Predicate + :return: The list type or empty string if not a list statement + """ + if s in self.__rdf_collection_subjects and p in {RDF.first, RDF.rest}: + return "_COLLECTION_BNODE" + + if s in self.__rdf_container_subjects: + return "_CONTAINER_BNODE" + + return "" + def write_adb_col_statements( self, rdf_graph: RDFGraph, @@ -1160,18 +1238,18 @@ def write_adb_col_statements( self.__adb_col_statements.bind("adb", self.__adb_ns) self.__rdf_graph = rdf_graph - self.__cntrl.rdf_graph = rdf_graph + self.controller.rdf_graph = rdf_graph with get_spinner_progress("(RDF → ADB): Write Col Statements") as rp: rp.add_task("") # 0. Add URI Collection statements if uri_map_collection_name: - if not self.__db.has_collection(uri_map_collection_name): + if not self.db.has_collection(uri_map_collection_name): m = f"URI collection '{uri_map_collection_name}' does not exist" raise ValueError(m) - for doc in self.__db.collection(uri_map_collection_name): + for doc in self.db.collection(uri_map_collection_name): uri = URIRef(doc[self.__rdf_uri_attr]) collection = str(doc["collection"]) self.__add_adb_col_statement(uri, collection, True) @@ -1197,6 +1275,13 @@ def write_adb_col_statements( if self.__contextualize_graph: self.__type_map = self.__combine_type_map_and_dr_map() + # If the resource collection is not None, we don't need to run the + # ArangoDB Collection Mapping Process to completion, since we will + # be using the resource collection for all RDF Resources except for + # Class and Property. + if self.__resource_collection is not None: + return self.__adb_col_statements + # 5. Finalize **adb_col_statements** for rdf_map in [self.__explicit_type_map, self.__domain_range_map]: for rdf_resource, class_set in rdf_map.items(): @@ -1204,7 +1289,7 @@ def write_adb_col_statements( if t in self.__adb_col_statements or len(class_set) == 0: continue # pragma: no cover # (false negative) - best_class = self.__cntrl.identify_best_class( + best_class = self.controller.identify_best_class( rdf_resource, class_set, self.__subclass_tree ) @@ -1246,8 +1331,8 @@ def migrate_unknown_resources( """ ur_collection_name = f"{graph_name}_UnknownResource" - ur_collection = self.__db.collection(ur_collection_name) - uri_map_collection = self.__db.collection(uri_map_collection_name) + ur_collection = self.db.collection(ur_collection_name) + uri_map_collection = self.db.collection(uri_map_collection_name) if ur_collection.count() == 0: logger.info("No Unknown Resources to migrate") @@ -1292,9 +1377,7 @@ def migrate_unknown_resources( "graph": graph_name, } - cursor = self.__db.aql.execute( - query, bind_vars=bind_vars, stream=True, **kwargs - ) + cursor = self.db.aql.execute(query, bind_vars=bind_vars, stream=True, **kwargs) edge_count = 0 @@ -1311,7 +1394,7 @@ def migrate_unknown_resources( edges_to_modify = edge_data["edges_to_modify"] edge_count += len(edges_to_modify) - result = self.__db.collection(edge_collection).update_many( + result = self.db.collection(edge_collection).update_many( edges_to_modify, merge=True, raise_on_document_error=True, @@ -1319,9 +1402,7 @@ def migrate_unknown_resources( logger.debug(result) - self.__db.collection(collection).update( - data, merge=True, silent=True - ) + self.db.collection(collection).update(data, merge=True, silent=True) cursor.batch().clear() if cursor.has_more(): @@ -1632,12 +1713,12 @@ def __fetch_adb_docs( if ignored_attributes: aql_return_value = f"UNSET(doc, {list(ignored_attributes)})" - col_size: int = self.__db.collection(col).count() + col_size: int = self.db.collection(col).count() with get_spinner_progress(f"(ADB → RDF): Export '{col}' ({col_size})") as sp: sp.add_task("") - cursor: Cursor = self.__db.aql.execute( + cursor: Cursor = self.db.aql.execute( f"FOR doc IN @@col RETURN {aql_return_value}", bind_vars={"@col": col}, **{**adb_export_kwargs, **{"stream": True}}, @@ -1711,7 +1792,7 @@ def __process_adb_vertex( term = self.__adb_doc_to_rdf_term(adb_v, v_col) self.__term_map[adb_v["_id"]] = term - if type(term) is Literal: + if isinstance(term, Literal): return term sg = URIRef(adb_v.get(self.__rdf_sub_graph_uri_attr, "")) or None @@ -1982,7 +2063,7 @@ def __adb_val_to_rdf_val( :type sg: URIRef | None """ - if type(val) is list: + if isinstance(val, list): if self.__list_conversion == "static": for v in val: self.__adb_val_to_rdf_val(col, s, p, v, sg) @@ -2011,7 +2092,7 @@ def __adb_val_to_rdf_val( val = json.dumps(val) self.__add_to_rdf_graph(s, p, Literal(val), sg) - elif type(val) is dict: + elif isinstance(val, dict): if self.__dict_conversion == "static": bnode = BNode() self.__add_to_rdf_graph(s, p, bnode, sg) @@ -2077,6 +2158,7 @@ def extract_adb_key_statements( def __rpt_process_subject_predicate_object( self, + adb_docs: ADBDocs, s: RDFTerm, p: URIRef, o: RDFTerm, @@ -2087,6 +2169,8 @@ def __rpt_process_subject_predicate_object( """RDF -> ArangoDB (RPT): Processes the RDF Statement (s, p, o) as an ArangoDB document for RPT. + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param s: The RDF Subject of the RDF Statement. :type s: URIRef | BNode :param p: The RDF Predicate of the RDF Statement. @@ -2106,19 +2190,23 @@ def __rpt_process_subject_predicate_object( """ sg_str = self.__get_subgraph_str(sg) - s_meta = self.__rpt_process_term(s) + s_meta = self.__rpt_process_term(adb_docs, s) - o_meta = self.__rpt_process_term(o) + o_meta = self.__rpt_process_term(adb_docs, o) - self.__rpt_process_statement(s_meta, p, o_meta, sg_str, reified_subject) + self.__rpt_process_statement( + adb_docs, s_meta, p, o_meta, sg_str, reified_subject + ) contextualize_statement_func(s_meta, p, o_meta, sg_str) - def __rpt_process_term(self, t: RDFTerm) -> RDFTermMeta: + def __rpt_process_term(self, adb_docs: ADBDocs, t: RDFTerm) -> RDFTermMeta: """RDF -> ArangoDB (RPT): Process an RDF Term as an ArangoDB document via RPT Standards. Returns the ArangoDB Collection & Document Key associated to the RDF term, along with its string representation. + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param t: The RDF Term to process :type t: URIRef | BNode | Literal :return: The RDF Term object, along with its associated ArangoDB @@ -2136,46 +2224,44 @@ def __rpt_process_term(self, t: RDFTerm) -> RDFTermMeta: # TODO: Populate adb docs? Or uncessary? - elif type(t) is URIRef: + elif isinstance(t, URIRef): t_col = self.__URIREF_COL t_label = self.rdf_id_to_adb_label(t_str) - self.__adb_docs[t_col][t_key] = { + adb_docs[t_col][t_key] = { "_key": t_key, self.__rdf_uri_attr: t_str, self.__rdf_label_attr: t_label, self.__rdf_type_attr: "URIRef", } - elif type(t) is BNode: + elif isinstance(t, BNode): t_col = self.__BNODE_COL - self.__adb_docs[t_col][t_key] = { + adb_docs[t_col][t_key] = { "_key": t_key, self.__rdf_label_attr: "", self.__rdf_type_attr: "BNode", } - elif type(t) is Literal: + elif isinstance(t, Literal): t_col = self.__LITERAL_COL t_value = self.__get_literal_val(t, t_str) t_label = t_value - self.__adb_docs[t_col][t_key] = { + adb_docs[t_col][t_key] = { self.__rdf_value_attr: t_value, self.__rdf_label_attr: t_label, # TODO: REVISIT self.__rdf_type_attr: "Literal", } if self.__use_hashed_literals_as_keys: - self.__adb_docs[t_col][t_key]["_key"] = t_key + adb_docs[t_col][t_key]["_key"] = t_key if t.language: - self.__adb_docs[t_col][t_key][self.__rdf_lang_attr] = t.language + adb_docs[t_col][t_key][self.__rdf_lang_attr] = t.language elif t.datatype: - self.__adb_docs[t_col][t_key][self.__rdf_datatype_attr] = str( - t.datatype - ) + adb_docs[t_col][t_key][self.__rdf_datatype_attr] = str(t.datatype) else: raise ValueError(f"Unable to process {t}") # pragma: no cover @@ -2184,6 +2270,7 @@ def __rpt_process_term(self, t: RDFTerm) -> RDFTermMeta: def __rpt_process_statement( self, + adb_docs: ADBDocs, s_meta: RDFTermMeta, p: URIRef, o_meta: RDFTermMeta, @@ -2193,6 +2280,8 @@ def __rpt_process_statement( """RDF -> ArangoDB (RPT): Processes the RDF Statement (s, p, o) as an ArangoDB edge for RPT. + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param s_meta: The RDF Term Metadata associated to **s**. :type s_meta: arango_rdf.typings.RDFTermMeta :param p: The RDF Predicate URIRef of the statement (s, p, o). @@ -2224,6 +2313,7 @@ def __rpt_process_statement( e_key = self.hash(f"{s_key}-{p_key}-{o_key}") self.__add_adb_edge( + adb_docs, self.__STATEMENT_COL, e_key, _from, @@ -2234,10 +2324,17 @@ def __rpt_process_statement( ) def __rpt_contextualize_statement( - self, s_meta: RDFTermMeta, p: URIRef, o_meta: RDFTermMeta, sg_str: str + self, + adb_docs: ADBDocs, + s_meta: RDFTermMeta, + p: URIRef, + o_meta: RDFTermMeta, + sg_str: str, ) -> None: """RDF -> ArangoDB (RPT): Contextualizes the RDF Statement (s, p, o). + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param s_meta: The RDF Term Metadata associated to **s**. :type s_meta: arango_rdf.typings.RDFTermMeta :param p: The RDF Predicate URIRef of the statement (s, p, o). @@ -2248,8 +2345,10 @@ def __rpt_contextualize_statement( to this statement (if any). :type sg_str: str """ - p_meta = self.__rpt_process_term(p) - self.__contextualize_statement(s_meta, p_meta, o_meta, sg_str, is_pgt=False) + p_meta = self.__rpt_process_term(adb_docs, p) + self.__contextualize_statement( + adb_docs, s_meta, p_meta, o_meta, sg_str, is_pgt=False + ) def __rpt_create_adb_graph(self, name: str) -> ADBGraph: """RDF -> ArangoDB (RPT): Create an ArangoDB graph based on @@ -2328,6 +2427,10 @@ def __pgt_remove_blacklisted_statements(self) -> None: def __pgt_parse_literal_statements( self, + adb_docs: ADBDocs, + literal_statements: DefaultDict[ + Tuple[RDFTerm, URIRef], List[Tuple[Literal, List[Any]]] + ], pgt_contextualize_statement_func: Callable[..., None], batch_size: Optional[int], adb_import_kwargs: Dict[str, Any], @@ -2349,76 +2452,121 @@ def __pgt_parse_literal_statements( `ArangoRDF.__insert_adb_docs()`. :type adb_import_kwargs: Dict[str, Any] """ - # TODO: Revisit FILTER clauses - # We rely on the FILTER clauses to make sure no literal - # statements belonging to RDF Lists are processed, - # as that is handled in another step (which also needs work...) - query = f""" - PREFIX rdf: <{RDF}> - - SELECT ?subject ?predicate - WHERE {{ - ?subject ?predicate ?object . - FILTER isLiteral(?object) - FILTER NOT EXISTS {{ ?subject rdf:first ?anyObject }} - FILTER NOT EXISTS {{ ?subject rdf:rest ?anyObject }} - FILTER (!regex(str(?predicate), "^{RDF}_[0-9]+$")) - FILTER (!regex(str(?predicate), "^{RDF}li$")) - }} - GROUP BY ?subject ?predicate - """ - - with get_spinner_progress("(RDF → ADB): PGT [RDF Literals (Query)]") as sp: - sp.add_task("") - - data = self.__rdf_graph.query(query) s: RDFTerm p: URIRef o: Literal - total = len(data) + total = len(literal_statements) batch_size = batch_size or total - bar_progress = get_bar_progress("(RDF → ADB): PGT [RDF Literals]", "#EF7D00") - bar_progress_task = bar_progress.add_task("", total=total - 1) + bar_progress = get_bar_progress("(RDF → ADB): PGT [Literals]", "#EF7D00") + bar_progress_task = bar_progress.add_task("", total=total) spinner_progress = get_import_spinner_progress(" ") - statements = ( - self.__rdf_graph.quads - if isinstance(self.__rdf_graph, RDFConjunctiveGraph) - else self.__rdf_graph.triples - ) - with Live(Group(bar_progress, spinner_progress)): - for i, (s, p) in enumerate(data, 1): - bar_progress.update(bar_progress_task, advance=1) + for i, (k, v) in enumerate(literal_statements.items(), 1): + s, p = k s_meta = self.__pgt_get_term_metadata(s) - self.__pgt_process_rdf_term(s_meta) + self.__pgt_process_rdf_term(adb_docs, s_meta) p_meta = self.__pgt_get_term_metadata(p) - self.__pgt_process_rdf_term(p_meta) + self.__pgt_process_rdf_term(adb_docs, p_meta) _, s_col, s_key, _ = s_meta _, _, _, p_label = p_meta - for _, _, o, *sg in statements((s, p, None)): + for o, sg in v: sg_str = self.__get_subgraph_str(sg) o_meta = self.__pgt_get_term_metadata(o) - self.__pgt_process_rdf_literal(o, s_col, s_key, p_label, sg_str) + self.__pgt_process_rdf_literal( + adb_docs, o, s_col, s_key, p_label, sg_str + ) pgt_contextualize_statement_func(s_meta, p_meta, o_meta, sg_str) - self.__rdf_graph.remove((s, p, o)) + if i % batch_size == 0: + bar_progress.update(bar_progress_task, advance=batch_size) + self.__insert_adb_docs( + adb_docs, spinner_progress, **adb_import_kwargs + ) + + last_advance = total % batch_size if batch_size > 0 else 0 + bar_progress.update(bar_progress_task, advance=last_advance) + self.__insert_adb_docs(adb_docs, spinner_progress, **adb_import_kwargs) + + def __pgt_parse_non_literal_statements( + self, + adb_docs: ADBDocs, + non_literal_statements: DefaultDict[ + Tuple[RDFTerm, URIRef], List[Tuple[RDFTerm, List[Any]]] + ], + contextualize_statement_func: Callable[..., None], + batch_size: Optional[int], + adb_import_kwargs: Dict[str, Any], + ) -> None: + """RDF -> ArangoDB (PGT): Processes all non-literal RDF statements. + + :param non_literal_statements: Dictionary mapping (s,p) pairs to lists + of (o, sg) tuples for non-literal objects. + :type non_literal_statements: DefaultDict[ + Tuple[RDFTerm, URIRef], List[Tuple[RDFTerm, List[Any]]]] + ] + :param contextualize_statement_func: A function that contextualizes + an RDF Statement. A no-op function is used if Graph Contextualization + is disabled. + :type contextualize_statement_func: Callable[..., None] + :param batch_size: The batch size to use when inserting ArangoDB Documents. + Defaults to None. + :type batch_size: int | None + :param adb_import_kwargs: The keyword arguments to pass to + `ArangoRDF.__insert_adb_docs()`. + :type adb_import_kwargs: Dict[str, Any] + """ + + s: RDFTerm # Subject + p: URIRef # Predicate + o: RDFTerm # Object + + total = len(non_literal_statements) + batch_size = batch_size or total + bar_progress = get_bar_progress("(RDF → ADB): PGT [Non-Literals]", "#08479E") + bar_progress_task = bar_progress.add_task("", total=total) + spinner_progress = get_import_spinner_progress(" ") + + with Live(Group(bar_progress, spinner_progress)): + for i, (k, v) in enumerate(non_literal_statements.items(), 1): + s, p = k + + rdf_list_col = self.__is_rdf_list_statement(s, p) + + if rdf_list_col: + doc = self.__rdf_list_data[rdf_list_col][s] + predicate_label = self.rdf_id_to_adb_label(str(p)) + + for o, sg in v: + self.__pgt_rdf_val_to_adb_val(doc, predicate_label, o) + + else: + for o, sg in v: + self.__pgt_process_subject_predicate_object( + adb_docs, s, p, o, sg, None, contextualize_statement_func + ) if i % batch_size == 0: - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + bar_progress.update(bar_progress_task, advance=batch_size) + self.__insert_adb_docs( + adb_docs, spinner_progress, **adb_import_kwargs + ) - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + last_advance = total % batch_size if batch_size > 0 else 0 + bar_progress.update(bar_progress_task, advance=last_advance) + self.__insert_adb_docs(adb_docs, spinner_progress, **adb_import_kwargs) def __pgt_process_subject_predicate_object( self, + adb_docs: ADBDocs, s: RDFTerm, p: URIRef, o: RDFTerm, @@ -2429,6 +2577,8 @@ def __pgt_process_subject_predicate_object( """RDF -> ArangoDB (PGT): Processes the RDF Statement (s, p, o) as an ArangoDB document for PGT. + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param s: The RDF Subject of the RDF Statement. :type s: URIRef | BNode :param p: The RDF Predicate of the RDF Statement. @@ -2449,15 +2599,17 @@ def __pgt_process_subject_predicate_object( sg_str = self.__get_subgraph_str(sg) s_meta = self.__pgt_get_term_metadata(s) - self.__pgt_process_rdf_term(s_meta) + self.__pgt_process_rdf_term(adb_docs, s_meta) p_meta = self.__pgt_get_term_metadata(p) - self.__pgt_process_rdf_term(p_meta) + self.__pgt_process_rdf_term(adb_docs, p_meta) o_meta = self.__pgt_get_term_metadata(o) - self.__pgt_process_object(s_meta, p_meta, o_meta, sg_str) + self.__pgt_process_object(adb_docs, s_meta, p_meta, o_meta, sg_str) - self.__pgt_process_statement(s_meta, p_meta, o_meta, sg_str, reified_subject) + self.__pgt_process_statement( + adb_docs, s_meta, p_meta, o_meta, sg_str, reified_subject + ) contextualize_statement_func(s_meta, p_meta, o_meta, sg_str) @@ -2475,10 +2627,16 @@ def __pgt_get_term_metadata(self, t: Union[URIRef, BNode, Literal]) -> RDFTermMe Collection name, Document Key, and Document label. :rtype: Tuple[URIRef | BNode | Literal, str, str, str] """ - if type(t) is Literal: + # Quick return for Literals (no caching needed) + if isinstance(t, Literal): return t, "", "", "" # No other metadata needed t_str = str(t) + + if self.enable_pgt_cache: + if t_str in self.pgt_term_metadata_cache: + return self.pgt_term_metadata_cache[t_str] + t_col = "" t_key = self.rdf_id_to_adb_key(t_str, t) t_label = self.rdf_id_to_adb_label(t_str) @@ -2506,7 +2664,11 @@ def __pgt_get_term_metadata(self, t: Union[URIRef, BNode, Literal]) -> RDFTermMe logger.debug(f"Found unknown resource: {t} ({t_key})") t_col = self.__UNKNOWN_RESOURCE - return t, str(t_col), t_key, t_label + result = t, str(t_col), t_key, t_label + if self.enable_pgt_cache: + self.pgt_term_metadata_cache[t_str] = result + + return result def __pgt_rdf_val_to_adb_val( self, @@ -2544,20 +2706,21 @@ def __pgt_rdf_val_to_adb_val( # This flag is only active in ArangoRDF.__pgt_process_rdf_lists() if process_val_as_serialized_list: - doc[key] += f"'{val}'," if type(val) is str else f"{val}," + doc[key] += f"'{val}'," if isinstance(val, str) else f"{val}," return prev_val = doc.get(key) if prev_val is None: doc[key] = val - elif type(prev_val) is list: + elif isinstance(prev_val, list): prev_val.append(val) else: doc[key] = [prev_val, val] def __pgt_process_rdf_term( self, + adb_docs: ADBDocs, t_meta: RDFTermMeta, s_col: str = "", s_key: str = "", @@ -2589,7 +2752,7 @@ def __pgt_process_rdf_term( t, t_col, t_key, t_label = t_meta - if t_key in self.__adb_docs.get(t_col, {}): + if t_key in adb_docs.get(t_col, {}): return if t in self.__reified_subject_map: @@ -2598,14 +2761,14 @@ def __pgt_process_rdf_term( if self.__predicate_collection is not None: t_col = self.__predicate_collection.name - self.__adb_docs[t_col][t_key] = { + adb_docs[t_col][t_key] = { "_key": t_key, "_from": _from, "_to": _to, } - elif type(t) is URIRef: - self.__adb_docs[t_col][t_key] = { + elif isinstance(t, URIRef): + adb_docs[t_col][t_key] = { "_key": t_key, self.__rdf_uri_attr: str(t), self.__rdf_label_attr: t_label, @@ -2617,22 +2780,28 @@ def __pgt_process_rdf_term( and t_col != self.__UNKNOWN_RESOURCE ): uri_col = self.__uri_map_collection.name - self.__adb_docs[uri_col][t_key] = { + adb_docs[uri_col][t_key] = { "_key": t_key, "collection": t_col, self.__rdf_uri_attr: str(t), } - elif type(t) is BNode: - self.__adb_docs[t_col][t_key] = { + elif isinstance(t, BNode): + adb_docs[t_col][t_key] = { "_key": t_key, self.__rdf_label_attr: "", self.__rdf_type_attr: "BNode", } - elif type(t) is Literal and all([s_col, s_key, p_label]): + elif isinstance(t, Literal) and all([s_col, s_key, p_label]): self.__pgt_process_rdf_literal( - t, s_col, s_key, p_label, sg_str, process_val_as_serialized_list + adb_docs, + t, + s_col, + s_key, + p_label, + sg_str, + process_val_as_serialized_list, ) else: @@ -2640,6 +2809,7 @@ def __pgt_process_rdf_term( def __pgt_process_rdf_literal( self, + adb_docs: ADBDocs, literal: Literal, s_col: str, s_key: str, @@ -2650,6 +2820,8 @@ def __pgt_process_rdf_literal( """RDF -> ArangoDB (PGT): Process an RDF Literal as an ArangoDB document property. + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param literal: The RDF Literal to process :type literal: Literal :param s_col: The ArangoDB document collection of the Subject associated @@ -2669,7 +2841,7 @@ def __pgt_process_rdf_literal( property. Defaults to False. :type process_val_as_serialized_list: bool """ - doc = self.__adb_docs[s_col][s_key] + doc = adb_docs[s_col][s_key] val = self.__get_literal_val(literal, str(literal)) self.__pgt_rdf_val_to_adb_val(doc, p_label, val, process_val_as_serialized_list) @@ -2677,7 +2849,12 @@ def __pgt_process_rdf_literal( doc[self.__rdf_sub_graph_uri_attr] = sg_str def __pgt_process_object( - self, s_meta: RDFTermMeta, p_meta: RDFTermMeta, o_meta: RDFTermMeta, sg_str: str + self, + adb_docs: ADBDocs, + s_meta: RDFTermMeta, + p_meta: RDFTermMeta, + o_meta: RDFTermMeta, + sg_str: str, ) -> None: """RDF -> ArangoDB (PGT): Processes the RDF Object into ArangoDB. Given the possibily of the RDF Object being used as the "root" of @@ -2685,6 +2862,8 @@ def __pgt_process_object( function is used to prevent calling `__pgt_process_rdf_term` if it is not required. + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param s_meta: The RDF Term Metadata associated to the RDF Subject of the statement containing the RDF Object. :type s_meta: arango_rdf.typings.RDFTermMeta @@ -2707,10 +2886,13 @@ def __pgt_process_object( self.__rdf_list_heads[s][p] = head else: - self.__pgt_process_rdf_term(o_meta, s_col, s_key, p_label, sg_str=sg_str) + self.__pgt_process_rdf_term( + adb_docs, o_meta, s_col, s_key, p_label, sg_str=sg_str + ) def __pgt_process_statement( self, + adb_docs: ADBDocs, s_meta: RDFTermMeta, p_meta: RDFTermMeta, o_meta: RDFTermMeta, @@ -2724,6 +2906,8 @@ def __pgt_process_statement( 1) The RDF Object within the RDF Statement is not a Literal 2) The RDF Object is not the "root" node of an RDF List structure + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param s_meta: The RDF Term Metadata associated to the RDF Subject of the statement containing the RDF Object. :type s_meta: arango_rdf.typings.RDFTermMeta @@ -2742,7 +2926,7 @@ def __pgt_process_statement( """ o, o_col, o_key, _ = o_meta - if type(o) is Literal or self.__pgt_object_is_head_of_rdf_list(o): + if isinstance(o, Literal) or self.__pgt_object_is_head_of_rdf_list(o): return _, s_col, s_key, _ = s_meta @@ -2766,6 +2950,7 @@ def __pgt_process_statement( e_key = self.hash(f"{s_key}-{p_key}-{o_key}") self.__add_adb_edge( + adb_docs, e_col, e_key, _from, @@ -2790,56 +2975,16 @@ def __pgt_object_is_head_of_rdf_list(self, o: RDFTerm) -> bool: :return: Whether the object points to an RDF List or not. :rtype: bool """ - # TODO: Discuss repurcussions of this assumption - if type(o) is not BNode: + # Quick check: if not a BNode, it can't be an RDF list head + if not isinstance(o, BNode): return False - first = (o, RDF.first, None) - rest = (o, RDF.rest, None) - - if first in self.__rdf_graph or rest in self.__rdf_graph: - return True - - _n = (o, URIRef(f"{RDF}_1"), None) - li = (o, URIRef(f"{RDF}li"), None) - - if _n in self.__rdf_graph or li in self.__rdf_graph: - return True - - return False - - def __pgt_statement_is_part_of_rdf_list(self, s: RDFTerm, p: URIRef) -> str: - """RDF -> ArangoDB (PGT): Return the associated "Document Buffer" key - if the RDF Statement (s, p, _) is part of an RDF Collection or RDF Container - within the RDF Graph. Essential for unpacking the complicated data structure of - RDF Lists and re-building them as an ArangoDB Document Property. - - :param s: The RDF Subject. - :type s: URIRef | BNode - :param p: The RDF Predicate. - :type p: URIRef - :return: The **self.adb_docs** "Document Buffer" key associated - to the RDF Statement. If the statement is not part of an RDF - List, return an empty string. - :rtype: str - """ - # TODO: Discuss repurcussions of this assumption - if type(s) is not BNode: - return "" - - if p in {RDF.first, RDF.rest}: - return "_COLLECTION_BNODE" - - p_str = str(p) - _n = r"^http://www.w3.org/1999/02/22-rdf-syntax-ns#_[0-9]{1,}$" - li = r"^http://www.w3.org/1999/02/22-rdf-syntax-ns#li$" - - if re.match(_n, p_str) or re.match(li, p_str): - return "_CONTAINER_BNODE" - - return "" + # Use pre-computed RDF list subjects for O(1) lookup + return o in self.__rdf_list_subjects - def __pgt_process_rdf_lists(self, bar_progress: Progress) -> None: + def __pgt_process_rdf_lists( + self, adb_docs: ADBDocs, bar_progress: Progress + ) -> None: """RDF -> ArangoDB (PGT): Process all RDF Collections & Containers within the RDF Graph prior to inserting the documents into ArangoDB. @@ -2869,7 +3014,7 @@ def __pgt_process_rdf_lists(self, bar_progress: Progress) -> None: s_meta = self.__pgt_get_term_metadata(s) _, s_col, s_key, _ = s_meta - doc = self.__adb_docs[s_col][s_key] + doc = adb_docs[s_col][s_key] doc["_key"] = s_key for p, p_dict in s_dict.items(): @@ -2880,7 +3025,9 @@ def __pgt_process_rdf_lists(self, bar_progress: Progress) -> None: sg: str = p_dict["sub_graph"] doc[p_label] = "" - self.__pgt_process_rdf_list_object(doc, s_meta, p_meta, root, sg) + self.__pgt_process_rdf_list_object( + adb_docs, doc, s_meta, p_meta, root, sg + ) doc[p_label] = doc[p_label].rstrip(",") # Delete doc[p_key] if there are no Literals within the RDF List @@ -2892,6 +3039,7 @@ def __pgt_process_rdf_lists(self, bar_progress: Progress) -> None: def __pgt_process_rdf_list_object( self, + adb_docs: ADBDocs, doc: Json, s_meta: RDFTermMeta, p_meta: RDFTermMeta, @@ -2929,7 +3077,9 @@ def __pgt_process_rdf_list_object( doc[p_label] += "[" next_bnode_dict = self.__rdf_list_data["_COLLECTION_BNODE"][o] - self.__pgt_unpack_rdf_collection(doc, s_meta, p_meta, next_bnode_dict, sg) + self.__pgt_unpack_rdf_collection( + adb_docs, doc, s_meta, p_meta, next_bnode_dict, sg + ) doc[p_label] = doc[p_label].rstrip(",") + "]," @@ -2937,7 +3087,9 @@ def __pgt_process_rdf_list_object( doc[p_label] += "[" next_bnode_dict = self.__rdf_list_data["_CONTAINER_BNODE"][o] - self.__pgt_unpack_rdf_container(doc, s_meta, p_meta, next_bnode_dict, sg) + self.__pgt_unpack_rdf_container( + adb_docs, doc, s_meta, p_meta, next_bnode_dict, sg + ) doc[p_label] = doc[p_label].rstrip(",") + "]," @@ -2947,13 +3099,19 @@ def __pgt_process_rdf_list_object( # Process the RDF Object as an ArangoDB Document self.__pgt_process_rdf_term( - o_meta, s_col, s_key, p_label, process_val_as_serialized_list=True + adb_docs, + o_meta, + s_col, + s_key, + p_label, + process_val_as_serialized_list=True, ) # Process the RDF Statement as an ArangoDB Edge - self.__pgt_process_statement(s_meta, p_meta, o_meta, sg) + self.__pgt_process_statement(adb_docs, s_meta, p_meta, o_meta, sg) def __pgt_unpack_rdf_collection( self, + adb_docs: ADBDocs, doc: Json, s_meta: RDFTermMeta, p_meta: RDFTermMeta, @@ -2979,16 +3137,19 @@ def __pgt_unpack_rdf_collection( """ first: RDFTerm = bnode_dict["first"] - self.__pgt_process_rdf_list_object(doc, s_meta, p_meta, first, sg) + self.__pgt_process_rdf_list_object(adb_docs, doc, s_meta, p_meta, first, sg) if "rest" in bnode_dict and bnode_dict["rest"] != RDF.nil: rest = bnode_dict["rest"] next_bnode_dict = self.__rdf_list_data["_COLLECTION_BNODE"][rest] - self.__pgt_unpack_rdf_collection(doc, s_meta, p_meta, next_bnode_dict, sg) + self.__pgt_unpack_rdf_collection( + adb_docs, doc, s_meta, p_meta, next_bnode_dict, sg + ) def __pgt_unpack_rdf_container( self, + adb_docs: ADBDocs, doc: Json, s_meta: RDFTermMeta, p_meta: RDFTermMeta, @@ -3019,15 +3180,22 @@ def __pgt_unpack_rdf_container( # It is possible for the Container Membership Property # to be re-used in multiple statements (e.g rdf:li), # hence the reason why `value` can be a list or a single element. - value_as_list = value if type(value) is list else [value] + value_as_list = value if isinstance(value, list) else [value] for o in value_as_list: - self.__pgt_process_rdf_list_object(doc, s_meta, p_meta, o, sg) + self.__pgt_process_rdf_list_object(adb_docs, doc, s_meta, p_meta, o, sg) def __pgt_contextualize_statement( - self, s_meta: RDFTermMeta, p_meta: RDFTermMeta, o_meta: RDFTermMeta, sg_str: str + self, + adb_docs: ADBDocs, + s_meta: RDFTermMeta, + p_meta: RDFTermMeta, + o_meta: RDFTermMeta, + sg_str: str, ) -> None: """RDF -> ArangoDB (PGT): Contextualizes the RDF Statement (s, p, o). + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param s_meta: The RDF Term Metadata associated to **s**. :type s_meta: arango_rdf.typings.RDFTermMeta :param p_meta: The RDF Term Metadata associated to **p**. @@ -3038,7 +3206,9 @@ def __pgt_contextualize_statement( to this statement (if any). :type sg_str: str """ - self.__contextualize_statement(s_meta, p_meta, o_meta, sg_str, is_pgt=True) + self.__contextualize_statement( + adb_docs, s_meta, p_meta, o_meta, sg_str, is_pgt=True + ) def __pgt_create_adb_graph(self, name: str) -> ADBGraph: """RDF -> ArangoDB (PGT): Create an ArangoDB graph based @@ -3157,6 +3327,7 @@ def __load_meta_ontology(self, rdf_graph: RDFGraph) -> RDFConjunctiveGraph: def __flatten_reified_triples( self, + adb_docs: ADBDocs, process_subject_predicate_object: Callable[..., None], contextualize_statement_func: Callable[..., None], batch_size: Optional[int], @@ -3198,7 +3369,7 @@ def process_reified_subject( process_reified_subject(new_reified_subject, sg) process_subject_predicate_object( - s, p, o, sg, reified_subject, contextualize_statement_func + adb_docs, s, p, o, sg, reified_subject, contextualize_statement_func ) # Remove the reified triple from the RDF Graph @@ -3241,17 +3412,20 @@ def process_reified_subject( with Live(Group(bar_progress, spinner_progress)): for i, (reified_subject, *sg) in enumerate(data, 1): - bar_progress.advance(bar_progress_task) - # Only process the reified triple if it has not been processed yet # i.e recursion if reified_subject not in self.__reified_subject_map: process_reified_subject(reified_subject, sg) if i % batch_size == 0: - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + bar_progress.update(bar_progress_task, advance=batch_size) + self.__insert_adb_docs( + adb_docs, spinner_progress, **adb_import_kwargs + ) - self.__insert_adb_docs(spinner_progress, **adb_import_kwargs) + last_advance = total % batch_size if batch_size > 0 else 0 + bar_progress.update(bar_progress_task, advance=last_advance) + self.__insert_adb_docs(adb_docs, spinner_progress, **adb_import_kwargs) def __get_subgraph_str(self, possible_sg: Optional[List[Any]]) -> str: """RDF -> ArangoDB: Extract the sub-graph URIRef string of a quad (if any). @@ -3267,16 +3441,17 @@ def __get_subgraph_str(self, possible_sg: Optional[List[Any]]) -> str: sg = possible_sg[0] sg_identifier = sg.identifier if isinstance(sg, RDFGraph) else sg - if type(sg_identifier) is URIRef: + if isinstance(sg_identifier, URIRef): return str(sg_identifier) - if type(sg_identifier) is BNode: + if isinstance(sg_identifier, BNode): return "" # TODO: Revisit raise ValueError(f"Sub Graph Identifier is not a URIRef or BNode: {sg}") def __add_adb_edge( self, + adb_docs: ADBDocs, col: str, key: str, _from: str, @@ -3286,9 +3461,11 @@ def __add_adb_edge( _sg: str, ) -> None: """RDF -> ArangoDB: Insert the JSON-equivalent of an ArangoDB Edge - into `self.adb_docs` for temporary storage, until it gets + into `adb_docs` for temporary storage, until it gets ingested into the **col** ArangoDB Collection. + :param adb_docs: The ArangoDB documents buffer to populate. + :type adb_docs: ADBDocs :param col: The name of the ArangoDB Edge Collection. :type col: str :param key: The ArangoDB Key of the Edge. @@ -3308,8 +3485,8 @@ def __add_adb_edge( if self.__predicate_collection is not None: col = self.__predicate_collection.name - self.__adb_docs[col][key] = { - **self.__adb_docs[col][key], + adb_docs[col][key] = { + **adb_docs[col][key], "_key": key, "_from": _from, "_to": _to, @@ -3319,7 +3496,7 @@ def __add_adb_edge( } if _sg: - self.__adb_docs[col][key][self.__rdf_sub_graph_uri_attr] = _sg + adb_docs[col][key][self.__rdf_sub_graph_uri_attr] = _sg def __build_explicit_type_map( self, adb_adb_col_statement: Callable[..., None] = empty_func @@ -3636,10 +3813,12 @@ def __get_literal_val(self, t: Literal, t_str: str) -> Any: return t.value if t.value is not None else t_str def __insert_adb_docs( - self, spinner_progress: Progress, **adb_import_kwargs: Any + self, adb_docs: ADBDocs, spinner_progress: Progress, **adb_import_kwargs: Any ) -> None: """RDF -> ArangoDB: Insert ArangoDB documents into their ArangoDB collection. + :param adb_docs: The ArangoDB documents buffer to insert. + :type adb_docs: ADBDocs :param spinner_progress: The spinner progress bar. :type spinner_progress: rich.progress.Progress :param adb_import_kwargs: Keyword arguments to specify additional @@ -3647,19 +3826,22 @@ def __insert_adb_docs( https://docs.python-arango.com/en/main/specs.html#arango.collection.Collection.insert_many :param adb_import_kwargs: Any """ - if len(self.__adb_docs) == 0: + if len(adb_docs) == 0: return + db = self.async_db if self.insert_async else self.db + adb_import_kwargs["overwrite_mode"] = "update" adb_import_kwargs["merge"] = True + if "raise_on_document_error" not in adb_import_kwargs: adb_import_kwargs["raise_on_document_error"] = True # Avoiding "RuntimeError: dictionary changed size during iteration" - adb_cols = list(self.__adb_docs.keys()) + adb_cols = list(adb_docs.keys()) for col in adb_cols: - doc_list = self.__adb_docs[col].values() + doc_list = adb_docs[col].values() action = f"(RDF → ADB): Import '{col}' ({len(doc_list)})" spinner_progress_task = spinner_progress.add_task("", action=action) @@ -3671,9 +3853,7 @@ def __insert_adb_docs( logger.debug(f"Inserting Documents: {doc_list}") try: - result = self.db.collection(col).insert_many( - doc_list, **adb_import_kwargs - ) + result = db.collection(col).insert_many(doc_list, **adb_import_kwargs) except Exception as e: e_str = str(e) @@ -3682,13 +3862,14 @@ def __insert_adb_docs( logger.debug(f"Insert Result: {result}") - del self.__adb_docs[col] + del adb_docs[col] spinner_progress.stop_task(spinner_progress_task) spinner_progress.update(spinner_progress_task, visible=False) def __contextualize_statement( self, + adb_docs: ADBDocs, s_meta: RDFTermMeta, p_meta: RDFTermMeta, o_meta: RDFTermMeta, @@ -3722,6 +3903,7 @@ def __contextualize_statement( _to_col = "Class" if is_pgt else self.__URIREF_COL self.__add_adb_edge( + adb_docs, col=edge_col, key=self.hash(edge_key), _from=f"{_from_col}/{p_key}", @@ -3733,10 +3915,11 @@ def __contextualize_statement( # Run RDFS Domain/Range Inference & Introspection dr_meta = [(*s_meta, "domain"), (*o_meta, "range")] - self.__infer_and_introspect_dr(p, p_key, dr_meta, sg_str, is_pgt) + self.__infer_and_introspect_dr(adb_docs, p, p_key, dr_meta, sg_str, is_pgt) def __infer_and_introspect_dr( self, + adb_docs: ADBDocs, p: URIRef, p_key: str, dr_meta: List[Tuple[RDFTerm, str, str, str, str]], @@ -3794,7 +3977,7 @@ def __infer_and_introspect_dr( self.__e_col_map[e_col_type]["to"].add("Class") for t, t_col, t_key, t_label, dr_label in dr_meta: - if type(t) is Literal: + if isinstance(t, Literal): continue DR_COL = dr_label if is_pgt else self.__STATEMENT_COL @@ -3807,6 +3990,7 @@ def __infer_and_introspect_dr( for _, class_key in self.__predicate_scope[p][dr_label]: key = self.hash(f"{t_key}-{self.__rdf_type_key}-{class_key}") self.__add_adb_edge( + adb_docs, col=TYPE_COL, key=key, _from=f"{t_col}/{t_key}", @@ -3833,6 +4017,7 @@ def __infer_and_introspect_dr( class_key = self.rdf_id_to_adb_key(class_str) key = self.hash(f"{p_key}-{dr_key}-{class_key}") self.__add_adb_edge( + adb_docs, col=DR_COL, key=key, _from=f"{P_COL}/{p_key}", diff --git a/tests/conftest.py b/tests/conftest.py index 6a51a16..b880c56 100644 --- a/tests/conftest.py +++ b/tests/conftest.py @@ -3,6 +3,7 @@ from pathlib import Path from typing import Any, Dict, Set, Tuple +import pytest from arango import ArangoClient, DefaultHTTPClient from arango.database import StandardDatabase from rdflib import BNode @@ -53,6 +54,17 @@ class NoTimeoutHTTPClient(DefaultHTTPClient): adbrdf = ArangoRDF(db) +@pytest.fixture(autouse=True) +def reset_arango_db() -> None: + global db + for g in db.graphs(): + db.delete_graph(g["name"], drop_collections=True) + + for c in db.collections(): + if c["system"] is False: + db.delete_collection(c["name"]) + + def arango_restore(path_to_data: str) -> None: global con restore_prefix = "./tools/" if os.getenv("GITHUB_ACTIONS") else ""