Fetch Data from APIs

On this page

Pull data from external APIs into your DuckLake pipeline. The runtime handles pagination, retries, and incremental state — your lib function handles the API-specific I/O. SQL handles the transformation.

The pattern

Every API integration follows the same structure:

  1. API dict — declares auth, endpoints, rate limits, pagination
  2. Starlark fetch() — calls the API, parses responses, returns rows
  3. Raw SQL model — fetches with API column names
  4. Staging SQL model — transforms: casts types, renames, pivots, expands JSON
lib/my_api.star          → I/O logic
models/raw/data.sql      → FROM my_api('resource')
models/staging/data.sql   → SELECT ... FROM raw.data (transformation)

Quick start

1. Create a lib function

# lib/my_api.star
API = {
    "base_url": "https://api.example.com",
    "auth": {"env": "MY_API_KEY"},
    "fetch": {
        "args": ["resource"],
        "page_size": 100,
    },
}

def fetch(resource, page):
    resp = http.get("/v1/" + resource, params={
        "limit": page.size,
        "cursor": page.cursor,
    })
    if not resp.ok:
        fail("API error: " + str(resp.status_code))

    data = resp.json
    return {"rows": data["items"], "next": data.get("next_cursor")}

2. Write SQL models

Raw — fetch with API column names:

-- models/raw/users.sql
-- @kind: table
-- @fetch

SELECT id::BIGINT AS id, name::VARCHAR AS name, email::VARCHAR AS email,
       metadata::JSON AS metadata
FROM my_api('users')

Staging — transform in SQL:

-- models/staging/users.sql
-- @kind: table

SELECT
    id::BIGINT AS id,
    name,
    LOWER(email) AS email,
    metadata->>'$.city' AS city
FROM raw.users

3. Run

ondatrasql run

Both models participate in the DAG. The staging model runs automatically after the raw model.

SQL casts control the API

The runtime extracts column names and DuckDB-native types from your SELECT casts via DuckDB AST. These are passed to fetch() as the columns kwarg:

SQLtype value
total::DECIMAL(18,3)"DECIMAL(18,3)"
count::INTEGER"INTEGER"
count::BIGINT"BIGINT"
items::JSON"JSON"
date::DATE"DATE"
name::VARCHAR"VARCHAR"
tags::VARCHAR[]["VARCHAR"]
address::STRUCT(street VARCHAR, zip INTEGER){"street": "VARCHAR", "zip": "INTEGER"}

Blueprints can use the type value to adapt their API requests — for example, the GAM blueprint splits columns into dimensions (VARCHAR) and metrics (numeric types). For blueprints that need JSON Schema, use lib_helpers.to_json_schema(col["type"]) to convert.

Pagination patterns

The runtime calls fetch() in a loop until "next" is None. Your blueprint owns the pagination logic:

Cursor-based

def fetch(resource, page):
    params = {"limit": page.size}
    if page.cursor:
        params["starting_after"] = page.cursor

    resp = http.get("/v1/" + resource, params=params)
    if not resp.ok:
        fail("API error: " + str(resp.status_code))
    data = resp.json

    next_cursor = None
    if data.get("has_more") and len(data["data"]) > 0:
        next_cursor = data["data"][-1]["id"]

    return {"rows": data["data"], "next": next_cursor}

Offset-based

def fetch(resource, page):
    offset = page.cursor or 0
    resp = http.get("/v1/" + resource, params={"limit": page.size, "offset": offset})
    if not resp.ok:
        fail("API error: " + str(resp.status_code))
    rows = resp.json["items"]
    next_offset = offset + page.size if len(rows) == page.size else None
    return {"rows": rows, "next": next_offset}

Date-range

Fetch one time window per page:

def fetch(series, page, is_backfill=True, initial_value="", last_value=""):
    if is_backfill:
        start_date = initial_value
    else:
        start_date = _next_day(last_value)

    cursor_date = start_date if page.cursor == None else page.cursor
    end_date = min(_add_days(cursor_date, 365), yesterday)

    resp = http.get("/data/" + series + "/" + cursor_date + "/" + end_date)
    if not resp.ok:
        fail("API error: " + str(resp.status_code))
    rows = [{"date": obs["date"], "value": obs["value"]} for obs in resp.json]

    next_cursor = _next_day(end_date) if end_date < yesterday else None
    return {"rows": rows, "next": next_cursor}

Dict cursor

When you need to carry multiple values between pages, return a dict — it passes through directly:

next_cursor = {"url": fetch_url, "token": next_token, "series_idx": 3}
# On next page:
series_idx = page.cursor["series_idx"]

Async fetch (report-style APIs)

For APIs that create reports asynchronously (submit → poll → fetch results), declare async: True and implement three functions:

API = {
    "fetch": {
        "async": True,
        "poll_interval": "5s",
        "poll_timeout": "5m",
        "poll_backoff": 2,
    },
}

def submit(columns=[], is_backfill=True, last_value=""):
    resp = http.post("/reports", json={...})
    if not resp.ok:
        fail("submit failed: " + str(resp.status_code))
    return {"job_id": resp.json["id"]}

def check(job_ref):
    resp = http.get("/reports/" + job_ref["job_id"])
    if not resp.ok:
        fail("poll failed: " + str(resp.status_code))
    if resp.json["status"] == "complete":
        return {"url": resp.json["result_url"]}
    return None  # keep polling

def fetch_result(result_ref, page):
    resp = http.get(result_ref["url"] + "?pageSize=" + str(page.size))
    if not resp.ok:
        fail("result failed: " + str(resp.status_code))
    return {"rows": resp.json["data"], "next": resp.json.get("next_page")}

The runtime handles the poll loop. See Fetch Contract — Async fetch.

Nested and structured data

APIs return nested JSON. Blueprints return it as-is — SQL handles expansion in a downstream model.

json_each — simplest, no schema needed

-- Raw: fetch with arrays as JSON
SELECT id::BIGINT AS id, name::VARCHAR AS name, tags::JSON AS tags FROM my_api('users')

-- Staging: expand
SELECT id, name, j.value->>'name' AS tag
FROM raw.users, LATERAL json_each(tags) AS j

json_transform — typed expansion

SELECT id, x.*
FROM raw.users,
LATERAL (SELECT unnest(json_transform(tags,
    '[{"name":"VARCHAR","color":"VARCHAR"}]')) AS x)

Arrow operators — direct access

SELECT id, metadata->>'$.address.city' AS city
FROM raw.users

Pivot — normalize then pivot in SQL

When the API returns one row per measurement:

-- Raw: normalized (series, date, value)
SELECT series::VARCHAR AS series, date::DATE AS date, value::DECIMAL(18,6) AS value
FROM riksbank('SEKEURPMI,SEKUSDPMI')

-- Staging: pivoted
SELECT date::DATE,
    MAX(CASE WHEN series = 'SEKEURPMI' THEN value END)::DECIMAL AS eur,
    MAX(CASE WHEN series = 'SEKUSDPMI' THEN value END)::DECIMAL AS usd
FROM raw.exchange_rates
GROUP BY date

Or with DuckDB’s native PIVOT:

PIVOT raw.exchange_rates ON series USING MAX(value) GROUP BY date::DATE

Incremental loads

Use @incremental to fetch only new data on subsequent runs:

-- models/raw/events.sql
-- @kind: append
-- @fetch
-- @incremental: created_at

SELECT id::BIGINT AS id, name::VARCHAR AS name, created_at::TIMESTAMP AS created_at
FROM my_api('events')
def fetch(resource, page, is_backfill=True, last_value=""):
    params = {"limit": page.size, "cursor": page.cursor}
    if not is_backfill:
        params["created_after"] = last_value

    resp = http.get("/v1/" + resource, params=params)
    return {"rows": resp.json["items"], "next": resp.json.get("next")}
KwargDescription
is_backfillTrue on first run — fetch everything
last_valueMAX(cursor_column) from previous run
initial_valueStarting value from @incremental_initial

Declare only the kwargs you use — the runtime filters automatically.

Authentication

Configured in the API dict. Injected into every http.* call automatically.

# API key
"auth": {"env": "API_KEY"}

# OAuth provider (browser consent, or inject ONDATRA_OAUTH_TOKEN_<PREFIX>)
"auth": {"provider": "hubspot"}

# Basic auth
"auth": {"user": {"env": "USER"}, "pass": {"env": "PASS"}}

# Google service account
"auth": {"service_account": {"env": "KEY_FILE"}, "scope": "https://..."}

Error handling

Always check resp.ok. http.get/http.post/http.request only raise on 5xx, 429, and transport failures (those retry, then fail the run). A 4xx (404, 400, 401, 403) does not raise — it returns a response with resp.ok == False. If you skip the check and read resp.json anyway, a client error is silently parsed into 0 rows, and for a table or scd2 model that 0-row result wipes / closes the target. OndatraSQL emits a warning at run time for any fetch that calls http.* without checking resp.ok / resp.status_code; silence an intentional case (e.g. an API that returns 404 for an empty collection) with a # ondatracheck:allow-unchecked-status <reason> comment in the lib file.

Fail — stops the pipeline:

if not resp.ok:
    fail("API error: " + str(resp.status_code) + " " + resp.text)

Abort — clean exit, 0 rows, no error:

if start_date > yesterday:
    abort()

Retry — automatic. Configure in API dict:

API = {"retry": 3, "backoff": 2, "timeout": 30}

Rate limiting — declare in API dict, runtime handles it:

API = {"rate_limit": {"requests": 2, "per": "1s"}}

Options via JSON arg

For API-specific configuration beyond columns:

FROM gam_report('{"custom_dimensions": [11678108], "currency": "USD"}')
API = {"fetch": {"args": ["options"]}}

def fetch(options="", page=None, columns=[]):
    opts = json.decode(options) if options else {}

Reference