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.
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.
FlowExecutionRequest / FlowAuditRecordRun 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.
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
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"}))
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"}'
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.
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.
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.
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.
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.
# 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.
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.
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.
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.
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.
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.
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.
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")