rig major updates
This commit is contained in:
81
rig/examples/data/dags/items_to_postgres.py
Normal file
81
rig/examples/data/dags/items_to_postgres.py
Normal file
@@ -0,0 +1,81 @@
|
||||
"""Pull items from the items API, rename them into the app's names, upsert into postgres.
|
||||
|
||||
Three links, kept apart on purpose:
|
||||
|
||||
- the API client (`fetch_items`) talks to the wire as it is — here the overlay's
|
||||
own simulator, `items-api`, whose field names are the API's;
|
||||
- the adapter (`to_app_row`) is the one place the wire's names become the app's:
|
||||
`id` -> `item_id`, `name` -> `item_name`, `price.amount_cents` -> `price_cents`.
|
||||
It belongs to whoever owns the app's model, so it lives in the overlay, not in rig;
|
||||
- the load writes through the `app_db` connection (AIRFLOW_CONN_APP_DB, built by
|
||||
addons/airflow.sh from the postgres secret) and is idempotent: an upsert keyed
|
||||
on `item_id`, so a retry or a rerun never duplicates a row.
|
||||
|
||||
Operational logic is explicit and minimal: hourly, no backfill, one retry.
|
||||
"""
|
||||
import json
|
||||
import urllib.request
|
||||
from datetime import datetime, timedelta
|
||||
|
||||
from airflow import DAG
|
||||
from airflow.operators.python import PythonOperator
|
||||
|
||||
ITEMS_URL = "http://items-api/v1/items"
|
||||
|
||||
CREATE = """
|
||||
CREATE TABLE IF NOT EXISTS items (
|
||||
item_id text PRIMARY KEY,
|
||||
item_name text NOT NULL,
|
||||
price_cents integer NOT NULL,
|
||||
currency text NOT NULL,
|
||||
loaded_at timestamptz NOT NULL DEFAULT now()
|
||||
)
|
||||
"""
|
||||
|
||||
UPSERT = """
|
||||
INSERT INTO items (item_id, item_name, price_cents, currency)
|
||||
VALUES (%(item_id)s, %(item_name)s, %(price_cents)s, %(currency)s)
|
||||
ON CONFLICT (item_id) DO UPDATE
|
||||
SET item_name = EXCLUDED.item_name,
|
||||
price_cents = EXCLUDED.price_cents,
|
||||
currency = EXCLUDED.currency,
|
||||
loaded_at = now()
|
||||
"""
|
||||
|
||||
|
||||
def fetch_items():
|
||||
"""The API client: the wire, as the API returns it."""
|
||||
with urllib.request.urlopen(ITEMS_URL, timeout=10) as response:
|
||||
return json.load(response)["items"]
|
||||
|
||||
|
||||
def to_app_row(item):
|
||||
"""The adapter: the API's names in, the app's names out."""
|
||||
return {
|
||||
"item_id": item["id"],
|
||||
"item_name": item["name"],
|
||||
"price_cents": item["price"]["amount_cents"],
|
||||
"currency": item["price"]["currency"],
|
||||
}
|
||||
|
||||
|
||||
def load_items():
|
||||
from airflow.providers.postgres.hooks.postgres import PostgresHook
|
||||
|
||||
rows = [to_app_row(item) for item in fetch_items()]
|
||||
hook = PostgresHook(postgres_conn_id="app_db")
|
||||
hook.run(CREATE)
|
||||
for row in rows:
|
||||
hook.run(UPSERT, parameters=row)
|
||||
print(f"upserted {len(rows)} items")
|
||||
|
||||
|
||||
with DAG(
|
||||
dag_id="items_to_postgres",
|
||||
schedule="@hourly",
|
||||
start_date=datetime(2026, 1, 1),
|
||||
catchup=False,
|
||||
default_args={"retries": 1, "retry_delay": timedelta(minutes=1)},
|
||||
tags=["example"],
|
||||
) as dag:
|
||||
PythonOperator(task_id="load_items", python_callable=load_items)
|
||||
Reference in New Issue
Block a user