GraphIngest SDK Documentation
The GraphIngest SDK provides decorators and wrappers that turn ordinary functions into orchestrated, observable pipeline steps. Each node is a single unit of work that runs in isolation. Each graph coordinates nodes with retries, timeouts, fan-out, and subgraph nesting.
Architecture Overview
extract
@node
transform
@node
load
@node
Quick Start
A complete, copy-paste-ready pipeline. This file does everything: defines nodes, composes them into a graph with retries, fans out in parallel, and deploys to the platform.
# pipeline.py β save this file and run: python pipeline.py
from graphingest import node, graph, deploy, RetryPolicy
import requests
# ββ Step 1: Define your nodes (individual tasks) ββ
@node(name="fetch-page", cache_ttl=3600, max_retries=3)
def fetch_page(url: str) -> dict:
"""Fetch a web page. Cached for 1 hour. Retries 3x on failure."""
resp = requests.get(url, timeout=10)
resp.raise_for_status()
return {"url": url, "status": resp.status_code, "length": len(resp.text)}
@node(name="summarize")
def summarize(page: dict) -> str:
"""Summarize a fetched page."""
return f"Page {page['url']}: {page['length']} chars, status {page['status']}"
# ββ Step 2: Compose nodes into a graph (pipeline) ββ
@graph(
name="web-scraper",
retry_policy=RetryPolicy(
max_retries=2,
delay_seconds=1,
backoff_factor=2,
jitter=True,
),
timeout_seconds=300,
)
def scrape_pipeline(urls: list[str]):
# Fan-out: fetch all URLs in parallel
pages = fetch_page.map(urls)
# Process each result
summaries = [summarize(page) for page in pages]
return {"total": len(summaries), "summaries": summaries}
# ββ Step 3: Deploy and run ββ
deploy() # push code to platform
# Execute the pipeline
result = scrape_pipeline([
"https://example.com",
"https://httpbin.org/get",
"https://jsonplaceholder.typicode.com/posts/1",
])
print(result)
# β {"total": 3, "summaries": ["Page ...: 1256 chars, status 200", ...]}Installation
pip install graphingest
# With AI agent support:
pip install graphingest[react]
# With LangGraph integration:
pip install graphingest[langgraph]@node / node()
A node is a single unit of work that runs in isolation. The SDK handles lifecycle management, logging, result serialization, and automatic retries so you can focus on your business logic.
Parameters
| Parameter | Type | Default | Description |
|---|---|---|---|
name | string | Function name | Unique key identifying this node |
cache_ttl | int / number | None | Cache TTL in seconds |
max_retries | int / number | 3 | Max automatic retries on failure |
tags | list / string[] | [] | Metadata tags for dashboard |
version | string | None | Semantic version string |
@node(name="extract-data", cache_ttl=3600, max_retries=5, tags=["etl"])
def extract(url: str) -> dict:
response = requests.get(url)
return response.json()
# Async support
@node(name="async-fetch")
async def fetch(url: str) -> dict:
async with httpx.AsyncClient() as client:
resp = await client.get(url)
return resp.json()@graph / graph()
A graph is a pipeline entrypoint β the top-level function that orchestrates nodes with retries (exponential backoff), timeouts, parameter validation, state hooks, run context, and streaming logs.
Parameters
| Parameter | Type | Default | Description |
|---|---|---|---|
name | string | Function name | Graph name for dashboard |
retry_policy | RetryPolicy | None | Retry configuration |
timeout | int / number | None | Max execution time (sec / ms) |
on_completion | hook[] | [] | Success callbacks |
on_failure | hook[] | [] | Failure callbacks |
on_cancellation | hook[] | [] | Timeout/cancel callbacks |
tags | list / string[] | [] | Metadata tags |
version | string | None | Semantic version |
.map() β Parallel Fan-Out
Fan-out a node across multiple inputs in parallel. Each input runs as a separate isolated execution. Results are collected in order. Must be called from within a @graph function.
@graph(name="parallel-pipeline")
def pipeline(urls: list[str]):
# Fan-out: dispatches len(urls) parallel workers
results = extract.map(urls)
# results[0] corresponds to urls[0], etc.
return results.submit() β Async Dispatch
Dispatch a single node asynchronously and get a NodeFuture back. The node runs in the background while your graph continues.
@graph(name="async-pipeline")
def pipeline(data: dict):
future = slow_transform.submit(data) # returns immediately
do_other_stuff() # do work while it runs
result = future.result(timeout=120) # block until ready
return resultNodeFuture
| Parameter | Type | Default | Description |
|---|---|---|---|
result() | Any / Promise | β | Block until done, return result |
RetryPolicy
Configurable retry strategy with exponential backoff and jitter to avoid thundering herd.
RetryPolicy(
max_retries=4, # Total retry attempts
delay_seconds=2, # Initial delay
backoff_factor=2.0, # Multiplier per attempt
max_delay_seconds=120, # Upper bound
jitter=True, # Β±50% randomization
)
# Delay formula:
# min(delay Γ backoff^attempt, max_delay) Γ random(0.5, 1.5)Presets
RetryPolicy(max_retries=3, delay_seconds=1, backoff_factor=2)Delays: ~1s, ~2s, ~4sRetryPolicy(max_retries=6, delay_seconds=0.5, backoff_factor=3, max_delay_seconds=60)Delays: ~0.5s, ~1.5s, ~4.5s, ~13.5s, ~40.5s, ~60sRetryPolicy(max_retries=3, delay_seconds=5, backoff_factor=1, jitter=False)Delays: 5s, 5s, 5sSubgraphs
Call a @graph from within another @graph. The child graph gets its own run ID, retry policy, timeout, and hooks β with parent_graph_run_id automatically linked.
@graph(name="sub-etl", retry_policy=RetryPolicy(max_retries=2))
def sub_etl(url: str) -> dict:
data = extract(url)
return transform(data)
@graph(name="main-pipeline")
def main_pipeline(urls: list[str]):
results = []
for url in urls:
result = sub_etl(url) # each creates a child graph run
results.append(result)
return resultsRun Contexts
Thread-safe (Python) / async-safe (TypeScript) context objects available from anywhere inside a running graph or node.
GraphRunContext
@graph(name="my-pipeline", version="2.0", tags=["prod"])
def pipeline(source: str):
ctx = GraphRunContext.get()
ctx.graph_run_id # "uuid-..."
ctx.graph_name # "my-pipeline"
ctx.graph_version # "2.0"
ctx.parameters # {"source": "..."}
ctx.tags # ["prod"]
ctx.parent_graph_run_id # None (or parent's ID)NodeRunContext
@node(name="extract")
def extract(url: str):
ctx = NodeRunContext.get()
ctx.node_run_id # "uuid-..."
ctx.node_key # "extract"
ctx.graph_run_id # parent graph's run ID
ctx.map_index # None (or int if .map())
ctx.retry_count # 0Streaming Logger
Real-time log streaming to the dashboard. Logs appear automatically as your pipeline runs.
import logging
logger = logging.getLogger(__name__)
@node(name="extract")
def extract(url: str):
logger.info("Starting extraction") # β appears in dashboard
logger.warning("Rate limit hit")
logger.error("Connection failed")
# Or use the dedicated logger:
from graphingest import get_run_logger
log = get_run_logger()
log.info("Pipeline started")LangGraph / AI Agent Orchestration
GraphIngest integrates with LangGraph to orchestrate AI agents at scale. Wrap any LangGraph agent as a @node, then use .map() to fan-out agents in parallel.
from graphingest.langgraph import agent_node, agent_graph, AgentConfig
researcher = agent_node(
name="researcher",
graph_builder=build_research_agent,
config=AgentConfig(
model="gpt-4o",
temperature=0.0,
max_iterations=10,
system_prompt="You are a research assistant.",
stream_steps=True,
),
cache_ttl=600,
retries=2,
)
# Use like any node:
result = researcher("What is quantum computing?")
# Fan-out: run 50 agents in parallel
results = researcher.map(["query1", "query2", ..., "query50"])Multi-Agent Pattern
Planner
breaks into N sub-tasks
Researcher
Agent 1
Researcher
Agent 2
Researcher
Agent N
Synthesizer
combines all findings
Start a job from GitHub
You keep a small file in your own GitHub repo. When it runs, GitHub sends a short note to graphingest.io: please start. Then GitHubβs computer turns off. The long work happens here. When it finishes, the mark on that save turns green or red.
Connect GitHub once in Settings. Add two secrets in the repo: GRAPHINGEST_API_KEY and GRAPHINGEST_FLOW_ID. The key is from Settings. The flow id is on graphingest.io/flows.
Save this as .github/workflows/graphingest.yml. The action lives at github.com/graphingest/run.
name: Start a GraphIngest job
on:
workflow_dispatch:
jobs:
start:
runs-on: ubuntu-latest
steps:
- uses: graphingest/run@v1
with:
api-key: ${{ secrets.GRAPHINGEST_API_KEY }}
flow-id: ${{ secrets.GRAPHINGEST_FLOW_ID }}Run it from the Actions tab. The dashboard run says it started from GitHub, and the link opens that save. Change workflow_dispatch to push when you want every save to start the job.
A Slack message when a job fails
The job fails while you are away. Slack gets one message: the job name, the error, and a link to that run on graphingest.io/runs. Open the link and continue from the row that failed.
In Slack, create an incoming webhook for the channel that should hear about failures. The address starts with https://hooks.slack.com/services/. Then send it to GraphIngest with your API key from Settings. The same form is on that Settings page if you would rather paste the address there.
curl -X POST https://www.graphingest.io/api/settings/notifications \
-H "Authorization: Bearer $GRAPHINGEST_API_KEY" \
-H "Content-Type: application/json" \
-d '{
"name": "Slack",
"channel_type": "slack",
"alert_level": "failures_only",
"config": { "webhook_url": "https://hooks.slack.com/services/..." }
}'failures_only is the usual choice. A finished job stays quiet. Use all when you also want a message for a job that finishes. The reply hides most of the webhook address. List what you saved with GET /api/settings/notifications. Send a test with POST /api/settings/notifications/test and the channel id. Remove one with DELETE /api/settings/notifications and { "id": "..." }.
Configuration
Set these environment variables before running your pipeline. You can find your API URL and key in the dashboard settings.
| Variable | Required | Description |
|---|---|---|
GRAPHINGEST_API_URL | Yes | Your GraphIngest API endpoint |
GRAPHINGEST_API_KEY | Yes | Your API key (found in dashboard) |
Full Examples
Complete, runnable examples for real-world use cases. Each file is self-contained β copy, set your env vars, and run.
Python12 examples
π Quick Start β Web Scraper
Fan-out 3 URLs, caching, retries, RetryPolicy
β° Vercel Cron Timeout Fix
Cron returns <1s, pipeline runs for hours on Cloud Run
π PDF Report Generation
Background job, job ID, frontend polling pattern
πΌοΈ Image Processing at Scale
10K images in parallel with .map() β minutes not hours
π§ LLM Content Pipeline
ThrottlePolicy + cache to avoid rate limits and $ waste
πΎ Database Migration
Cache makes restart instant β skip already-migrated batches
π ETL Pipeline
Multi-source extract β transform β load with subgraphs
π€ AI Agent
ReAct research agent with search/scrape/summarize tools
π§ LangGraph Fan-Out
Parallel agents, five at a time, with a throttle on the next batch
π Flow Control
Per-user concurrency, throttling, priority queuing
π FastAPI Backend
Dispatch jobs + status polling endpoints + frontend JS
π₯ Multi-Tenant SaaS
Free vs Pro tiers with different limits
TypeScript8 examples
π Quick Start β Web Scraper
Fan-out 3 URLs, caching, retries, RetryPolicy
β° Vercel Cron Timeout Fix
Cron returns <1s, pipeline runs for hours
β±οΈ Vercel Cron Route
Next.js route that returns immediately and dispatches the pipeline
π¦ No Workers
deploy() a TypeScript graph. Nothing stays running when it's idle
π PDF Report Generation
Background job with async dispatch and polling
πΌοΈ Image Processing
10K images in parallel with .map()
π§ LLM Content Pipeline
Throttle + cache for OpenAI rate limits
π ETL Pipeline
Multi-source ETL with subgraphs
Go5 examples
π Quick Start β Web Scraper
Fan-out URLs, caching, RetryPolicy
π ETL Pipeline
Multi-source ETL with .Map() fan-out
β³ Async Jobs
Background jobs with .Submit() and .Map()
π¦ No Workers
deploy() a Go graph. Nothing stays running when it's idle
π― Full Feature Demo
All SDK features: nodes, graphs, subgraphs, futures
Ready to get started?
Deploy your first pipeline in under 5 minutes.