diff --git a/arango_rdf/main.py b/arango_rdf/main.py index c16bb87..7879da3 100644 --- a/arango_rdf/main.py +++ b/arango_rdf/main.py @@ -40,6 +40,8 @@ ) from .utils import ( Node, + NoOpLive, + NoOpProgress, Tree, empty_func, get_bar_progress, @@ -75,6 +77,10 @@ class ArangoRDF(AbstractArangoRDF): repeated computations. Defaults to False. Not always useful, especially when terms are not repeated alot in the RDF graph. :type enable_pgt_cache: bool + :param enable_rich: If True, will enable rich progress bars and spinners. + Defaults to True. Set to False when using multiprocessing or concurrent + modules, as rich can interfere with them. + :type enable_rich: bool :raise TypeError: On invalid parameter types """ @@ -86,6 +92,7 @@ def __init__( rdf_attribute_prefix: str = "_", insert_async: bool = False, enable_pgt_cache: bool = False, + enable_rich: bool = True, ): self.set_logging(logging_lvl) @@ -143,6 +150,9 @@ def __init__( self.enable_pgt_cache = enable_pgt_cache self.pgt_term_metadata_cache: Dict[str, RDFTermMeta] = {} + # Rich progress bar configuration + self.enable_rich = enable_rich + # RDF Graph for maintaining the ArangoDB Collections & Keys # of the RDF Resources self.__adb_col_statements = RDFGraph() @@ -188,6 +198,30 @@ def rdf_attribute_prefix(self) -> str: def set_logging(self, level: Union[int, str]) -> None: logger.setLevel(level) + def _get_spinner_progress(self, text: str) -> Union[Progress, NoOpProgress]: + """Get a spinner progress bar or no-op version based on enable_rich.""" + if self.enable_rich: + return get_spinner_progress(text) + return NoOpProgress() + + def _get_bar_progress(self, text: str, color: str) -> Union[Progress, NoOpProgress]: + """Get a bar progress or no-op version based on enable_rich.""" + if self.enable_rich: + return get_bar_progress(text, color) + return NoOpProgress() + + def _get_import_spinner_progress(self, text: str) -> Union[Progress, NoOpProgress]: + """Get an import spinner progress or no-op version based on enable_rich.""" + if self.enable_rich: + return get_import_spinner_progress(text) + return NoOpProgress() + + def _live_context(self, *renderables: Any) -> Union[Live, NoOpLive]: + """Get a Live context manager or no-op version based on enable_rich.""" + if self.enable_rich: + return Live(Group(*renderables)) + return NoOpLive() + ########################### # Public: ArangoDB -> RDF # ########################### @@ -749,7 +783,9 @@ def contextualize_statement_func( self.__rdf_graph = self.__load_meta_ontology(self.__rdf_graph) - with get_spinner_progress("(RDF → ADB): Graph Contextualization") as rp: + with self._get_spinner_progress( + "(RDF → ADB): Graph Contextualization" + ) as rp: rp.add_task("") self.__explicit_type_map = self.__build_explicit_type_map() @@ -782,9 +818,9 @@ def contextualize_statement_func( total = len(self.__rdf_graph) batch_size = batch_size or total - bar_progress = get_bar_progress("(RDF → ADB): RPT", "#BF23C4") + bar_progress = self._get_bar_progress("(RDF → ADB): RPT", "#BF23C4") bar_progress_task = bar_progress.add_task("", total=total) - spinner_progress = get_import_spinner_progress(" ") + spinner_progress = self._get_import_spinner_progress(" ") statements = ( self.__rdf_graph.quads @@ -792,7 +828,7 @@ def contextualize_statement_func( else self.__rdf_graph.triples ) - with Live(Group(bar_progress, spinner_progress)): + with self._live_context(bar_progress, spinner_progress): for i, (s, p, o, *sg) in enumerate(statements((None, None, None)), 1): logger.debug(f"RPT: {s} {p} {o} {sg}") @@ -956,21 +992,21 @@ def rdf_to_arangodb_by_pgt( raise ValueError(m) if not self.db.has_collection(uri_map_collection_name): - self.db.create_collection(uri_map_collection_name) + self.__create_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) + self.__create_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) + self.__create_collection(predicate_collection_name, edge=True) self.__predicate_collection = self.db.collection(predicate_collection_name) @@ -1092,7 +1128,7 @@ def contextualize_statement_func( literal_statements = defaultdict(list) non_literal_statements = defaultdict(list) - with get_spinner_progress("(RDF → ADB): PGT [Prepare Statements]") as rp: + with self._get_spinner_progress("(RDF → ADB): PGT [Prepare Statements]") as rp: rp.add_task("") for s, p, o, *sg in statements((None, None, None)): @@ -1129,9 +1165,9 @@ def contextualize_statement_func( # PGT: RDF Lists # ################## - bar_progress = get_bar_progress("(RDF → ADB): PGT [Lists]", "#EF7D00") - spinner_progress = get_import_spinner_progress(" ") - with Live(Group(bar_progress, spinner_progress)): + bar_progress = self._get_bar_progress("(RDF → ADB): PGT [Lists]", "#EF7D00") + spinner_progress = self._get_import_spinner_progress(" ") + with self._live_context(bar_progress, spinner_progress): self.__pgt_process_rdf_lists(adb_docs, bar_progress) self.__insert_adb_docs(adb_docs, spinner_progress, **adb_import_kwargs) @@ -1141,7 +1177,7 @@ def contextualize_statement_func( if namespace_collection_name: if not self.db.has_collection(namespace_collection_name): - self.db.create_collection(namespace_collection_name) + self.__create_collection(namespace_collection_name) docs = [ {"prefix": prefix, "uri": uri, "_key": self.hash(uri)} @@ -1240,7 +1276,7 @@ def write_adb_col_statements( self.__rdf_graph = rdf_graph self.controller.rdf_graph = rdf_graph - with get_spinner_progress("(RDF → ADB): Write Col Statements") as rp: + with self._get_spinner_progress("(RDF → ADB): Write Col Statements") as rp: rp.add_task("") # 0. Add URI Collection statements @@ -1381,7 +1417,7 @@ def migrate_unknown_resources( edge_count = 0 - with get_spinner_progress("(RDF → ADB): Migrate Unknown Resources") as sp: + with self._get_spinner_progress("(RDF → ADB): Migrate Unknown Resources") as sp: sp.add_task("") while not cursor.empty(): @@ -1715,7 +1751,9 @@ def __fetch_adb_docs( col_size: int = self.db.collection(col).count() - with get_spinner_progress(f"(ADB → RDF): Export '{col}' ({col_size})") as sp: + with self._get_spinner_progress( + f"(ADB → RDF): Export '{col}' ({col_size})" + ) as sp: sp.add_task("") cursor: Cursor = self.db.aql.execute( @@ -1751,10 +1789,10 @@ def __process_adb_cursor( :type col_uri: URIRef """ - progress = get_bar_progress(f"(ADB → RDF): '{col}'", progress_color) + progress = self._get_bar_progress(f"(ADB → RDF): '{col}'", progress_color) progress_task_id = progress.add_task("", total=col_size) - with Live(Group(progress)): + with self._live_context(progress): while not cursor.empty(): for doc in cursor.batch(): process_adb_doc(doc, col, col_uri) @@ -2363,23 +2401,23 @@ def __rpt_create_adb_graph(self, name: str) -> ADBGraph: if self.db.has_graph(name): # pragma: no cover return self.db.graph(name) - return self.db.create_graph( - name, - edge_definitions=[ - { - "edge_collection": self.__STATEMENT_COL, - "from_vertex_collections": [ - self.__URIREF_COL, - self.__BNODE_COL, - ], - "to_vertex_collections": [ - self.__URIREF_COL, - self.__BNODE_COL, - self.__LITERAL_COL, - ], - } - ], - ) + edge_definitions = [ + { + "edge_collection": self.__STATEMENT_COL, + "from_vertex_collections": [ + self.__URIREF_COL, + self.__BNODE_COL, + ], + "to_vertex_collections": [ + self.__URIREF_COL, + self.__BNODE_COL, + self.__LITERAL_COL, + ], + } + ] + + self.__create_graph(name, edge_definitions=edge_definitions) + return self.db.graph(name) ################################## # Private: RDF -> ArangoDB (PGT) # @@ -2459,11 +2497,11 @@ def __pgt_parse_literal_statements( total = len(literal_statements) batch_size = batch_size or total - bar_progress = get_bar_progress("(RDF → ADB): PGT [Literals]", "#EF7D00") + bar_progress = self._get_bar_progress("(RDF → ADB): PGT [Literals]", "#EF7D00") bar_progress_task = bar_progress.add_task("", total=total) - spinner_progress = get_import_spinner_progress(" ") + spinner_progress = self._get_import_spinner_progress(" ") - with Live(Group(bar_progress, spinner_progress)): + with self._live_context(bar_progress, spinner_progress): for i, (k, v) in enumerate(literal_statements.items(), 1): s, p = k @@ -2531,11 +2569,13 @@ def __pgt_parse_non_literal_statements( total = len(non_literal_statements) batch_size = batch_size or total - bar_progress = get_bar_progress("(RDF → ADB): PGT [Non-Literals]", "#08479E") + bar_progress = self._get_bar_progress( + "(RDF → ADB): PGT [Non-Literals]", "#08479E" + ) bar_progress_task = bar_progress.add_task("", total=total) - spinner_progress = get_import_spinner_progress(" ") + spinner_progress = self._get_import_spinner_progress(" ") - with Live(Group(bar_progress, spinner_progress)): + with self._live_context(bar_progress, spinner_progress): for i, (k, v) in enumerate(non_literal_statements.items(), 1): s, p = k @@ -3258,7 +3298,13 @@ def __pgt_create_adb_graph(self, name: str) -> ADBGraph: orphan_v_cols = orphan_v_cols ^ {self.__UNKNOWN_RESOURCE} if not self.db.has_graph(name): - return self.db.create_graph(name, edge_definitions, list(orphan_v_cols)) + self.__create_graph( + name, + edge_definitions=edge_definitions, + orphan_collections=list(orphan_v_cols), + ) + + return self.db.graph(name) old_edge_definitions = { edge_def["edge_collection"]: edge_def @@ -3296,6 +3342,33 @@ def __pgt_create_adb_graph(self, name: str) -> ADBGraph: # Private: RDF -> ArangoDB (RPT, PGT, LPG) # ############################################ + def __create_collection(self, col: str, edge: bool = False) -> None: + """RDF -> ArangoDB: Create an ArangoDB Collection.""" + try: + self.db.create_collection(col, edge=edge) + except Exception: + # Collection may have been created by another thread + if not self.db.has_collection(col): + raise + + def __create_graph( + self, + name: str, + edge_definitions: List[Dict[str, Any]], + orphan_collections: List[str] = [], + ) -> None: + """RDF -> ArangoDB: Create an ArangoDB Graph.""" + try: + self.db.create_graph( + name, + edge_definitions=edge_definitions, + orphan_collections=orphan_collections, + ) + except Exception: + # Graph may have been created by another thread + if not self.db.has_graph(name): + raise + def __load_meta_ontology(self, rdf_graph: RDFGraph) -> RDFConjunctiveGraph: """RDF -> ArangoDB: Load the RDF, RDFS, and OWL Ontologies into **rdf_graph** as 3 sub-graphs. This method returns @@ -3338,6 +3411,9 @@ def __flatten_reified_triples( NOTE: This modifies the RDF Graph in-place. TODO: Revisit + NOTE: This function is NOT thread-safe due to thread-safety issues with + rdflib's SPARQL parser. Therefore it should ONLY be called from a single thread. + :param process_subject_predicate_object: A function that processes the RDF Statement (s, p, o) as an ArangoDB document. Either `__rpt_process_subject_predicate_object` or @@ -3398,7 +3474,7 @@ def process_reified_subject( """ text = "(RDF → ADB): PGT [Flatten Reified Triples (Query)]" - with get_spinner_progress(text) as sp: + with self._get_spinner_progress(text) as sp: sp.add_task("") data = self.__rdf_graph.query(query) @@ -3406,11 +3482,11 @@ def process_reified_subject( total = len(data) batch_size = batch_size or total m = "(RDF → ADB): Flatten Reified Triples" - bar_progress = get_bar_progress(m, "#FFFFFF") + bar_progress = self._get_bar_progress(m, "#FFFFFF") bar_progress_task = bar_progress.add_task("", total=total) - spinner_progress = get_import_spinner_progress(" ") + spinner_progress = self._get_import_spinner_progress(" ") - with Live(Group(bar_progress, spinner_progress)): + with self._live_context(bar_progress, spinner_progress): for i, (reified_subject, *sg) in enumerate(data, 1): # Only process the reified triple if it has not been processed yet # i.e recursion @@ -3831,9 +3907,10 @@ def __insert_adb_docs( db = self.async_db if self.insert_async else self.db - adb_import_kwargs["overwrite_mode"] = "update" - adb_import_kwargs["merge"] = True - + if "overwrite_mode" not in adb_import_kwargs: + adb_import_kwargs["overwrite_mode"] = "update" + if "merge" not in adb_import_kwargs: + adb_import_kwargs["merge"] = True if "raise_on_document_error" not in adb_import_kwargs: adb_import_kwargs["raise_on_document_error"] = True @@ -3848,7 +3925,7 @@ def __insert_adb_docs( if not self.db.has_collection(col): is_edge = col in self.__e_col_map - self.db.create_collection(col, edge=is_edge) + self.__create_collection(col, edge=is_edge) logger.debug(f"Inserting Documents: {doc_list}") @@ -4079,7 +4156,9 @@ def __extract_statements( _, p, _ = triple - with get_spinner_progress(f"(RDF ↔ ADB): Extract Statements '{str(p)}'") as sp: + with self._get_spinner_progress( + f"(RDF ↔ ADB): Extract Statements '{str(p)}'" + ) as sp: sp.add_task("") for t in rdf_graph.triples(triple): diff --git a/arango_rdf/utils.py b/arango_rdf/utils.py index 905b0ac..f257bd6 100644 --- a/arango_rdf/utils.py +++ b/arango_rdf/utils.py @@ -54,6 +54,52 @@ def get_bar_progress(text: str, color: str) -> Progress: ) +class NoOpProgress: + """A no-op Progress class that mimics the rich Progress interface. + + Used when enable_rich=False to avoid interference with + multiprocessing or concurrent modules. + """ + + def __init__(self) -> None: + pass + + def __enter__(self) -> "NoOpProgress": + return self + + def __exit__(self, *args: Any) -> None: + pass + + def add_task(self, description: str = "", total: int = 0, **kwargs: Any) -> int: + return 0 + + def advance(self, task_id: int, advance: int = 1) -> None: + pass + + def update(self, task_id: int, **kwargs: Any) -> None: + pass + + def stop_task(self, task_id: int) -> None: + pass + + +class NoOpLive: + """A no-op Live context manager that mimics the rich Live interface. + + Used when enable_rich=False to avoid interference with + multiprocessing or concurrent modules. + """ + + def __init__(self, *args: Any, **kwargs: Any) -> None: + pass + + def __enter__(self) -> "NoOpLive": + return self + + def __exit__(self, *args: Any) -> None: + pass + + class Node: def __init__(self, name: str, depth: int = 0) -> None: self.name = name diff --git a/docs/arangodb_to_rdf.rst b/docs/arangodb_to_rdf.rst new file mode 100644 index 0000000..613663e --- /dev/null +++ b/docs/arangodb_to_rdf.rst @@ -0,0 +1,351 @@ +ArangoDB to RDF +--------------- + +ArangoRDF provides three methods to export ArangoDB Graphs to RDF format: + +1. ``arangodb_graph_to_rdf`` - Export by graph name (simplest) +2. ``arangodb_collections_to_rdf`` - Export by collection names +3. ``arangodb_to_rdf`` - Export with fine-grained control via metagraph + +All three methods return an ``rdflib.Graph`` object containing the RDF representation +of your ArangoDB data. + + +Quick Start +=========== + +.. code-block:: python + + from rdflib import Graph + from arango import ArangoClient + from arango_rdf import ArangoRDF + + db = ArangoClient().db() + adbrdf = ArangoRDF(db) + + # Export entire graph by name + rdf_graph = adbrdf.arangodb_graph_to_rdf("MyGraph", rdf_graph=Graph()) + + # Serialize to different formats + print(rdf_graph.serialize(format="turtle")) + rdf_graph.serialize("output.ttl", format="turtle") + rdf_graph.serialize("output.xml", format="xml") + + +Export Methods +============== + +1. Export by Graph Name +~~~~~~~~~~~~~~~~~~~~~~~ + +The simplest approach - exports all vertex and edge collections defined in the +ArangoDB graph: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + name="MyGraph", + rdf_graph=Graph() + ) + +2. Export by Collection Names +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Export specific vertex and edge collections: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_collections_to_rdf( + name="MyGraph", + rdf_graph=Graph(), + v_cols={"Person", "Company"}, # Vertex collections + e_cols={"knows", "worksAt"}, # Edge collections + ) + +3. Export with Metagraph +~~~~~~~~~~~~~~~~~~~~~~~~ + +Fine-grained control over which attributes to export from each collection: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_to_rdf( + name="MyGraph", + rdf_graph=Graph(), + metagraph={ + "vertexCollections": { + "Person": {"name", "age", "email"}, # Only these attributes + "Company": {"name", "founded"}, + }, + "edgeCollections": { + "knows": {"since"}, + "worksAt": {"role", "startDate"}, + }, + }, + explicit_metagraph=True, # Only include specified attributes + ) + +Set ``explicit_metagraph=False`` to include all attributes while still filtering +collections: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_to_rdf( + name="MyGraph", + rdf_graph=Graph(), + metagraph={ + "vertexCollections": { + "Person": set(), # Empty set = all attributes + "Company": set(), + }, + "edgeCollections": { + "knows": set(), + }, + }, + explicit_metagraph=False, # Include all attributes + ) + + +Conversion Options +================== + +List Conversion Mode +~~~~~~~~~~~~~~~~~~~~ + +Controls how ArangoDB arrays are converted to RDF: + ++---------------+------------------------------------------------------------------+ +| Mode | Description | ++===============+==================================================================+ +| ``static`` | Each array element becomes a separate triple (default) | ++---------------+------------------------------------------------------------------+ +| ``collection``| Uses RDF Collection structure (``rdf:first``, ``rdf:rest``) | ++---------------+------------------------------------------------------------------+ +| ``container`` | Uses RDF Container structure (``rdf:_1``, ``rdf:_2``, etc.) | ++---------------+------------------------------------------------------------------+ +| ``serialize`` | Serializes array as JSON string literal (best for round-tripping)| ++---------------+------------------------------------------------------------------+ + +.. code-block:: python + + # Example: Array [1, 2, 3] with different modes + + # static (default): Creates 3 separate triples + # :subject :predicate 1 . + # :subject :predicate 2 . + # :subject :predicate 3 . + + # collection: RDF Collection structure + # :subject :predicate _:b1 . + # _:b1 rdf:first 1 ; rdf:rest _:b2 . + # _:b2 rdf:first 2 ; rdf:rest _:b3 . + # _:b3 rdf:first 3 ; rdf:rest rdf:nil . + + # serialize: JSON string + # :subject :predicate "[1, 2, 3]" . + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + list_conversion_mode="collection", + ) + +Dict Conversion Mode +~~~~~~~~~~~~~~~~~~~~ + +Controls how nested ArangoDB objects are converted to RDF: + ++---------------+------------------------------------------------------------------+ +| Mode | Description | ++===============+==================================================================+ +| ``static`` | Creates BNode with properties for each key (default) | ++---------------+------------------------------------------------------------------+ +| ``serialize`` | Serializes object as JSON string literal (best for round-tripping)| ++---------------+------------------------------------------------------------------+ + +.. code-block:: python + + # Example: {"city": "NYC", "zip": "10001"} with different modes + + # static (default): Creates BNode structure + # :subject :address _:b1 . + # _:b1 :city "NYC" ; :zip "10001" . + + # serialize: JSON string + # :subject :address "{\"city\": \"NYC\", \"zip\": \"10001\"}" . + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + dict_conversion_mode="serialize", + ) + + +Round-Tripping Support +====================== + +ArangoRDF supports round-tripping (ArangoDB → RDF → ArangoDB) with special options +to preserve ArangoDB-specific information. + +Preserving Collection Names +~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Use ``include_adb_v_col_statements=True`` to generate ``adb:collection`` statements: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + include_adb_v_col_statements=True, + ) + + # Generates statements like: + # "Person" . + +Preserving Document Keys +~~~~~~~~~~~~~~~~~~~~~~~~ + +Use ``include_adb_v_key_statements=True`` to preserve vertex document keys: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + include_adb_v_key_statements=True, + ) + + # Generates statements like: + # "doc123" . + +Preserving Edge Keys +~~~~~~~~~~~~~~~~~~~~ + +Use ``include_adb_e_key_statements=True`` to preserve edge document keys. +Note: This imposes triple reification on all edges. + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + include_adb_e_key_statements=True, + ) + +Preserving Namespace Prefixes +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Store namespace prefixes in an ArangoDB collection for round-trip reconstruction: + +.. code-block:: python + + # When exporting from ArangoDB (after a previous RDF import with PGT) + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + namespace_collection_name="Namespaces", # Collection created during PGT import + ) + + +Inferring RDF Types +=================== + +Use ``infer_type_from_adb_v_col=True`` to generate ``rdf:type`` statements based +on the ArangoDB collection name: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + infer_type_from_adb_v_col=True, + ) + + # A document in the "Person" collection generates: + # rdf:type . + + +Ignoring Attributes +=================== + +Exclude specific attributes from the export: + +.. code-block:: python + + rdf_graph = adbrdf.arangodb_graph_to_rdf( + "MyGraph", + rdf_graph=Graph(), + metagraph={ + "vertexCollections": {"Person": set()}, + "edgeCollections": {"knows": set()}, + }, + explicit_metagraph=False, + ignored_attributes={"_internal_field", "created_at", "updated_at"}, + ) + +Note: ``ignored_attributes`` cannot be used when ``explicit_metagraph=True``. + + +Complete Example +================ + +.. code-block:: python + + from rdflib import Graph + from arango import ArangoClient + from arango_rdf import ArangoRDF + + # Connect to ArangoDB + db = ArangoClient(hosts="http://localhost:8529").db( + "_system", username="root", password="" + ) + adbrdf = ArangoRDF(db) + + # Export with all options for round-tripping + rdf_graph = adbrdf.arangodb_graph_to_rdf( + name="MyGraph", + rdf_graph=Graph(), + list_conversion_mode="serialize", # Best for round-tripping + dict_conversion_mode="serialize", # Best for round-tripping + include_adb_v_col_statements=True, # Preserve collection names + include_adb_v_key_statements=True, # Preserve vertex keys + ) + + # Print as Turtle + print(rdf_graph.serialize(format="turtle")) + + # Save to file + rdf_graph.serialize("export.ttl", format="turtle") + + # Get statistics + print(f"Exported {len(rdf_graph)} triples") + + +Parameter Reference +=================== + ++--------------------------------+-------------------+------------------------------------------------+ +| Parameter | Default | Description | ++================================+===================+================================================+ +| ``name`` | (required) | ArangoDB graph name | ++--------------------------------+-------------------+------------------------------------------------+ +| ``rdf_graph`` | (required) | Target rdflib Graph object | ++--------------------------------+-------------------+------------------------------------------------+ +| ``list_conversion_mode`` | ``"static"`` | How to convert arrays | ++--------------------------------+-------------------+------------------------------------------------+ +| ``dict_conversion_mode`` | ``"static"`` | How to convert nested objects | ++--------------------------------+-------------------+------------------------------------------------+ +| ``infer_type_from_adb_v_col`` | ``False`` | Generate rdf:type from collection name | ++--------------------------------+-------------------+------------------------------------------------+ +| ``include_adb_v_col_statements``| ``False`` | Include adb:collection statements | ++--------------------------------+-------------------+------------------------------------------------+ +| ``include_adb_v_key_statements``| ``False`` | Include adb:key for vertices | ++--------------------------------+-------------------+------------------------------------------------+ +| ``include_adb_e_key_statements``| ``False`` | Include adb:key for edges (uses reification) | ++--------------------------------+-------------------+------------------------------------------------+ +| ``namespace_collection_name`` | ``None`` | Collection storing namespace prefixes | ++--------------------------------+-------------------+------------------------------------------------+ +| ``ignored_attributes`` | ``None`` | Set of attributes to exclude | ++--------------------------------+-------------------+------------------------------------------------+ + diff --git a/docs/index.rst b/docs/index.rst index af03342..5185b63 100644 --- a/docs/index.rst +++ b/docs/index.rst @@ -78,4 +78,6 @@ Contents rdf_to_arangodb_rpt rdf_to_arangodb_pgt rdf_to_arangodb_lpg + rdf_to_arangodb_concurrent + arangodb_to_rdf specs diff --git a/docs/rdf_to_arangodb_concurrent.rst b/docs/rdf_to_arangodb_concurrent.rst new file mode 100644 index 0000000..a0064f6 --- /dev/null +++ b/docs/rdf_to_arangodb_concurrent.rst @@ -0,0 +1,214 @@ +Concurrent Imports +------------------ + +ArangoRDF supports concurrent imports using Python's ``concurrent.futures`` module, +allowing you to import multiple RDF graphs into the same ArangoDB graph in parallel. + +This can significantly speed up imports when you have multiple RDF files to process. + +Basic Usage +=========== + +.. code-block:: python + + from concurrent.futures import ThreadPoolExecutor, as_completed + from rdflib import Graph, Namespace, Literal + from arango import ArangoClient + from arango_rdf import ArangoRDF + + db = ArangoClient().db() + + # Create multiple RDF graphs + EX = Namespace("http://example.org/") + + g1 = Graph() + g1.add((EX.Alice, EX.knows, EX.Bob)) + g1.add((EX.Alice, EX.name, Literal("Alice"))) + + g2 = Graph() + g2.add((EX.Bob, EX.knows, EX.Charlie)) + g2.add((EX.Bob, EX.name, Literal("Bob"))) + + graphs = [g1, g2] + + def import_rdf(rdf_graph: Graph) -> None: + # Each thread MUST create its own ArangoRDF instance + adbrdf = ArangoRDF(db, enable_rich=False) + adbrdf.rdf_to_arangodb_by_pgt( + "MyGraph", + rdf_graph, + overwrite_graph=False, + flatten_reified_triples=False, + overwrite_mode="ignore", + raise_on_document_error=False, + ) + + with ThreadPoolExecutor(max_workers=4) as executor: + futures = [executor.submit(import_rdf, g) for g in graphs] + for future in as_completed(futures): + future.result() # Raises exception if import failed + + +Requirements & Limitations +========================== + +There are several important requirements and limitations to be aware of when +using concurrent imports: + +1. Disable Rich Progress Bars (``enable_rich=False``) +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Rich progress bars interfere with concurrent/multiprocessing modules. Always +set ``enable_rich=False`` when creating ``ArangoRDF`` instances for concurrent use: + +.. code-block:: python + + adbrdf = ArangoRDF(db, enable_rich=False) + +2. Create Separate ArangoRDF Instances Per Thread +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +Each thread **must** create its own ``ArangoRDF`` instance. Do not share a single +instance across threads, as this will cause race conditions and errors. + +.. code-block:: python + + # CORRECT: Create instance inside the thread function + def import_rdf(rdf_graph): + adbrdf = ArangoRDF(db, enable_rich=False) # New instance per thread + adbrdf.rdf_to_arangodb_by_pgt(...) + + # WRONG: Sharing instance across threads + adbrdf = ArangoRDF(db, enable_rich=False) + def import_rdf(rdf_graph): + adbrdf.rdf_to_arangodb_by_pgt(...) # Shared instance - will fail! + +3. Disable Triple Reification (``flatten_reified_triples=False``) +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +The triple reification flattening process uses SPARQL queries internally, which +rely on ``pyparsing`` - a library that is **not thread-safe**. When using concurrent +imports, you must disable this feature: + +.. code-block:: python + + adbrdf.rdf_to_arangodb_by_pgt( + ..., + flatten_reified_triples=False, # Required for thread safety + ) + +If your RDF data contains reified triples and you need to flatten them, you must +process those graphs sequentially (not concurrently). + +4. Handle Write-Write Conflicts +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +When multiple threads insert data into the same ArangoDB graph concurrently, +they may attempt to insert the same document (e.g., a shared predicate like +``ex:name``). This causes write-write conflicts. + +To handle this gracefully, use: + +.. code-block:: python + + adbrdf.rdf_to_arangodb_by_pgt( + ..., + overwrite_mode="ignore", # Skip documents that already exist + raise_on_document_error=False, # Don't raise on conflicts + ) + +This tells ArangoDB to silently skip duplicate documents rather than failing. + +5. Don't Overwrite the Graph (``overwrite_graph=False``) +~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~~ + +When importing to the same graph from multiple threads, ensure you don't +overwrite the graph on each import: + +.. code-block:: python + + adbrdf.rdf_to_arangodb_by_pgt( + "MyGraph", + rdf_graph, + overwrite_graph=False, # Append to existing graph + ) + + +Complete Example +================ + +Here's a complete example that follows all the requirements: + +.. code-block:: python + + from concurrent.futures import ThreadPoolExecutor, as_completed + from rdflib import Graph, Namespace, Literal + from arango import ArangoClient + from arango_rdf import ArangoRDF + + # Setup + db = ArangoClient().db() + graph_name = "ConcurrentImportExample" + + # Clean up any existing graph + if db.has_graph(graph_name): + db.delete_graph(graph_name, drop_collections=True) + + # Create sample RDF graphs + EX = Namespace("http://example.org/") + + graphs = [] + for i in range(10): + g = Graph() + person = EX[f"Person{i}"] + g.add((person, EX.name, Literal(f"Person {i}"))) + g.add((person, EX.knows, EX.Person0)) + graphs.append(g) + + def import_rdf(rdf_graph: Graph) -> str: + """Import a single RDF graph - called from each thread.""" + # Create a new ArangoRDF instance for this thread + adbrdf = ArangoRDF(db, enable_rich=False) + + adbrdf.rdf_to_arangodb_by_pgt( + graph_name, + rdf_graph, + overwrite_graph=False, + flatten_reified_triples=False, + overwrite_mode="ignore", + raise_on_document_error=False, + ) + return "success" + + # Import all graphs concurrently + results = [] + with ThreadPoolExecutor(max_workers=4) as executor: + futures = [executor.submit(import_rdf, g) for g in graphs] + for future in as_completed(futures): + try: + results.append(future.result()) + except Exception as e: + print(f"Import failed: {e}") + + print(f"Successfully imported {len(results)} graphs") + + +Summary of Required Parameters +============================== + +When using concurrent imports, always use these parameters: + ++-------------------------------+-------------------+----------------------------------------+ +| Parameter | Value | Reason | ++===============================+===================+========================================+ +| ``enable_rich`` | ``False`` | Rich interferes with concurrency | ++-------------------------------+-------------------+----------------------------------------+ +| ``flatten_reified_triples`` | ``False`` | SPARQL parser is not thread-safe | ++-------------------------------+-------------------+----------------------------------------+ +| ``overwrite_graph`` | ``False`` | Append to graph, don't recreate | ++-------------------------------+-------------------+----------------------------------------+ +| ``overwrite_mode`` | ``"ignore"`` | Skip duplicate documents | ++-------------------------------+-------------------+----------------------------------------+ +| ``raise_on_document_error`` | ``False`` | Don't fail on write conflicts | ++-------------------------------+-------------------+----------------------------------------+ + diff --git a/tests/test_main.py b/tests/test_main.py index 32684fc..a075cfb 100644 --- a/tests/test_main.py +++ b/tests/test_main.py @@ -1,4 +1,5 @@ import json +from concurrent.futures import ThreadPoolExecutor, as_completed from typing import Any, Dict, List import pytest @@ -6,7 +7,7 @@ from rdflib import RDF, RDFS, BNode from rdflib import ConjunctiveGraph as RDFConjunctiveGraph from rdflib import Graph as RDFGraph -from rdflib import Literal, URIRef +from rdflib import Literal, Namespace, URIRef from arango_rdf import ArangoRDF from arango_rdf.exception import ArangoRDFImportException @@ -5655,3 +5656,49 @@ def test_lpg_case_12_1(name: str, rdf_graph: RDFGraph) -> None: assert edge["_to"].split("/")[0] in {"Class", "Node"} db.delete_graph("Test", drop_collections=True) + + +def test_pgt_concurrent() -> None: + db.delete_graph("Test", drop_collections=True, ignore_missing=True) + + EX = Namespace("http://example.org/") + + g1 = RDFGraph() + g1.add((EX.Alice, EX.knows, EX.Bob)) + g1.add((EX.Alice, EX.name, Literal("Alice"))) + + g2 = RDFGraph() + g2.add((EX.Bob, EX.knows, EX.Charlie)) + g2.add((EX.Bob, EX.name, Literal("Bob"))) + + graph_specs = [("Test", g1), ("Test", g2)] + results = [] + + def import_rdf(graph_name: str, rdf_graph: RDFGraph) -> str: + # Disable rich progress bars to avoid interference with concurrent modules + adbrdf = ArangoRDF(db, enable_rich=False) + adbrdf.rdf_to_arangodb_by_pgt( + graph_name, + rdf_graph, + overwrite_graph=False, + # Triple Reification is **not** thread-safe due to SPARQL queries + flatten_reified_triples=False, + resource_collection_name="Node", + # For concurrent inserts: ignore duplicates, don't raise on conflicts + overwrite_mode="ignore", + raise_on_document_error=False, + ) + + return graph_name + + with ThreadPoolExecutor(max_workers=2) as executor: + futures = [ + executor.submit(import_rdf, name, graph) for name, graph in graph_specs + ] + for future in as_completed(futures): + results.append(future.result()) + + assert db.has_graph("Test") + assert db.collection("Node").count() == 3 + assert db.collection("Property").count() == 2 + assert db.collection("knows").count() == 2