Aller au contenu

Gremlin

ddigraph prend en charge les bases de données graphe compatibles Gremlin via le framework Apache TinkerPop.

Bases de données prises en charge

Base de données Connexion Cas d'utilisation
Apache TinkerGraph En mémoire Tests locaux, développement
JanusGraph WebSocket/HTTP Production, distribué
Amazon Neptune WebSocket Cloud AWS, service géré
Azure Cosmos DB WebSocket Cloud Azure, API Gremlin

Dépendances

GremlinPython est un extra optionnel :

pip install "ddigraph[gremlin]"

Utilisation de base

Connexion au serveur Gremlin

from gremlin_python.process.anonymous_traversal import traversal
from gremlin_python.driver.driver_remote_connection import DriverRemoteConnection

# Serveur TinkerPop local
connection = DriverRemoteConnection("ws://localhost:8182/gremlin", "g")
g = traversal().withRemote(connection)

# AWS Neptune
# connection = DriverRemoteConnection(
#     'wss://your-neptune-endpoint:8182/gremlin',
#     'g'
# )

Charger des données DDI

from ddigraph import iter_graph


def node_key(node):
    """Une clé stable tirée de toute l'identité, pas du premier champ."""
    return "|".join(f"{k}={v}" for k, v in sorted(node.identity.items()))


vertices = {}
edges = []

# Première passe : tous les nœuds arrivent avant la première relation.
# On collecte donc les sommets avant de câbler quoi que ce soit.
for chunk in iter_graph("survey.xml"):
    for node in chunk.nodes:
        vertices[node_key(node)] = (
            g.addV(node.label)
            .property("id", node_key(node))
            .property("name", node.properties.get("label", ""))
            .next()
        )
    edges.extend(
        (node_key(edge.start), edge.type, node_key(edge.end)) for edge in chunk.relationships
    )

# Seconde passe : les extrémités sont désormais toutes dans `vertices`.
for start, rel_type, end in edges:
    if start in vertices and end in vertices:
        g.V(vertices[start]).addE(rel_type).to(vertices[end]).iterate()

connection.close()

Exemple complet

Consultez demo/load_gremlin.py pour un exemple complet :

"""Load DDI into Gremlin-compatible graph database."""

from gremlin_python.process.anonymous_traversal import traversal
from gremlin_python.driver.driver_remote_connection import DriverRemoteConnection
from gremlin_python.process.graph_traversal import __

from ddigraph import iter_graph


def load_ddi_to_gremlin(ddi_path: str, gremlin_endpoint: str = "ws://localhost:8182/gremlin"):
    """Parse DDI-L file and load into Gremlin database."""

    connection = DriverRemoteConnection(gremlin_endpoint, "g")
    g = traversal().withRemote(connection)

    try:
        # Clear existing data (optional)
        g.V().drop().iterate()

        fragment_ids = set()

        # First pass: create all vertices
        for chunk in iter_graph(ddi_path):
            for fragment in chunk.nodes:
                props = fragment.to_dict()

                vertex = g.addV(fragment.element_type)
                vertex = vertex.property("fragment_id", fragment.fragment_id)

                if fragment.label:
                    vertex = vertex.property("label", fragment.label)
                if fragment.urn:
                    vertex = vertex.property("urn", fragment.urn)
                if fragment.agency:
                    vertex = vertex.property("agency", fragment.agency)
                if fragment.version:
                    vertex = vertex.property("version", fragment.version)

                # Add type-specific properties
                if fragment.element_type == "QuestionItem":
                    if props.get("question_text"):
                        vertex = vertex.property("question_text", props["question_text"])
                elif fragment.element_type == "Category":
                    if props.get("category_label"):
                        vertex = vertex.property("category_label", props["category_label"])

                vertex.next()
                fragment_ids.add(fragment.fragment_id)

        print(f"Created {len(fragment_ids)} vertices")

        # Second pass: create edges
        edge_count = 0
        for chunk in iter_graph(ddi_path):
            for fragment in chunk.nodes:
                for rel_type, ref in fragment.references:
                    if ref.id in fragment_ids:
                        g.V().has("fragment_id", fragment.fragment_id).addE(rel_type).to(
                            __.V().has("fragment_id", ref.id)
                        ).iterate()
                        edge_count += 1

        print(f"Created {edge_count} edges")

    finally:
        connection.close()


def query_examples(gremlin_endpoint: str = "ws://localhost:8182/gremlin"):
    """Example Gremlin queries."""

    connection = DriverRemoteConnection(gremlin_endpoint, "g")
    g = traversal().withRemote(connection)

    try:
        # Count vertices by type
        counts = g.V().groupCount().by(__.label()).next()
        print("Vertex counts by type:")
        for label, count in counts.items():
            print(f"  {label}: {count}")

        # Find all questions
        questions = g.V().hasLabel("QuestionItem").valueMap(True).toList()
        print(f"\nFound {len(questions)} questions")

        # Traverse from instrument to questions
        paths = (
            g.V()
            .hasLabel("Instrument")
            .repeat(__.out("HAS_CONSTRUCT"))
            .until(__.hasLabel("QuestionConstruct"))
            .out("ASKS_QUESTION")
            .path()
            .toList()
        )

        print(f"\nFound {len(paths)} paths from Instrument to Question")

    finally:
        connection.close()


if __name__ == "__main__":
    load_ddi_to_gremlin("data/Ireland_LabourSurvey.xml")
    query_examples()

Requêtes Gremlin

Traversées de base

// Compter les sommets par libellé
g.V().groupCount().by(label)

// Trouver tous les QuestionItems
g.V().hasLabel('QuestionItem').valueMap(true)

// Obtenir un fragment spécifique par ID
g.V().has('fragment_id', 'abc-123').valueMap(true)

Traversées de relations

// De l'Instrument à tous les constructs
g.V().hasLabel('Instrument').out('HAS_CONSTRUCT').valueMap(true)

// Questions avec leurs listes de codes
g.V().hasLabel('QuestionItem')
    .as('q')
    .out('USES_CODELIST')
    .as('cl')
    .select('q', 'cl')
    .by(valueMap('label', 'question_text'))
    .by(valueMap('label'))

// Catégories dans une liste de codes
g.V().has('fragment_id', 'codelist-123')
    .out('HAS_CATEGORY')
    .valueMap('category_label')

Requêtes de chemins

// Chemin complet de l'Instrument aux Questions
g.V().hasLabel('Instrument')
    .repeat(out('HAS_CONSTRUCT'))
    .until(hasLabel('QuestionConstruct'))
    .out('ASKS_QUESTION')
    .path()
    .by('label')

// Branchements conditionnels
g.V().hasLabel('IfThenElse')
    .project('condition', 'then', 'else')
    .by('condition')
    .by(out('THEN').values('label'))
    .by(out('ELSE').values('label'))

Configuration spécifique aux bases de données

JanusGraph

from gremlin_python.driver.serializer import GraphSONSerializersV3d0

connection = DriverRemoteConnection(
    "ws://localhost:8182/gremlin", "g", message_serializer=GraphSONSerializersV3d0()
)

Amazon Neptune

from gremlin_python.driver import client

# Authentification IAM
import boto3
from botocore.auth import SigV4Auth
from botocore.awsrequest import AWSRequest

# Connexion WebSocket avec SigV4
connection = DriverRemoteConnection(
    "wss://your-cluster.region.neptune.amazonaws.com:8182/gremlin", "g"
)

Azure Cosmos DB

from gremlin_python.driver import client, serializer

# Cosmos DB nécessite un sérialiseur spécifique
connection = DriverRemoteConnection(
    "wss://your-account.gremlin.cosmos.azure.com:443/",
    "g",
    username="/dbs/your-database/colls/your-graph",
    password="your-primary-key",
)

Chargement par lots

Pour les fichiers DDI volumineux, utilisez des écritures par lots :

BATCH_SIZE = 100

vertices_batch = []
for i, fragment in enumerate(
    node for chunk in iter_graph("large_survey.xml") for node in chunk.nodes
):
    vertices_batch.append(fragment)

    if len(vertices_batch) >= BATCH_SIZE:
        # Soumettre le lot
        for f in vertices_batch:
            g.addV(f.element_type).property("fragment_id", f.fragment_id).property(
                "label", f.label or ""
            ).iterate()
        vertices_batch = []
        print(f"Processed {i + 1} vertices")

# Dernier lot
for f in vertices_batch:
    g.addV(f.element_type).property("fragment_id", f.fragment_id).property(
        "label", f.label or ""
    ).iterate()

Voir aussi