Skip to content

How do I load a large graph into TigerGraph quickly?

You load a large graph into TigerGraph. By default GraFlo sends the records to TigerGraph's REST endpoint batch by batch, and for a large first load that overhead dominates. TigerGraph has a faster path of its own: a loading job that reads CSV files.

With bulk loading switched on, GraFlo writes the cast records to CSV files, uploads them to S3 (or an S3-compatible store such as MinIO), and runs one loading job at the end that reads them. The manifest does not change, and the S3 credentials stay out of it.

flowchart LR
    csv[relations.csv] --> cast[GraFlo casts rows]
    cast --> staged["CSV files<br/>in bulk_staging/"]
    staged -- upload --> s3[(MinIO / S3)]
    s3 -- LOADING JOB --> tg[(TigerGraph)]

The graph is the one of example 03: companies, and several kinds of relation between them read from a column.

What you need

  • GraFlo installed (pip install graflo).
  • A running TigerGraph and a running MinIO. The repository ships containers for both; see docker/README.md. TigerGraph must be able to reach MinIO, because it reads the uploaded files itself.

The scripts read connection settings from the environment files of those containers, not from your shell:

Setting Read from Variables
TigerGraph docker/tigergraph TG_WEB (GSQL port), TIGERGRAPH_HOSTNAME, TIGERGRAPH_USERNAME, TIGERGRAPH_PASSWORD (if unset there: the same variable in your shell, then tigergraph)
MinIO address, used by GraFlo docker/minio MINIO_ENDPOINT, or MINIO_HOSTNAME and MINIO_API_PORT
MinIO address, used by TigerGraph docker/minio MINIO_LOADER_ENDPOINT; without it, TigerGraph gets the address GraFlo uses
MinIO credentials docker/minio MINIO_ROOT_USER, MINIO_ROOT_PASSWORD
Bucket docker/minio MINIO_STAGING_BUCKET (default graflo-staging)

Three variables are read from your shell: BULK_USE_S3 (default 1; 0 keeps the files on local disk), BULK_S3_BUCKET (overrides the bucket) and BULK_S3_PREFIX (folder in the bucket, default demo).

The data

data/relations.csv:

company_a company_b relation date
Acme Beta partners 2024-01-01
Gamma Acme supplies 2024-02-01

Steps

1. Name the staging area in the manifest

manifest.yaml adds one entry to bindings: a staging area called bulk_s3, whose S3 settings are registered at run time under the label minio_bulk.

bindings:
    # ... connectors
    staging_proxy:
    -   name: bulk_s3
        conn_proxy: minio_bulk

2. Switch the target to bulk loading

Bulk loading is a setting of the TigerGraph connection. ingest.py sets it and registers the MinIO settings under minio_bulk:

minio_conf = minio_config()
conn_conf.bulk_load = TigergraphBulkLoadConfig(
    enabled=True,
    staging_dir=str(STAGING_DIR),
    s3_staging_name="bulk_s3",
    s3_bucket=minio_conf.bucket,
    s3_key_prefix=os.environ.get("BULK_S3_PREFIX", "demo"),
)
provider.register_generalized_config(
    conn_proxy="minio_bulk",
    config=minio_conf.to_s3_generalized_conn_config(),
)

Before ingesting, the script creates the bucket if it is missing, and fails at once if MinIO cannot be reached.

3. Start MinIO

From the repository root:

cd docker/minio
docker compose --env-file .env --profile graflo.minio up -d
cd ../..

4. Load and check

uv run python examples/13-tigergraph-bulk-s3/ingest.py
uv run python examples/13-tigergraph-bulk-s3/inspect_bulk.py

The scripts can run from any directory. inspect_bulk.py reports what was staged and what TigerGraph holds, and exits with an error when either is missing. With --staged-only it checks the files without connecting to TigerGraph.

What you should see

  • Staged, in the newest bulk_staging/<session>/: company.csv with 4 rows (one per company in each row of the CSV, so Acme twice) and edge_relates.csv with 2 rows. TigerGraph keeps both kinds of relation in one edge type, relates, with the kind in its relation attribute.
  • Loaded: 3 company vertices (Acme, Beta, Gamma; Acme is in both rows and is stored once) and 2 relates edges, both touching Acme.

Each run writes a new session directory; bulk_staging/ is not tracked by git.

What goes wrong

The ingest reports success and the graph is empty. The loading job reads the files itself, so an s3:// address that TigerGraph cannot reach is not an error GraFlo sees. This happens when TigerGraph runs in a container and MinIO's address is 127.0.0.1: that address is valid for GraFlo on the host but not inside the container. ingest.py warns about this case before it ingests. Set MINIO_LOADER_ENDPOINT in the MinIO settings to an address valid inside the TigerGraph container (for example http://172.17.0.1:9003 on Linux, or http://host.docker.internal:9003 where Docker provides it), or run with BULK_USE_S3=0, in which case TigerGraph must be able to read the local staging directory.

Connection refused on the MinIO port. GraFlo talks to MinIO's S3 API, not to its web console; they use different ports (MINIO_API_PORT and MINIO_CONSOLE_PORT). Check that docker ps shows graflo.minio as up. If the container never starts, the output of docker compose ... up names the reason; port is already allocated means another process uses that host port. Choose free ports in the MinIO settings, remove the container with docker rm -f graflo.minio and start it again.

Vertices are staged but edge_relates.csv is missing. No edge reached staging, so the problem is in the relations resource, not in S3 or the loading job.

Also possible

  • For AWS S3 or any other S3-compatible store, build an S3GeneralizedConnConfig yourself and register it under the same label; MinioConfig.to_s3_generalized_conn_config() is only a shortcut.
  • For other ways to emulate S3 in development, see Emulating S3 in development.

Files

The example lives in examples/13-tigergraph-bulk-s3.

manifest.yaml
# Companies and the relations between them, plus a staging entry for the bulk
# load. Bulk loading itself is switched on in ingest.py: it is a setting of the
# target database, not of the manifest.
schema:
    metadata:
        name: companies
        version: 1.0.0
    graph:
        vertex_config:
            vertices:
            -   name: company
                properties:
                -   name
                identity:
                -   name
        edge_config:
            edges:
            -   source: company
                target: company
                identities:
                -   -   relation
                properties:
                -   date
    db_profile: {}
ingestion_model:
    resources:
    -   name: relations
        pipeline:
        -   vertex: company
            from:
                name: company_a
        -   vertex: company
            from:
                name: company_b
        -   from: company
            to: company
            relation_field: relation
bindings:
    connectors:
    -   name: relations_files
        regex: "^relations.*\\.csv$"
        sub_path: data
        resource_name: relations
    # The staging area "bulk_s3" uses the S3 settings registered at run time
    # under the label "minio_bulk".
    staging_proxy:
    -   name: bulk_s3
        conn_proxy: minio_bulk
_common.py
"""Shared helpers for example 13 (TigerGraph bulk load + S3 staging)."""

from __future__ import annotations

import os
from collections.abc import Iterator
from contextlib import contextmanager
from pathlib import Path
from urllib.parse import urlparse

from suthing import FileHandle

from graflo import GraphManifest
from graflo.connections import TigergraphConfig
from graflo.object_storage import MinioConfig
from graflo.onto import DBType

EXAMPLE_DIR = Path(__file__).resolve().parent
MANIFEST_PATH = EXAMPLE_DIR / "manifest.yaml"
STAGING_DIR = EXAMPLE_DIR / "bulk_staging"

LOOPBACK_HOSTS = {"127.0.0.1", "localhost", "::1", "0.0.0.0"}


@contextmanager
def example_workdir() -> Iterator[Path]:
    """Manifest file connectors use ``sub_path: data`` relative to this example."""
    previous_cwd = os.getcwd()
    try:
        os.chdir(EXAMPLE_DIR)
        yield EXAMPLE_DIR
    finally:
        os.chdir(previous_cwd)


def load_manifest() -> GraphManifest:
    manifest = GraphManifest.from_config(FileHandle.load(MANIFEST_PATH))
    manifest.finish_init()
    return manifest


def physical_names() -> tuple[str, str]:
    """Storage names TigerGraph sees: ``(vertex type, edge relation)``.

    Derived rather than hard-coded — sanitization and the default relation name
    both belong to the schema, and the staged CSV file names follow them.
    """
    manifest = load_manifest()
    schema = manifest.graph_schema
    schema.db_profile.db_flavor = DBType.TIGERGRAPH
    schema_db = schema.resolve_db_aware(DBType.TIGERGRAPH)
    vertex = schema_db.vertex_config.vertex_dbname("company")
    edge = schema.core_schema.edge_config.edges[0]
    relation = schema_db.edge_config.runtime(edge).relation_name
    return vertex, relation or "relates"


def use_s3() -> bool:
    """``BULK_USE_S3`` — staging uploads to S3 (default) or local paths only."""
    return os.environ.get("BULK_USE_S3", "1").lower() in ("1", "true", "yes")


def tigergraph_config() -> TigergraphConfig:
    conn_conf = TigergraphConfig.from_docker_env()
    conn_conf.max_job_size = 5000
    return conn_conf


def minio_config() -> MinioConfig:
    conf = MinioConfig.from_docker_env()
    if bucket_override := os.environ.get("BULK_S3_BUCKET"):
        conf = conf.model_copy(update={"bucket": bucket_override})
    return conf


def loader_endpoint_is_loopback(conf: MinioConfig) -> bool:
    """True when TigerGraph would be handed a URL only this host can resolve.

    ``bulk_gsql`` falls back to ``endpoint_url`` when ``loader_endpoint_url`` is
    unset, so a loopback endpoint reaches the LOADING JOB unchanged. That is
    correct when TigerGraph runs on this host and wrong when it runs in Docker.
    """
    if conf.loader_endpoint_url is not None:
        return False
    return (urlparse(conf.endpoint_url).hostname or "") in LOOPBACK_HOSTS


def latest_session_dir() -> Path | None:
    """Newest ``bulk_staging/<session_id>/`` directory, if any run left one."""
    if not STAGING_DIR.is_dir():
        return None
    sessions = [p for p in STAGING_DIR.iterdir() if p.is_dir()]
    if not sessions:
        return None
    return max(sessions, key=lambda p: p.stat().st_mtime)
ingest.py
"""How do I load a large graph into TigerGraph quickly?

Switches the TigerGraph target to bulk loading: ingestion writes CSV files,
uploads them to S3 (MinIO here) and runs one TigerGraph loading job that reads
them. With ``BULK_USE_S3=0`` the files stay on local disk. Run it from any
directory:

    uv run python ingest.py
    uv run python inspect_bulk.py    # check what was staged and what was loaded
"""

from __future__ import annotations

import logging
import os

from _common import (
    STAGING_DIR,
    example_workdir,
    load_manifest,
    loader_endpoint_is_loopback,
    minio_config,
    tigergraph_config,
    use_s3,
)

from graflo.connections import (
    InMemoryConnectionProvider,
    TigergraphBulkLoadConfig,
    TigergraphConfig,
)
from graflo.hq import GraphEngine
from graflo.hq.caster import IngestionParams
from graflo.object_storage import ensure_staging_bucket_for_config

logger = logging.getLogger(__name__)


def configure_bulk_load(
    conn_conf: TigergraphConfig, provider: InMemoryConnectionProvider
) -> None:
    """Enable CSV staging on *conn_conf*, uploading to S3 unless BULK_USE_S3=0."""
    STAGING_DIR.mkdir(parents=True, exist_ok=True)

    if not use_s3():
        conn_conf.bulk_load = TigergraphBulkLoadConfig(
            enabled=True,
            staging_dir=str(STAGING_DIR),
        )
        return

    minio_conf = minio_config()
    conn_conf.bulk_load = TigergraphBulkLoadConfig(
        enabled=True,
        staging_dir=str(STAGING_DIR),
        s3_staging_name="bulk_s3",
        s3_bucket=minio_conf.bucket,
        s3_key_prefix=os.environ.get("BULK_S3_PREFIX", "demo"),
    )
    provider.register_generalized_config(
        conn_proxy="minio_bulk",
        config=minio_conf.to_s3_generalized_conn_config(),
    )

    # Preflight: raises with an actionable message when the S3 API is unreachable,
    # so a missing MinIO fails here rather than as an empty graph after a
    # "successful" ingest.
    ensure_staging_bucket_for_config(minio_conf)

    if loader_endpoint_is_loopback(minio_conf):
        logger.warning(
            "MinIO endpoint %s is a loopback address and MINIO_LOADER_ENDPOINT is "
            "unset, so TigerGraph's CREATE DATA_SOURCE will point there too. That "
            "works only if TigerGraph runs on this host; from a container the "
            "LOADING JOB reads nothing and the ingest still reports success. Set "
            "MINIO_LOADER_ENDPOINT for the MinIO docker stack, or run with "
            "BULK_USE_S3=0. Confirm either way with inspect_bulk.py.",
            minio_conf.endpoint_url,
        )


def main() -> None:
    logging.basicConfig(level=logging.WARNING, handlers=[logging.StreamHandler()])
    logging.getLogger("graflo").setLevel(logging.INFO)
    logger.setLevel(logging.INFO)

    with example_workdir():
        manifest = load_manifest()

        conn_conf = tigergraph_config()
        provider = InMemoryConnectionProvider()
        configure_bulk_load(conn_conf, provider)

        engine = GraphEngine(target_db_flavor=conn_conf.connection_type)
        engine.define_and_ingest(
            manifest=manifest,
            target_db_config=conn_conf,
            ingestion_params=IngestionParams(clear_data=True),
            recreate_schema=True,
            connection_provider=provider,
        )
    logger.info("Ingest finished (bulk_load=%s)", conn_conf.bulk_load)


if __name__ == "__main__":
    main()
inspect_bulk.py
"""
Check what bulk staging produced and what TigerGraph actually loaded.

Bulk ingest can report success and leave the graph empty — the LOADING JOB reads
the staged files itself, so an S3 URL the TigerGraph process cannot resolve is
not an error GraFlo ever sees. Run this after ingest.py:

    uv run python inspect_bulk.py            # staged files + loaded graph
    uv run python inspect_bulk.py --staged-only    # no database connection

Expected for this example: 3 companies (Acme, Beta, Gamma — the duplicate Acme
row upserts) and 2 relations, staged as company.csv and edge_relates.csv.
"""

from __future__ import annotations

import csv
import sys
from pathlib import Path

import click
from _common import (
    example_workdir,
    latest_session_dir,
    physical_names,
    tigergraph_config,
)

from graflo.architecture.graph_types import EdgeDirection
from graflo.db import ConnectionManager

EXPECTED_VERTICES = 3
EXPECTED_EDGES = 2


def _row_count(path: Path) -> int:
    """Data rows, excluding the header the staging writer emits."""
    with path.open(encoding="utf-8", newline="") as handle:
        return max(sum(1 for _ in csv.reader(handle)) - 1, 0)


def _report_staged(vertex_type: str, edge_type: str) -> bool:
    session = latest_session_dir()
    if session is None:
        click.echo("staged: nothing — bulk_staging/ has no session directory.")
        click.echo("        Run ingest.py first.")
        return False

    click.echo(f"staged: {session.name}/")
    staged = sorted(session.glob("*.csv"))
    for path in staged:
        click.echo(f"        {path.name:<24} {_row_count(path)} rows")
    if not staged:
        click.echo("        (empty)")

    names = {path.name for path in staged}
    vertex_file = f"{vertex_type}.csv"
    edge_file = f"edge_{edge_type}.csv"
    complete = True
    if vertex_file not in names:
        click.echo(f"        MISSING {vertex_file} — no vertices were staged.")
        complete = False
    if edge_file not in names:
        click.echo(
            f"        MISSING {edge_file} — vertices staged but edges did not. "
            "The relation is extracted per row, so an empty edge batch means the "
            "pipeline's edge step never fired."
        )
        complete = False
    return complete


def _report_graph(vertex_type: str, edge_type: str) -> bool:
    conn_conf = tigergraph_config()
    with ConnectionManager(connection_config=conn_conf) as db_client:
        vertices = db_client.fetch_docs(vertex_type)
        edges = db_client.fetch_edges(
            from_type=vertex_type,
            from_id="Acme",
            edge_type=edge_type,
            direction=EdgeDirection.ANY,
        )

    names = sorted(str(doc.get("name", "?")) for doc in vertices)
    click.echo(f"loaded: {len(vertices)} {vertex_type} ({', '.join(names) or '-'})")
    click.echo(f"        {len(edges)} {edge_type} incident to Acme")

    if not vertices:
        click.echo(
            "\nThe graph is empty. The LOADING JOB ran but read nothing — usually "
            "an s3:// URL that TigerGraph itself cannot reach. Set "
            "MINIO_LOADER_ENDPOINT to an address valid inside the TigerGraph "
            "container, or re-run with BULK_USE_S3=0 so the job reads local paths."
        )
        return False
    return len(vertices) == EXPECTED_VERTICES


@click.command()
@click.option(
    "--staged-only",
    is_flag=True,
    help="Only inspect the staged CSV files; do not connect to TigerGraph.",
)
def main(staged_only: bool) -> None:
    """Report staged CSVs and the loaded graph, and fail loudly on either gap."""
    with example_workdir():
        vertex_type, edge_type = physical_names()
        staged_ok = _report_staged(vertex_type, edge_type)
        if staged_only:
            sys.exit(0 if staged_ok else 1)
        click.echo("")
        graph_ok = _report_graph(vertex_type, edge_type)
    sys.exit(0 if staged_ok and graph_ok else 1)


if __name__ == "__main__":
    main()