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¶
| 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.
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:
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.csvwith 4 rows (one per company in each row of the CSV, so Acme twice) andedge_relates.csvwith 2 rows. TigerGraph keeps both kinds of relation in one edge type,relates, with the kind in itsrelationattribute. - Loaded: 3
companyvertices (Acme, Beta, Gamma; Acme is in both rows and is stored once) and 2relatesedges, 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
S3GeneralizedConnConfigyourself 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.
What to read next¶
- TigerGraph bulk load: every
setting of
bulk_loadand the limits of bulk mode. - Object storage.
- A graph on disk, without a database.
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()