Skip to content

Example 05: CSV Import - Graph

View source code

Production-ready graph creation from MovieLens dataset with comprehensive performance analysis

Overview

This example demonstrates creating a graph database from the MovieLens dataset, transforming documents into vertices and edges. You'll learn production-ready patterns for:

  • Graph modeling - Users and Movies as vertices, ratings and tags as edges
  • SQL pipeline - End-to-end graph creation with SQL DDL/DML
  • GraphBatch vs synchronous transactions - Compare the two Java-API vertex paths
  • Index optimization - Create indexes BEFORE bulk edge creation
  • Export & roundtrip validation - Verify data integrity through complete cycle
  • Performance benchmarking - Measure and compare 6 different configurations
  • Query validation - 10 graph queries with result verification
  • Directed graph semantics - Traversals should be interpreted as directional

What You'll Learn

  • Graph schema design (vertex types, edge types, properties)
  • Foreign key resolution (userId → User vertex, movieId → Movie vertex)
  • Bulk vertex creation through db.graph_batch(...), which crosses the Python/Java boundary once per batch
  • Why the key indexes come before edge creation: each edge batch looks up its endpoints by key
  • Export and roundtrip import of the graph
  • Graph query patterns (MATCH, collaborative filtering, recommendations)
  • Production import patterns for large-scale graphs

Prerequisites

1. Install ArcadeDB Python bindings:

pip install arcadedb-embedded

2. Source database from Example 04:

The example reads Example 04's document database (./my_test_databases/movielens_<size>_db, or --source-db), or imports one from an Example 04 JSONL export with --import-jsonl. It does not download anything: if the source database is missing it prints the Example 04 command to run and exits.

Two dataset sizes available:

  • movielens-large: ~330K users, ~86K movies, ~33M edges (~971 MB CSV) - Realistic performance testing
  • movielens-small: ~610 users, ~9,700 movies, ~100K edges (~3.2 MB CSV) - Quick testing

Usage

# Recommended SQL configuration (synchronous)
python 05_csv_import_graph.py --dataset movielens-small --method sql

# Java API with GraphBatch vertices, same graph as --method sql
python 05_csv_import_graph.py --dataset movielens-small --method java

# Run with export for roundtrip validation
python 05_csv_import_graph.py --dataset movielens-small --batch-size 5000 --method java --no-async --export

# See all options
python 05_csv_import_graph.py --help

Key options:

  • --dataset {movielens-small,movielens-large} - Dataset size (default: movielens-small)
  • --method {java,sql} - Creation method (DSL-first recommendation: sql)
  • --no-async - Build vertices in synchronous transactions instead of GraphBatch
  • --no-index - Skip creating indexes (slower, for comparison)
  • --batch-size BATCH_SIZE - Batch size for operations (default: 5000)
  • --parallel PARALLEL - Above 1, enable GraphBatch parallel flush (default: 4)
  • --export - Export graph database to JSONL after creation
  • --export-path EXPORT_PATH - Export filename (default: {db_name}.jsonl.tgz)
  • --db-name DB_NAME - Database name (default derived from dataset, e.g. movielens_graph_small_db)
  • --source-db SOURCE_DB - Custom source database path
  • --import-jsonl IMPORT_JSONL - Import from JSONL export instead of using the source document DB

Recommendations:

  • Method: Use sql for consistency with DSL-first examples and guides
  • Vertex path: With --method java, leave GraphBatch on. It is the repository's recommended bulk graph path, and --no-async exists to compare it against synchronous transactions. --method sql always builds vertices in synchronous transactions
  • Indexes: Keep enabled; each edge batch looks up its endpoints by userId and movieId
  • Batch size: 5000 for small, 50000 for large datasets
  • Export: Use --export to validate data integrity via roundtrip

All three of --method sql, --method java, and --method java --no-async produce the same graph: 610 Users, 9,742 Movies, 97,823 RATED, and 3,436 TAGGED on movielens-small.

Graph Schema

Vertex Types

User

  • Properties: userId (INTEGER, indexed)
  • Count: 610 (small) / 330,000 (large)

Movie

  • Properties: movieId (INTEGER, indexed), title (STRING), genres (STRING), imdbId (STRING), tmdbId (INTEGER)
  • Count: 9,742 (small) / 86,537 (large)

Edge Types

RATED (User → Movie)

  • Properties: rating (FLOAT), timestamp (LONG)
  • Count: 97,823 (small) / 33,155,309 (large)

TAGGED (User → Movie)

  • Properties: tag (STRING), timestamp (LONG)
  • Count: 3,436 (small) / 2,212,213 (large)

The small-dataset counts are below the source row counts (100,836 ratings, 3,683 tags) on purpose. download_data.py injects NULLs into the MovieLens CSVs, and this example only builds edges from rows that have the properties the edge declares: timestamp IS NOT NULL for RATED, plus tag IS NOT NULL for TAGGED. Example 04 keeps all the rows, NULLs included, because importing NULLs is what that example demonstrates.

Graph Queries

The example includes 10 comprehensive graph queries (8 SQL + 2 OpenCypher), run by run_and_validate_queries() and validated against the per-dataset baselines in EXPECTED_RESULTS:

Query 1: Movies rated by User #1 (SQL - Basic Traversal)

SELECT expand(out('RATED')) FROM User WHERE userId = 1

Query 2: Movies rated 5.0 by User #1 (SQL - Edge Property Filter)

SELECT expand(outE('RATED')[rating = 5.0].inV())
FROM User WHERE userId = 1

Query 3: Rating statistics for top 5 active users (SQL - Aggregations)

SELECT u.userId as userId,
       COUNT(e) as num_ratings,
       AVG(e.rating) as avg_rating,
       MIN(e.rating) as min_rating,
       MAX(e.rating) as max_rating
FROM (
  MATCH {type: User, as: u}.outE('RATED'){as: e} RETURN u,e
)
GROUP BY u.userId
ORDER BY num_ratings DESC
LIMIT 5

Query 4: Top 10 most rated movies (SQL - Aggregations)

SELECT m.movieId as movieId, m.title as title,
       COUNT(e) as num_ratings, AVG(e.rating) as avg_rating
FROM (
  MATCH {type: Movie, as: m}.inE('RATED'){as: e} RETURN m, e
)
GROUP BY m.movieId, m.title
ORDER BY num_ratings DESC
LIMIT 10

Query 5: Top 10 most tagged movies (SQL - Aggregations)

SELECT m.movieId as movieId, m.title as title, COUNT(e) as num_tags
FROM (
  MATCH {type: Movie, as: m}.inE('TAGGED'){as: e} RETURN m, e
)
GROUP BY m.movieId, m.title
ORDER BY num_tags DESC
LIMIT 10

Query 6: Users who rated same movies as User #1 (SQL - MATCH Pattern)

SELECT friend.userId as other_user, movie.title as common_movie,
       a.rating as my_rating, b.rating as their_rating
FROM (
  MATCH {type: User, where: (userId = 1), as: me}
        .outE('RATED'){as: a}
        .inV(){as: movie}
        .inE('RATED'){as: b}
        .outV(){as: friend, where: (userId != 1)}
  RETURN me, friend, movie, a, b
)

Query 7: Users with similar taste to User #1 (SQL - MATCH + Aggregation)

Same as Query 6 but filters both edges to rating >= 4.5 and aggregates count(*) as shared_high_ratings grouped by friend.userId.

Query 8: Rating distribution across all ratings (SQL - Aggregation)

SELECT rating, count(*) as frequency
FROM RATED
WHERE rating IS NOT NULL
GROUP BY rating
ORDER BY rating

Query 9: User #1's top-rated movies (OpenCypher - Basic Pattern)

MATCH (u:User {userId: 1})-[r:RATED]->(m:Movie)
WHERE r.rating >= 4.0
RETURN m.title as title, r.rating as rating
ORDER BY rating DESC

Query 10: Users who rated same movies as User #1 (OpenCypher - Pattern)

MATCH (u:User {userId: 1})-[:RATED]->(m:Movie)<-[:RATED]-(other:User)
WHERE other.userId <> 1
RETURN other.userId as other_user, count(*) as shared_movies
ORDER BY shared_movies DESC

Note: Queries 1-8 use SQL (including SQL MATCH patterns); Queries 9-10 use OpenCypher as an alternative pattern syntax.

Code Walkthrough

Step 1: Create Graph Schema

# create_schema() runs each DDL statement wrapped in try/except so the
# schema build is idempotent (re-running a CREATE is a no-op):
schema_commands = [
    "CREATE VERTEX TYPE User",
    "CREATE VERTEX TYPE Movie",
    "CREATE EDGE TYPE RATED",
    "CREATE EDGE TYPE TAGGED",
    "CREATE PROPERTY User.userId INTEGER",
    "CREATE PROPERTY Movie.movieId INTEGER",
    "CREATE PROPERTY Movie.title STRING",
    "CREATE PROPERTY Movie.genres STRING",
    "CREATE PROPERTY Movie.imdbId STRING",
    "CREATE PROPERTY Movie.tmdbId INTEGER",
    "CREATE PROPERTY RATED.rating FLOAT",
    "CREATE PROPERTY RATED.timestamp LONG",
    "CREATE PROPERTY TAGGED.tag STRING",
    "CREATE PROPERTY TAGGED.timestamp LONG",
]
for command in schema_commands:
    try:
        db.command("sql", command)
    except Exception:
        pass

# Indexes (unless --no-index): created BEFORE bulk edge creation
db.command("sql", "CREATE INDEX ON User (userId) UNIQUE_HASH")
db.command("sql", "CREATE INDEX ON Movie (movieId) UNIQUE_HASH")

Because the graph is directed, user-to-movie traversals should normally be expressed as outgoing traversals from User to Movie. If you need incoming traversals, express that explicitly in SQL MATCH or OpenCypher.

Step 2: Create Vertices

# The example uses a VertexCreator class for batch creation. Users come from
# the distinct userIds in the source Rating type. The synchronous SQL path
# inserts user vertices in batched transactions:
class VertexCreator:
    def _create_users(self, total_users: int):
        """Create User vertices from distinct Rating.userId values."""
        with arcadedb.open_database(str(source_db_path)) as source_db:
            query = "SELECT userId FROM Rating GROUP BY userId ORDER BY userId"
            batch_user_ids = []
            for record in source_db.query("sql", query):
                batch_user_ids.append(record.get("userId"))
                if len(batch_user_ids) >= self.batch_size:
                    with self.db.transaction():
                        for uid in batch_user_ids:
                            self.db.command(
                                "sql", "INSERT INTO User SET userId = ?", uid
                            )
                    batch_user_ids = []
            # ... handle remaining batch ...

# See full implementation in the Python file for the Java-API paths.

The --method java path builds the same vertices through GraphBatch, which hands each batch to Java in one call and returns one RID per row:

with self.db.graph_batch(parallel_flush=self.parallel_level > 1) as batch:
    created = batch.create_vertices("User", rows)
    if len(created) != len(rows):
        raise RuntimeError(
            f"GraphBatch returned {len(created)} RIDs for {len(rows)} User rows"
        )

The RID count is checked against the row count on every flush, so a short write raises instead of producing a quietly incomplete graph.

Step 3: Create Edges with Foreign Key Resolution

# EdgeCreator paginates the source Rating type (database-level streaming),
# resolves User/Movie RIDs once per batch via a cache (`IN :ids` lookups, with
# the ids bound as one list parameter so every chunk reuses the same query text
# and the userId/movieId indexes), then creates edges directly between RIDs.
class EdgeCreator:
    def _create_rated_edges(self, total_ratings: int):
        """Create RATED edges from Rating records."""
        last_rid = "#-1:-1"
        while True:
            query = f"""
                SELECT *, @rid as rid FROM Rating
                WHERE timestamp IS NOT NULL AND @rid > {last_rid}
                LIMIT {self.batch_size}
            """
            chunk = list(source_db.query("sql", query))
            if not chunk:
                break

            with self.db.transaction():
                # Build {userId: rid} and {movieId: rid} caches for this batch
                user_cache, movie_cache = self._build_vertex_cache(chunk)
                for record in chunk:
                    user_rid = user_cache.get(record.get("userId"))
                    movie_rid = movie_cache.get(record.get("movieId"))
                    if user_rid and movie_rid:
                        # RIDs and values are bound as ? parameters, so every
                        # edge runs the same statement text
                        self.db.command(
                            "sql",
                            "CREATE EDGE RATED FROM ? TO ? SET rating = ?, timestamp = ?",
                            user_rid,
                            movie_rid,
                            record.get("rating"),
                            record.get("timestamp"),
                        )

            last_rid = chunk[-1].get("rid")

# See full implementation in the Python file (Java-API and SQL variants).

Step 4: Run Validation Queries

# Verify vertex counts (validate_counts_and_samples)
user_count = list(db.query("sql", "SELECT count(*) as count FROM User"))[0].get("count")
movie_count = list(db.query("sql", "SELECT count(*) as count FROM Movie"))[0].get("count")

# Verify edge counts
rated_count = list(db.query("sql", "SELECT count(*) as count FROM RATED"))[0].get("count")
tagged_count = list(db.query("sql", "SELECT count(*) as count FROM TAGGED"))[0].get("count")

print(f"  Users:  {user_count:,}")
print(f"  Movies: {movie_count:,}")
print(f"  RATED:  {rated_count:,}")
print(f"  TAGGED: {tagged_count:,}")

Running the Benchmark

Single Configuration

# Java API with GraphBatch vertices
python 05_csv_import_graph.py --dataset movielens-small --method java

# Java API with synchronous vertex transactions, for comparison
python 05_csv_import_graph.py --dataset movielens-small --method java --no-async

# Compare with SQL
python 05_csv_import_graph.py --dataset movielens-small --method sql

# With export for roundtrip validation
python 05_csv_import_graph.py --dataset movielens-small --batch-size 5000 --method java --no-async --export

Comprehensive Benchmark (All 6 Configurations)

Each configuration is one run of 05_csv_import_graph.py with the flags shown:

  1. java (--method java) - Java API with GraphBatch vertices and indexes
  2. java_noasync (--method java --no-async) - Java API with synchronous vertex transactions and indexes
  3. java_noindex (--method java --no-index) - Java API with GraphBatch vertices, no indexes
  4. java_noindex_noasync (--method java --no-index --no-async) - Java API with synchronous vertex transactions, no indexes
  5. sql (--method sql) - SQL with indexes (always synchronous)
  6. sql_noindex (--method sql --no-index) - SQL without indexes (always synchronous)

Configurations 1 and 3 build vertices with GraphBatch and have not been timed yet, so the tables above have no rows for them.

Benchmark Configuration

JVM Settings

Prefer configuring heap inside the script before it creates the first database:

import arcadedb_embedded as arcadedb
from arcadedb_embedded.jvm import start_jvm

start_jvm(heap_size="8g", jvm_args="-Xms8g")

Memory Planning:

  • Heap size: 8 GB (-Xmx8g)
  • Total process memory: 15-19 GB (heap + non-heap)
  • Plan for 2× heap size in actual RAM

Batch Sizes

--batch-size 5000    # Small dataset (default)
--batch-size 50000   # Large dataset (recommended)

Larger batches = fewer commits = faster imports

Parallel Processing

--parallel 4         # Above 1, enable GraphBatch parallel flush (default 4)

Note: This affects only the Java-API GraphBatch vertex path. The value is not a parallel level: the script passes parallel_flush=args.parallel > 1 to db.graph_batch(...), so any value above 1 has the same effect. It has no effect with --no-async or with --method sql, both of which build vertices in synchronous transactions. GraphBatch lands every row with parallel_flush on or off.

Next Steps

Example 06 - Vector Search: Movie Recommendations

Semantic similarity search with MovieLens data:

  • Generate embeddings from movie titles/genres
  • Build HNSW (JVector) index for nearest-neighbor search
  • Find similar movies using cosine distance
  • Compare vector similarity with graph collaborative filtering on the rating data