REST API

Code examples

Eight end-to-end recipes plus three patterns. Each recipe shows the same task in cURL, Python, Bash, and PowerShell: pick the tab that matches your stack.

Setup

Every example assumes you've set ETLWORKS_URL (e.g. https://app.etlworks.com/rest) and ETLWORKS_API_KEY in your environment. See the Quickstart for getting these.

Run a flow and wait for completion

Trigger a flow, then poll the executions endpoint until it reaches a terminal status. The most common pattern in production scripts.

Real shapes from FlowExecutionRequest / FlowAuditRecord

Run body is a flat map of string→string parameters, not wrapped in {"params": …}. Send {} if your flow takes no parameters. The response is a FlowExecutionResponse with flowId, auditId, and an int status code. To track progress, use auditId (not the int code) against GET /v1/executions/{flowId}?auditId={auditId}, which returns a FlowAuditRecord with a string status: queued, running, success, warning, error, or canceled.

Shell
FLOW_ID=12345

# Run the flow. Body is a flat map of parameters (or {} for none).
AUDIT_ID=$(curl -s -X POST "$ETLWORKS_URL/v1/flows/$FLOW_ID/run" \
  -H "Authorization: Bearer $ETLWORKS_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{"runDate":"2026-05-09"}' | jq -r '.auditId')

# Poll the audit record until status is terminal.
while :; do
  STATUS=$(curl -s "$ETLWORKS_URL/v1/executions/$FLOW_ID?auditId=$AUDIT_ID" \
    -H "Authorization: Bearer $ETLWORKS_API_KEY" | jq -r '.status')
  echo "  $STATUS"
  case "$STATUS" in success|warning|error|canceled) break ;; esac
  sleep 5
done
Python
import os, time, requests

BASE = os.environ["ETLWORKS_URL"]
H = {"Authorization": f"Bearer {os.environ['ETLWORKS_API_KEY']}"}
TERMINAL = {"success", "warning", "error", "canceled"}

def run_and_wait(flow_id, parameters=None, timeout=900):
    # Body is a flat map of {string: string} parameters (or {} for none).
    r = requests.post(f"{BASE}/v1/flows/{flow_id}/run",
                      headers={**H, "Content-Type": "application/json"},
                      json=parameters or {})
    r.raise_for_status()
    audit_id = r.json()["auditId"]   # NOT executionId

    deadline = time.time() + timeout
    while time.time() < deadline:
        rec = requests.get(f"{BASE}/v1/executions/{flow_id}",
                           headers=H, params={"auditId": audit_id}).json()
        if rec["status"] in TERMINAL:
            return rec
        time.sleep(5)
    raise TimeoutError(f"Flow {flow_id} (audit {audit_id}) still running after {timeout}s")

print(run_and_wait(12345, parameters={"runDate": "2026-05-09"}))
Bash
run_and_wait () {
  local flow_id="$1" body="${2:-{}}"
  local audit_id status
  audit_id=$(curl -s -X POST "$ETLWORKS_URL/v1/flows/$flow_id/run" \
    -H "Authorization: Bearer $ETLWORKS_API_KEY" \
    -H "Content-Type: application/json" -d "$body" | jq -r '.auditId')
  while :; do
    status=$(curl -s "$ETLWORKS_URL/v1/executions/$flow_id?auditId=$audit_id" \
      -H "Authorization: Bearer $ETLWORKS_API_KEY" | jq -r '.status')
    case "$status" in success|warning|error|canceled) echo "$status"; return ;; esac
    sleep 5
  done
}

run_and_wait 12345 '{"runDate":"2026-05-09"}'
PowerShell
function Invoke-Flow ($FlowId, $Parameters = @{}) {
  $headers = @{ "Authorization" = "Bearer $env:ETLWORKS_API_KEY" }
  # Body is the parameters map directly, not wrapped in another object.
  $body    = $Parameters | ConvertTo-Json -Compress
  $started = Invoke-RestMethod -Method Post -Uri "$env:ETLWORKS_URL/v1/flows/$FlowId/run" `
              -Headers $headers -ContentType "application/json" -Body $body
  $auditId = $started.auditId
  $terminal = "success","warning","error","canceled"
  do {
    Start-Sleep 5
    $rec = Invoke-RestMethod -Uri "$env:ETLWORKS_URL/v1/executions/$FlowId`?auditId=$auditId" -Headers $headers
    Write-Host "  $($rec.status)"
  } while ($rec.status -notin $terminal)
  $rec
}
Invoke-Flow -FlowId 12345 -Parameters @{ runDate = "2026-05-09" }

Create a connection

Add a Postgres connection programmatically: useful when provisioning environments from CI/CD or Terraform-style scripts.

Connection types are metadata-driven

The connectionType discriminator and the keys inside properties are not arbitrary strings, they're defined per-connector in the platform's connector library (under src/main/resources/di/…/connections/<type>.json). For PostgreSQL, the type is connection.db.postgres and the property keys are field.connection.template.host, field.connection.template.port, field.connection.template.database, etc. The same pattern holds for every connector.

Easiest way to discover the right shape for any connector: create the connection once in the Etlworks UI, then GET /v1/connections/{id} to see the exact connectionType and properties keys it uses. Copy that shape into your script and parameterize the values.

Shell
curl -s -X POST "$ETLWORKS_URL/v1/connections" \
  -H "Authorization: Bearer $ETLWORKS_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "name": "warehouse-prod",
    "description": "Production warehouse",
    "connectionType": "connection.db.postgres",
    "properties": {
      "field.connection.template.host": "warehouse.internal",
      "field.connection.template.port": "5432",
      "field.connection.template.database": "analytics",
      "field.connection.user": "etlworks",
      "field.connection.password": "{{secret:warehouse_pw}}"
    },
    "tags": ["production", "warehouse"]
  }'

The {{secret:…}} placeholder pulls from your tenant's secret store at runtime: the literal value is never sent over the API. Test the connection without saving via POST /v1/connections/test with the same body.

Schedule a flow

Schedule an existing flow. The body matches the FlowSchedule model: flowId, name, the schedule string (the platform's scheduling syntax: the UI's schedule picker is the source of truth for valid expressions), and enabled. Note: pause/resume use /disable and /enable, not /pause//resume.

Shell
curl -s -X POST "$ETLWORKS_URL/v1/schedules" \
  -H "Authorization: Bearer $ETLWORKS_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "name": "Daily warehouse refresh",
    "description": "Refresh the warehouse from staging every day at 4am",
    "flowId": 12345,
    "schedule": "0 0 4 * * ?",
    "enabled": true,
    "tags": ["production"]
  }'

# Pause:
curl -s -X POST "$ETLWORKS_URL/v1/schedules/$SCHEDULE_ID/disable" \
  -H "Authorization: Bearer $ETLWORKS_API_KEY"

# Resume:
curl -s -X POST "$ETLWORKS_URL/v1/schedules/$SCHEDULE_ID/enable" \
  -H "Authorization: Bearer $ETLWORKS_API_KEY"

Query the audit log

The audit log records every API call: who made it, the method, when, how long it took, and whether it failed. Pull the last 7 days of failed calls. Administrators only. For what changed inside a flow, use its version history, GET /v1/flows/{flowId}/history.

Python
from datetime import date, timedelta

records = requests.post(
    f"{BASE}/v1/audit",
    headers=H,
    params={"exception": "true"},                       # failed calls only
    json={"fromDate": (date.today() - timedelta(days=7)).isoformat()},
).json()

# Keys are abbreviated on the wire: dt date, ur user, md method, dn duration (ms), ex exception
for r in records:
    print(f"{r['dt']}  {r['ur']:<28} {r['md']:<40} {r['dn']:>6} ms  {r.get('ex') or ''}")

Export and import (env promotion)

Export a flow from staging and import it into production. The export is a text bundle that carries the flow and the connections and formats it uses; you can store it, diff it, or import it into another tenant. Import takes the bundle as the dump field of a JSON body, so preview first, then apply.

Shell
# Export one or more flows (comma-separated ids)
curl -s "$ETLWORKS_URL/v1/flows/export?flows=12345&excludeCredentials=true" \
  -H "Authorization: Bearer $ETLWORKS_API_KEY" \
  > flow-bundle.txt

# Wrap the bundle in an ImportRequest and preview it in production
jq -Rs '{dump: ., collisionPolicy: "Replace"}' flow-bundle.txt > import.json

curl -s -X POST "$ETLWORKS_PROD_URL/v1/flows/import/preview" \
  -H "Authorization: Bearer $PROD_API_KEY" \
  -H "Content-Type: application/json" \
  --data-binary @import.json

# Apply if the preview looks right
curl -s -X POST "$ETLWORKS_PROD_URL/v1/flows/import" \
  -H "Authorization: Bearer $PROD_API_KEY" \
  -H "Content-Type: application/json" \
  --data-binary @import.json

collisionPolicy decides what happens when an object already exists: Add, Replace, Keep, or Keep All, Replace Flows and Macros. With excludeCredentials=true the production connections keep their own passwords; set importCredentials only if you mean to overwrite them. Connections and formats on their own move through GET /v1/io/export and POST /v1/io/import the same way.

Bulk-tag resources

Find the flows by name, then add a tag to all of them in one call with POST /v1/flows/bulk. The bulk endpoint changes only what the actions name, so there is no risk of overwriting other fields. Connections, formats and listeners have the same endpoint with ioIds, and schedules with scheduleIds.

Python
flows = requests.get(f"{BASE}/v1/flows", headers=H).json()
ids = [f["id"] for f in flows if "warehouse" in f["name"].lower()]

r = requests.post(f"{BASE}/v1/flows/bulk", headers=H, json={
    "flowIds": ids,
    "actions": [{"type": "tags.add", "tags": ["warehouse-pipeline"]}],
    "modifiedMessage": "Tag warehouse flows",
})
r.raise_for_status()
print(r.json())   # {"status": ..., "totalFlows": ..., "updatedFlows": ..., "updatedFlowIds": [...]}

Build a health check

Count executions by status for each flow you care about, the data behind an internal "is the platform healthy?" board. GET /v1/executions/{flowId}/statuses returns a map from status to count.

Python
flows = requests.get(f"{BASE}/v1/flows", headers=H).json()

for f in flows:
    if "production" not in (f.get("tags") or []):
        continue
    counts = requests.get(f"{BASE}/v1/executions/{f['id']}/statuses", headers=H).json()
    # e.g. {"success": 1809, "error": 33, "warning": 4}
    failed = counts.get("error", 0)
    total = sum(counts.values())
    print(f"{f['name']:<48} {total:>6} runs  {failed:>4} failed")

Subscribe to flow events (outbound webhooks)

Etlworks fires webhooks, it doesn't receive them. Register a webhook to be POSTed to your URL when the listed events fire on the flows you care about. Body matches the Webhook model.

Need inbound triggers (something else → run a flow)?

That's a listener connection, not a webhook. Create an HTTP listener (or a queue listener for Kafka, SQS or MQTT) with POST /v1/listeners, then bind a flow to it with a schedule. The listener URL Etlworks returns is what your upstream system POSTs to.

Shell
curl -s -X POST "$ETLWORKS_URL/v1/webhooks" \
  -H "Authorization: Bearer $ETLWORKS_API_KEY" \
  -H "Content-Type: application/json" \
  -d '{
    "name": "warehouse-flow-events",
    "description": "Notify ops on warehouse-pipeline failures and successes",
    "events": ["Flow.Failed", "Flow.Executed"],
    "url": "https://your-app.example.com/etlworks-events",
    "method": "POST",
    "secret": "use-this-for-hmac-verification",
    "flowIncludes": [12345],
    "enabled": true
  }'

Event names are the ones the UI shows, such as Flow.Failed, Flow.Executed, Flow.Warning or Connection.Updated, or * for every event; GET /v1/webhooks/metadata lists them, and an unknown name is rejected with 400. When secret is set, each delivery carries an HMAC-SHA1 signature of the payload in X-Etl-Webhook-Signature, alongside X-Etl-Webhook-Event and X-Etl-Webhook-Event-Id. Limit it to specific flows with flowIncludes, or send every flow except a deny-list with flowExcludes.


Pattern: walk execution history

List endpoints return everything in one response and need no paging. Execution history is the exception: it caps each response and reports how many older records it left out in X-Truncated-Count. Walk back by passing the oldest auditId you have as lastAuditId.

Python
def executions(flow_id, **filters):
    last = None
    while True:
        params = dict(filters, **({"lastAuditId": last} if last else {}))
        r = requests.get(f"{BASE}/v1/executions/{flow_id}/history", headers=H, params=params)
        r.raise_for_status()
        page = r.json()
        yield from page
        if not page or int(r.headers.get("X-Truncated-Count", 0)) == 0:
            return
        last = min(e["auditId"] for e in page)

for e in executions(12345, status="error", fromDate="2026-09-01"):
    print(e["auditId"], e["status"], e["started"])

Pattern: retry with exponential backoff

Wraps a request in retry-on-5xx and retry-on-429 logic, honoring the Retry-After header.

Python
import time, random

def call_with_retry(method, url, attempts=5, **kwargs):
    for i in range(attempts):
        r = requests.request(method, url, headers=H, **kwargs)
        if r.status_code < 500 and r.status_code != 429:
            r.raise_for_status()
            return r.json() if r.content else None
        delay = int(r.headers.get("Retry-After", 0)) or min(60, (2 ** i) + random.random())
        time.sleep(delay)
    raise RuntimeError(f"Gave up after {attempts} attempts: {url}")

Pattern: parallel fan-out

Run twenty flows at once, wait for all of them. Useful for nightly batch jobs that have no shared dependency.

Python
from concurrent.futures import ThreadPoolExecutor

flow_ids = [12345, 12346, 12347, 12348, …]

with ThreadPoolExecutor(max_workers=8) as pool:
    results = list(pool.map(lambda fid: run_and_wait(fid), flow_ids))

failed = [r for r in results if r["status"] != "success"]
print(f"  {len(results)} run, {len(failed)} failed")