From fda13a7f28a0a6e80a448225397d84d43d7e4e9b Mon Sep 17 00:00:00 2001 From: Waldemar Hummer Date: Mon, 5 Oct 2026 03:45:05 +0900 Subject: [PATCH 1/3] add localstack-bruin: Bruin data pipeline on LocalStack S3 + Snowflake Sample Bruin pipeline that loads raw CSVs from LocalStack S3 into the LocalStack Snowflake emulator (external stage + COPY INTO), transforms them with data quality checks, and writes a report back to S3. Includes a CI workflow that runs the pipeline end to end and checks the report. Co-Authored-By: Claude Opus 5.5 --- .github/workflows/test-localstack-bruin.yml | 58 ++++++++++++++ .gitignore | 4 + localstack-bruin/.bruin.yml | 21 +++++ localstack-bruin/.gitignore | 7 ++ localstack-bruin/Makefile | 58 ++++++++++++++ localstack-bruin/README.md | 76 +++++++++++++++++++ localstack-bruin/data/customers.csv | 9 +++ localstack-bruin/data/orders.csv | 17 +++++ localstack-bruin/docker-compose.yml | 11 +++ localstack-bruin/pipeline/assets/common.py | 31 ++++++++ .../pipeline/assets/orders_enriched.py | 54 +++++++++++++ localstack-bruin/pipeline/assets/raw_load.py | 45 +++++++++++ .../pipeline/assets/requirements.txt | 2 + .../pipeline/assets/revenue_by_country.py | 52 +++++++++++++ localstack-bruin/pipeline/pipeline.yml | 1 + 15 files changed, 446 insertions(+) create mode 100644 .github/workflows/test-localstack-bruin.yml create mode 100644 localstack-bruin/.bruin.yml create mode 100644 localstack-bruin/.gitignore create mode 100644 localstack-bruin/Makefile create mode 100644 localstack-bruin/README.md create mode 100644 localstack-bruin/data/customers.csv create mode 100644 localstack-bruin/data/orders.csv create mode 100644 localstack-bruin/docker-compose.yml create mode 100644 localstack-bruin/pipeline/assets/common.py create mode 100644 localstack-bruin/pipeline/assets/orders_enriched.py create mode 100644 localstack-bruin/pipeline/assets/raw_load.py create mode 100644 localstack-bruin/pipeline/assets/requirements.txt create mode 100644 localstack-bruin/pipeline/assets/revenue_by_country.py create mode 100644 localstack-bruin/pipeline/pipeline.yml diff --git a/.github/workflows/test-localstack-bruin.yml b/.github/workflows/test-localstack-bruin.yml new file mode 100644 index 0000000..aa59d33 --- /dev/null +++ b/.github/workflows/test-localstack-bruin.yml @@ -0,0 +1,58 @@ +name: Bruin on LocalStack + +on: + pull_request: + branches: [main] + paths: + - 'localstack-bruin/**' + - '.github/workflows/test-localstack-bruin.yml' + push: + branches: [main] + paths: + - 'localstack-bruin/**' + - '.github/workflows/test-localstack-bruin.yml' + workflow_dispatch: + +env: + LOCALSTACK_AUTH_TOKEN: ${{ secrets.TEST_LOCALSTACK_AUTH_TOKEN }} + +jobs: + test-localstack-bruin: + name: Bruin pipeline on LocalStack S3 + Snowflake + runs-on: ubuntu-latest + timeout-minutes: 20 + defaults: + run: + working-directory: localstack-bruin + + steps: + - name: Checkout + uses: actions/checkout@v4 + + - name: Install Bruin CLI + run: | + curl -LsSf https://getbruin.com/install/cli | sh + echo "$HOME/.local/bin" >> "$GITHUB_PATH" + + - name: Start LocalStack + run: make start + + - name: Seed S3 + run: make seed + + - name: Validate pipeline + run: make validate + + - name: Run pipeline + run: make run + + - name: Check report + run: make test + + - name: Dump LocalStack logs on failure + if: failure() + run: make logs || true + + - name: Tear down + if: always() + run: make stop diff --git a/.gitignore b/.gitignore index f82a459..9c956da 100644 --- a/.gitignore +++ b/.gitignore @@ -104,3 +104,7 @@ venv.bak/ .mypy_cache/ .idea/ + +logs/runs +logs/*.log +logs/queries diff --git a/localstack-bruin/.bruin.yml b/localstack-bruin/.bruin.yml new file mode 100644 index 0000000..8422094 --- /dev/null +++ b/localstack-bruin/.bruin.yml @@ -0,0 +1,21 @@ +# Bruin connections, all pointing at LocalStack. The credentials are LocalStack's +# dummy defaults, so this file is safe to commit. +default_environment: local +environments: + local: + connections: + # LocalStack Snowflake. Bruin's native Snowflake connection always targets + # *.snowflakecomputing.com, so the Python assets get this connection injected + # as a secret and connect to the host below themselves. + snowflake: + - name: localstack-snowflake + account: test + username: test + password: test + database: SHOP + warehouse: TEST + generic: + - name: SNOWFLAKE_HOST + value: snowflake.localhost.localstack.cloud + - name: AWS_ENDPOINT_URL + value: http://localhost:4566 diff --git a/localstack-bruin/.gitignore b/localstack-bruin/.gitignore new file mode 100644 index 0000000..adc4018 --- /dev/null +++ b/localstack-bruin/.gitignore @@ -0,0 +1,7 @@ +__pycache__/ +logs/ +.venv/ + +# Bruin auto-adds `.bruin.yml` here; this sample commits it on purpose (dummy LocalStack creds only) +.bruin.yml +!.bruin.yml diff --git a/localstack-bruin/Makefile b/localstack-bruin/Makefile new file mode 100644 index 0000000..f12e255 --- /dev/null +++ b/localstack-bruin/Makefile @@ -0,0 +1,58 @@ +export AWS_ACCESS_KEY_ID ?= test +export AWS_SECRET_ACCESS_KEY ?= test +export AWS_DEFAULT_REGION ?= us-east-1 + +AWS := aws --endpoint-url=http://localhost:4566 +BRUIN := bruin +BRUIN_FLAGS := --config-file .bruin.yml + +.PHONY: help start wait seed validate lineage run report test all stop logs + +help: ## Show available targets + @grep -E '^[a-zA-Z_-]+:.*##' $(MAKEFILE_LIST) | \ + awk 'BEGIN{FS=":.*##"}{printf " %-10s %s\n", $$1, $$2}' + +start: ## Start LocalStack (AWS + Snowflake emulator) + docker compose up -d + @$(MAKE) --no-print-directory wait + +wait: ## Wait until S3 and Snowflake are available + @echo "Waiting for LocalStack..." + @for i in $$(seq 1 90); do \ + curl -s localhost:4566/_localstack/health | grep -q '"snowflake": "available"' && exit 0; sleep 2; \ + done; echo "Timed out waiting for LocalStack"; exit 1 + +seed: ## Create S3 buckets and upload the raw CSV files + $(AWS) s3 mb s3://shop-raw 2>/dev/null || true + $(AWS) s3 mb s3://shop-reports 2>/dev/null || true + $(AWS) s3 cp data/customers.csv s3://shop-raw/customers/customers.csv + $(AWS) s3 cp data/orders.csv s3://shop-raw/orders/orders.csv + +validate: ## Validate the Bruin pipeline + $(BRUIN) validate $(BRUIN_FLAGS) pipeline + +lineage: ## Show the lineage of the report asset + $(BRUIN) lineage --full pipeline/assets/revenue_by_country.py + +run: ## Run the full Bruin pipeline + $(BRUIN) run $(BRUIN_FLAGS) pipeline + +report: ## Print the report the pipeline wrote to S3 + $(AWS) s3 cp s3://shop-reports/revenue_by_country.json - + +test: ## Check the report in S3 has the expected totals + @$(AWS) s3 cp s3://shop-reports/revenue_by_country.json - | python3 -c 'import json, sys; \ + rows = {r["country"]: r for r in json.load(sys.stdin)}; \ + assert len(rows) == 7, f"expected 7 countries, got {len(rows)}"; \ + assert rows["US"]["orders"] == 4 and rows["US"]["revenue"] == 910.0, rows["US"]; \ + total = round(sum(r["revenue"] for r in rows.values()), 2); \ + assert total == 2372.16, f"unexpected total revenue {total}"; \ + print(f"OK: {len(rows)} countries, total revenue {total}")' + +all: start seed run report ## Start LocalStack, seed data, run the pipeline, show the report + +logs: ## Show LocalStack container logs + docker compose logs localstack + +stop: ## Stop and remove the LocalStack container + docker compose down diff --git a/localstack-bruin/README.md b/localstack-bruin/README.md new file mode 100644 index 0000000..2e518e5 --- /dev/null +++ b/localstack-bruin/README.md @@ -0,0 +1,76 @@ +# Bruin on LocalStack (AWS + Snowflake) + +A small [Bruin](https://getbruin.com) data pipeline that runs fully locally against +[LocalStack](https://localstack.cloud). Raw data lives in LocalStack S3, gets loaded and +transformed in the LocalStack Snowflake emulator, and the final report is written back +to LocalStack S3. No cloud accounts needed. + +## What the pipeline does + +``` + LocalStack S3 LocalStack Snowflake LocalStack S3 + ───────────── ──────────────────────────────────────────────── ───────────── + s3://shop-raw ─stage──► RAW.CUSTOMERS ─┐ + customers/*.csv + COPY RAW.ORDERS ─┴─► ANALYTICS.ORDERS_ENRICHED ──► ANALYTICS.REVENUE_BY_COUNTRY ──► s3://shop-reports + orders/*.csv (+ data quality checks) revenue_by_country.json +``` + +| Asset | File | What it shows | +|---|---|---| +| `raw.shop_data` | `raw_load.py` | Snowflake external stage on a LocalStack S3 bucket, `COPY INTO` from CSV | +| `analytics.orders_enriched` | `orders_enriched.py` | Snowflake CTAS join plus quality checks that fail the run on bad data | +| `analytics.revenue_by_country` | `revenue_by_country.py` | Snowflake aggregation, JSON report published to LocalStack S3 | + +Bruin handles the DAG (`depends`), injects the LocalStack connection details as secrets, +and gives each Python asset an isolated environment built with `uv` from +`pipeline/assets/requirements.txt`. + +## Prerequisites + +- Docker +- A `LOCALSTACK_AUTH_TOKEN` with access to the Snowflake emulator +- [Bruin CLI](https://getbruin.com/docs/bruin/getting-started/introduction/installation.html): + `curl -LsSf https://getbruin.com/install/cli | sh` +- AWS CLI (used by the Makefile to seed S3) + +## Quick start + +```bash +export LOCALSTACK_AUTH_TOKEN=ls-... +make all # start LocalStack, seed S3, run the pipeline, print the report +make stop # stop and remove the LocalStack container +``` + +Or step by step: + +```bash +make start # docker compose up (localstack/snowflake image: AWS + Snowflake on :4566) +make seed # create buckets, upload data/*.csv to s3://shop-raw +make validate # bruin validate +make lineage # bruin lineage for the report asset +make run # bruin run +make report # print s3://shop-reports/revenue_by_country.json +make test # assert the report has the expected totals (used in CI) +``` + +To see the quality checks stop the pipeline, upload an order for a customer that +doesn't exist and run again: + +```bash +(cat data/orders.csv; echo "1017,99,2025-06-09,Mouse,1,24.50") | \ + aws --endpoint-url=http://localhost:4566 s3 cp - s3://shop-raw/orders/orders.csv +make run # analytics.orders_enriched fails, revenue_by_country is skipped +make seed # restore the original data +``` + +## How Bruin is pointed at LocalStack + +All connections live in [`.bruin.yml`](.bruin.yml), using LocalStack's dummy `test` +credentials. + +Bruin's native Snowflake connection has no `host` option (its Go driver always connects +to `.snowflakecomputing.com`), so the Snowflake steps are Python assets. Each one +declares the `localstack-snowflake` connection under `secrets`, Bruin injects it as JSON, +and [`common.py`](pipeline/assets/common.py) connects with `snowflake-connector-python` +using `host=snowflake.localhost.localstack.cloud`, `port=4566`. The S3 endpoint for boto3 +comes from a `generic` connection in the same file. diff --git a/localstack-bruin/data/customers.csv b/localstack-bruin/data/customers.csv new file mode 100644 index 0000000..d7b7953 --- /dev/null +++ b/localstack-bruin/data/customers.csv @@ -0,0 +1,9 @@ +customer_id,name,country,signup_date +1,Alice Martin,DE,2025-01-12 +2,Bruno Costa,BR,2025-02-03 +3,Chen Wei,SG,2025-02-18 +4,Dana Levi,US,2025-03-01 +5,Emil Novak,AT,2025-03-22 +6,Fatima Zahra,MA,2025-04-09 +7,George Brown,US,2025-04-30 +8,Hana Sato,JP,2025-05-14 diff --git a/localstack-bruin/data/orders.csv b/localstack-bruin/data/orders.csv new file mode 100644 index 0000000..0de4db3 --- /dev/null +++ b/localstack-bruin/data/orders.csv @@ -0,0 +1,17 @@ +order_id,customer_id,order_date,product,quantity,unit_price +1001,1,2025-06-01,Keyboard,1,89.90 +1002,2,2025-06-01,Mouse,2,24.50 +1003,3,2025-06-02,Monitor,1,329.00 +1004,4,2025-06-02,Laptop Stand,1,49.00 +1005,1,2025-06-03,USB-C Hub,1,39.99 +1006,5,2025-06-03,Keyboard,1,89.90 +1007,6,2025-06-04,Headset,1,129.00 +1008,7,2025-06-04,Monitor,2,329.00 +1009,8,2025-06-05,Mouse,1,24.50 +1010,4,2025-06-05,Webcam,1,74.00 +1011,2,2025-06-06,Keyboard,1,89.90 +1012,3,2025-06-06,USB-C Hub,3,39.99 +1013,7,2025-06-07,Headset,1,129.00 +1014,5,2025-06-07,Laptop Stand,2,49.00 +1015,8,2025-06-08,Monitor,1,329.00 +1016,1,2025-06-08,Webcam,1,74.00 diff --git a/localstack-bruin/docker-compose.yml b/localstack-bruin/docker-compose.yml new file mode 100644 index 0000000..ea4d6a0 --- /dev/null +++ b/localstack-bruin/docker-compose.yml @@ -0,0 +1,11 @@ +services: + localstack: + container_name: localstack-bruin + image: localstack/snowflake:latest + ports: + - "127.0.0.1:4566:4566" + environment: + - LOCALSTACK_AUTH_TOKEN=${LOCALSTACK_AUTH_TOKEN:?LOCALSTACK_AUTH_TOKEN is required} + - DEBUG=${DEBUG:-0} + volumes: + - /var/run/docker.sock:/var/run/docker.sock diff --git a/localstack-bruin/pipeline/assets/common.py b/localstack-bruin/pipeline/assets/common.py new file mode 100644 index 0000000..bb8c971 --- /dev/null +++ b/localstack-bruin/pipeline/assets/common.py @@ -0,0 +1,31 @@ +"""Shared helpers for connecting the Python assets to LocalStack (Snowflake + AWS).""" + +import json +import os + +import boto3 +import snowflake.connector + + +def snowflake_connect(**kwargs): + # Bruin injects the `localstack-snowflake` connection as JSON (see `secrets` in each asset) + conn = json.loads(os.environ["SNOWFLAKE_CONN"]) + return snowflake.connector.connect( + host=os.environ["SNOWFLAKE_HOST"], + port=4566, + account=conn["account"], + user=conn["username"], + password=conn["password"], + warehouse=conn.get("warehouse"), + **kwargs, + ) + + +def aws_client(service): + return boto3.client( + service, + endpoint_url=os.environ["AWS_ENDPOINT_URL"], + aws_access_key_id="test", + aws_secret_access_key="test", + region_name="us-east-1", + ) diff --git a/localstack-bruin/pipeline/assets/orders_enriched.py b/localstack-bruin/pipeline/assets/orders_enriched.py new file mode 100644 index 0000000..27428de --- /dev/null +++ b/localstack-bruin/pipeline/assets/orders_enriched.py @@ -0,0 +1,54 @@ +"""@bruin +name: analytics.orders_enriched +type: python +description: | + Joins orders with customers in LocalStack Snowflake and computes line revenue, + then runs a few data quality checks that fail the pipeline if violated. + +depends: + - raw.shop_data + +secrets: + - key: localstack-snowflake + inject_as: SNOWFLAKE_CONN + - key: SNOWFLAKE_HOST +@bruin""" + +from .common import snowflake_connect + +# check name -> query returning the number of offending rows +CHECKS = { + "order_id is unique": "SELECT COUNT(*) - COUNT(DISTINCT order_id) FROM ORDERS_ENRICHED", + "every order has a customer": "SELECT COUNT(*) FROM ORDERS_ENRICHED WHERE country IS NULL", + "revenue is positive": "SELECT COUNT(*) FROM ORDERS_ENRICHED WHERE revenue <= 0", +} + + +def main(): + with snowflake_connect(database="SHOP") as sf, sf.cursor() as cur: + cur.execute("CREATE SCHEMA IF NOT EXISTS SHOP.ANALYTICS") + cur.execute("USE SCHEMA SHOP.ANALYTICS") + cur.execute( + """CREATE OR REPLACE TABLE ORDERS_ENRICHED AS + SELECT o.order_id, o.order_date, o.customer_id, + c.name AS customer_name, c.country, + o.product, o.quantity, o.unit_price, + o.quantity * o.unit_price AS revenue + FROM SHOP.RAW.ORDERS o + LEFT JOIN SHOP.RAW.CUSTOMERS c ON c.customer_id = o.customer_id""" + ) + + failed = [] + for name, query in CHECKS.items(): + cur.execute(query) + bad_rows = cur.fetchone()[0] + print(f"[{'PASS' if bad_rows == 0 else 'FAIL'}] {name}") + if bad_rows: + failed.append(name) + + if failed: + raise SystemExit(f"Data quality checks failed: {', '.join(failed)}") + + +if __name__ == "__main__": + main() diff --git a/localstack-bruin/pipeline/assets/raw_load.py b/localstack-bruin/pipeline/assets/raw_load.py new file mode 100644 index 0000000..d18005f --- /dev/null +++ b/localstack-bruin/pipeline/assets/raw_load.py @@ -0,0 +1,45 @@ +"""@bruin +name: raw.shop_data +type: python +description: | + Loads the raw CSV files from the LocalStack S3 data lake into LocalStack Snowflake, + using an external stage on the S3 bucket and COPY INTO. + +secrets: + - key: localstack-snowflake + inject_as: SNOWFLAKE_CONN + - key: SNOWFLAKE_HOST +@bruin""" + +from .common import snowflake_connect + +CSV_FORMAT = "FILE_FORMAT = (TYPE = CSV SKIP_HEADER = 1)" + + +def main(): + with snowflake_connect() as sf, sf.cursor() as cur: + for stmt in [ + "CREATE DATABASE IF NOT EXISTS SHOP", + "CREATE SCHEMA IF NOT EXISTS SHOP.RAW", + "USE SCHEMA SHOP.RAW", + # External stage pointing at the LocalStack S3 bucket + """CREATE OR REPLACE STAGE raw_stage + URL = 's3://shop-raw/' + CREDENTIALS = (AWS_KEY_ID = 'test' AWS_SECRET_KEY = 'test')""", + """CREATE OR REPLACE TABLE CUSTOMERS ( + customer_id INT, name VARCHAR, country VARCHAR, signup_date DATE)""", + """CREATE OR REPLACE TABLE ORDERS ( + order_id INT, customer_id INT, order_date DATE, product VARCHAR, + quantity INT, unit_price NUMBER(10, 2))""", + f"COPY INTO CUSTOMERS FROM @raw_stage/customers/ {CSV_FORMAT}", + f"COPY INTO ORDERS FROM @raw_stage/orders/ {CSV_FORMAT}", + ]: + cur.execute(stmt) + + for table in ("CUSTOMERS", "ORDERS"): + cur.execute(f"SELECT COUNT(*) FROM {table}") + print(f"Loaded {cur.fetchone()[0]} rows into SHOP.RAW.{table}") + + +if __name__ == "__main__": + main() diff --git a/localstack-bruin/pipeline/assets/requirements.txt b/localstack-bruin/pipeline/assets/requirements.txt new file mode 100644 index 0000000..7a6b1c5 --- /dev/null +++ b/localstack-bruin/pipeline/assets/requirements.txt @@ -0,0 +1,2 @@ +snowflake-connector-python>=3.12 +boto3>=1.35 diff --git a/localstack-bruin/pipeline/assets/revenue_by_country.py b/localstack-bruin/pipeline/assets/revenue_by_country.py new file mode 100644 index 0000000..3d4fdbb --- /dev/null +++ b/localstack-bruin/pipeline/assets/revenue_by_country.py @@ -0,0 +1,52 @@ +"""@bruin +name: analytics.revenue_by_country +type: python +description: | + Builds a revenue-by-country mart in LocalStack Snowflake, then publishes it as a + JSON report to the LocalStack S3 reports bucket. + +depends: + - analytics.orders_enriched + +secrets: + - key: localstack-snowflake + inject_as: SNOWFLAKE_CONN + - key: SNOWFLAKE_HOST + - key: AWS_ENDPOINT_URL +@bruin""" + +import json + +from .common import aws_client, snowflake_connect + +REPORT_BUCKET = "shop-reports" +REPORT_KEY = "revenue_by_country.json" + + +def main(): + with snowflake_connect(database="SHOP", schema="ANALYTICS") as sf, sf.cursor() as cur: + cur.execute( + """CREATE OR REPLACE TABLE REVENUE_BY_COUNTRY AS + SELECT country, + COUNT(DISTINCT customer_id) AS customers, + COUNT(*) AS orders, + ROUND(SUM(revenue), 2) AS revenue + FROM ORDERS_ENRICHED + GROUP BY country""" + ) + cur.execute("SELECT country, customers, orders, revenue FROM REVENUE_BY_COUNTRY ORDER BY revenue DESC") + rows = [ + {"country": c, "customers": cu, "orders": o, "revenue": float(r)} + for c, cu, o, r in cur.fetchall() + ] + + for row in rows: + print(f"{row['country']:>4} orders={row['orders']:<3} revenue={row['revenue']:>9.2f}") + + s3 = aws_client("s3") + s3.put_object(Bucket=REPORT_BUCKET, Key=REPORT_KEY, Body=json.dumps(rows, indent=2)) + print(f"Report written to s3://{REPORT_BUCKET}/{REPORT_KEY}") + + +if __name__ == "__main__": + main() diff --git a/localstack-bruin/pipeline/pipeline.yml b/localstack-bruin/pipeline/pipeline.yml new file mode 100644 index 0000000..4578af8 --- /dev/null +++ b/localstack-bruin/pipeline/pipeline.yml @@ -0,0 +1 @@ +name: shop-analytics From 80aa85c6bfc91045e62496ebe6c324d9c0efa372 Mon Sep 17 00:00:00 2001 From: Waldemar Hummer Date: Mon, 5 Oct 2026 04:19:53 +0900 Subject: [PATCH 2/3] localstack-bruin: use Snowflake-licensed token in CI, fail fast on missing license Co-Authored-By: Claude Opus 5.5 --- .github/workflows/test-localstack-bruin.yml | 2 +- localstack-bruin/Makefile | 5 ++++- 2 files changed, 5 insertions(+), 2 deletions(-) diff --git a/.github/workflows/test-localstack-bruin.yml b/.github/workflows/test-localstack-bruin.yml index aa59d33..e165159 100644 --- a/.github/workflows/test-localstack-bruin.yml +++ b/.github/workflows/test-localstack-bruin.yml @@ -14,7 +14,7 @@ on: workflow_dispatch: env: - LOCALSTACK_AUTH_TOKEN: ${{ secrets.TEST_LOCALSTACK_AUTH_TOKEN }} + LOCALSTACK_AUTH_TOKEN: ${{ secrets.TEST_LOCALSTACK_SNOWFLAKE_AUTH_TOKEN }} jobs: test-localstack-bruin: diff --git a/localstack-bruin/Makefile b/localstack-bruin/Makefile index f12e255..2b43e6b 100644 --- a/localstack-bruin/Makefile +++ b/localstack-bruin/Makefile @@ -19,7 +19,10 @@ start: ## Start LocalStack (AWS + Snowflake emulator) wait: ## Wait until S3 and Snowflake are available @echo "Waiting for LocalStack..." @for i in $$(seq 1 90); do \ - curl -s localhost:4566/_localstack/health | grep -q '"snowflake": "available"' && exit 0; sleep 2; \ + curl -s localhost:4566/_localstack/health | grep -q '"snowflake": "available"' && exit 0; \ + if docker compose logs localstack 2>/dev/null | grep -q "not covered by your license"; then \ + echo "Your LOCALSTACK_AUTH_TOKEN does not include the Snowflake emulator"; exit 1; fi; \ + sleep 2; \ done; echo "Timed out waiting for LocalStack"; exit 1 seed: ## Create S3 buckets and upload the raw CSV files From 673583d359e3b6ba50d4843b589c7536073cb2f7 Mon Sep 17 00:00:00 2001 From: Waldemar Hummer Date: Mon, 5 Oct 2026 04:47:13 +0900 Subject: [PATCH 3/3] localstack-bruin: switch to snowflake-next emulator Run LocalStack (S3) and localstack/snowflake-next side by side, with external stages reaching S3 via SF_S3_ENDPOINT. This drops the earlier emulator workarounds: COPY uses a stage-level file format with MATCH_BY_COLUMN_NAME, and the report is unloaded from Snowflake to S3 with COPY INTO @stage instead of boto3. Ports are configurable via LOCALSTACK_PORT / SNOWFLAKE_PORT. Co-Authored-By: Claude Opus 5.5 --- .github/workflows/test-localstack-bruin.yml | 2 +- localstack-bruin/.bruin.yml | 6 +-- localstack-bruin/Makefile | 29 +++++++------ localstack-bruin/README.md | 42 +++++++++++++------ localstack-bruin/docker-compose.yml | 24 ++++++++--- localstack-bruin/pipeline/assets/common.py | 17 ++------ .../pipeline/assets/orders_enriched.py | 1 + localstack-bruin/pipeline/assets/raw_load.py | 11 ++--- .../pipeline/assets/requirements.txt | 1 - .../pipeline/assets/revenue_by_country.py | 41 +++++++++--------- 10 files changed, 101 insertions(+), 73 deletions(-) diff --git a/.github/workflows/test-localstack-bruin.yml b/.github/workflows/test-localstack-bruin.yml index e165159..c5037ce 100644 --- a/.github/workflows/test-localstack-bruin.yml +++ b/.github/workflows/test-localstack-bruin.yml @@ -18,7 +18,7 @@ env: jobs: test-localstack-bruin: - name: Bruin pipeline on LocalStack S3 + Snowflake + name: Bruin pipeline on LocalStack S3 + Snowflake (snowflake-next) runs-on: ubuntu-latest timeout-minutes: 20 defaults: diff --git a/localstack-bruin/.bruin.yml b/localstack-bruin/.bruin.yml index 8422094..f66522a 100644 --- a/localstack-bruin/.bruin.yml +++ b/localstack-bruin/.bruin.yml @@ -6,7 +6,7 @@ environments: connections: # LocalStack Snowflake. Bruin's native Snowflake connection always targets # *.snowflakecomputing.com, so the Python assets get this connection injected - # as a secret and connect to the host below themselves. + # as a secret and connect to the host/port below themselves. snowflake: - name: localstack-snowflake account: test @@ -17,5 +17,5 @@ environments: generic: - name: SNOWFLAKE_HOST value: snowflake.localhost.localstack.cloud - - name: AWS_ENDPOINT_URL - value: http://localhost:4566 + - name: SNOWFLAKE_PORT + value: ${SNOWFLAKE_PORT} diff --git a/localstack-bruin/Makefile b/localstack-bruin/Makefile index 2b43e6b..2ce77bb 100644 --- a/localstack-bruin/Makefile +++ b/localstack-bruin/Makefile @@ -1,8 +1,10 @@ export AWS_ACCESS_KEY_ID ?= test export AWS_SECRET_ACCESS_KEY ?= test export AWS_DEFAULT_REGION ?= us-east-1 +export LOCALSTACK_PORT ?= 4566 +export SNOWFLAKE_PORT ?= 4567 -AWS := aws --endpoint-url=http://localhost:4566 +AWS := aws --endpoint-url=http://localhost:$(LOCALSTACK_PORT) BRUIN := bruin BRUIN_FLAGS := --config-file .bruin.yml @@ -12,15 +14,16 @@ help: ## Show available targets @grep -E '^[a-zA-Z_-]+:.*##' $(MAKEFILE_LIST) | \ awk 'BEGIN{FS=":.*##"}{printf " %-10s %s\n", $$1, $$2}' -start: ## Start LocalStack (AWS + Snowflake emulator) +start: ## Start LocalStack (S3) and the Snowflake emulator docker compose up -d @$(MAKE) --no-print-directory wait -wait: ## Wait until S3 and Snowflake are available +wait: ## Wait until LocalStack S3 and the Snowflake emulator are available @echo "Waiting for LocalStack..." @for i in $$(seq 1 90); do \ - curl -s localhost:4566/_localstack/health | grep -q '"snowflake": "available"' && exit 0; \ - if docker compose logs localstack 2>/dev/null | grep -q "not covered by your license"; then \ + curl -s localhost:$(LOCALSTACK_PORT)/_localstack/health | grep -q '"s3": "available"' && \ + curl -s localhost:$(SNOWFLAKE_PORT)/_localstack/health | grep -q '"snowflake":"available"' && exit 0; \ + if docker compose logs 2>/dev/null | grep -qi "not covered by your license"; then \ echo "Your LOCALSTACK_AUTH_TOKEN does not include the Snowflake emulator"; exit 1; fi; \ sleep 2; \ done; echo "Timed out waiting for LocalStack"; exit 1 @@ -40,22 +43,22 @@ lineage: ## Show the lineage of the report asset run: ## Run the full Bruin pipeline $(BRUIN) run $(BRUIN_FLAGS) pipeline -report: ## Print the report the pipeline wrote to S3 +report: ## Print the report Snowflake unloaded to S3 $(AWS) s3 cp s3://shop-reports/revenue_by_country.json - test: ## Check the report in S3 has the expected totals @$(AWS) s3 cp s3://shop-reports/revenue_by_country.json - | python3 -c 'import json, sys; \ - rows = {r["country"]: r for r in json.load(sys.stdin)}; \ + rows = {r["COUNTRY"]: r for r in map(json.loads, filter(str.strip, sys.stdin))}; \ assert len(rows) == 7, f"expected 7 countries, got {len(rows)}"; \ - assert rows["US"]["orders"] == 4 and rows["US"]["revenue"] == 910.0, rows["US"]; \ - total = round(sum(r["revenue"] for r in rows.values()), 2); \ + assert rows["US"]["ORDERS"] == 4 and float(rows["US"]["REVENUE"]) == 910.0, rows["US"]; \ + total = round(sum(float(r["REVENUE"]) for r in rows.values()), 2); \ assert total == 2372.16, f"unexpected total revenue {total}"; \ print(f"OK: {len(rows)} countries, total revenue {total}")' -all: start seed run report ## Start LocalStack, seed data, run the pipeline, show the report +all: start seed run report ## Start everything, seed data, run the pipeline, show the report -logs: ## Show LocalStack container logs - docker compose logs localstack +logs: ## Show the container logs + docker compose logs -stop: ## Stop and remove the LocalStack container +stop: ## Stop and remove the containers docker compose down diff --git a/localstack-bruin/README.md b/localstack-bruin/README.md index 2e518e5..6af48d6 100644 --- a/localstack-bruin/README.md +++ b/localstack-bruin/README.md @@ -2,29 +2,44 @@ A small [Bruin](https://getbruin.com) data pipeline that runs fully locally against [LocalStack](https://localstack.cloud). Raw data lives in LocalStack S3, gets loaded and -transformed in the LocalStack Snowflake emulator, and the final report is written back -to LocalStack S3. No cloud accounts needed. +transformed in the LocalStack Snowflake emulator, and the final report is unloaded from +Snowflake back to LocalStack S3. No cloud accounts needed. ## What the pipeline does ``` LocalStack S3 LocalStack Snowflake LocalStack S3 ───────────── ──────────────────────────────────────────────── ───────────── - s3://shop-raw ─stage──► RAW.CUSTOMERS ─┐ - customers/*.csv + COPY RAW.ORDERS ─┴─► ANALYTICS.ORDERS_ENRICHED ──► ANALYTICS.REVENUE_BY_COUNTRY ──► s3://shop-reports - orders/*.csv (+ data quality checks) revenue_by_country.json + s3://shop-raw ─stage──► RAW.CUSTOMERS ─┐ s3://shop-reports + customers/*.csv + COPY RAW.ORDERS ─┴─► ANALYTICS.ORDERS_ENRICHED revenue_by_country.json + orders/*.csv (+ data quality checks) ▲ + │ │ COPY INTO @stage + ▼ │ (unload) + ANALYTICS.REVENUE_BY_COUNTRY ─────────────┘ ``` | Asset | File | What it shows | |---|---|---| -| `raw.shop_data` | `raw_load.py` | Snowflake external stage on a LocalStack S3 bucket, `COPY INTO` from CSV | +| `raw.shop_data` | `raw_load.py` | Snowflake external stage on a LocalStack S3 bucket, `COPY INTO ... MATCH_BY_COLUMN_NAME` from CSV | | `analytics.orders_enriched` | `orders_enriched.py` | Snowflake CTAS join plus quality checks that fail the run on bad data | -| `analytics.revenue_by_country` | `revenue_by_country.py` | Snowflake aggregation, JSON report published to LocalStack S3 | +| `analytics.revenue_by_country` | `revenue_by_country.py` | Snowflake aggregation, unloaded as JSON to LocalStack S3 with `COPY INTO @stage` | Bruin handles the DAG (`depends`), injects the LocalStack connection details as secrets, and gives each Python asset an isolated environment built with `uv` from `pipeline/assets/requirements.txt`. +## Setup + +[`docker-compose.yml`](docker-compose.yml) starts two containers: + +| Container | Image | Host port | Purpose | +|---|---|---|---| +| `localstack-bruin-aws` | `localstack/localstack-pro` | `LOCALSTACK_PORT` (4566) | S3 | +| `localstack-bruin-snowflake` | `localstack/snowflake-next` | `SNOWFLAKE_PORT` (4567) | Snowflake emulator | + +The Snowflake emulator reaches S3 for external stages through +`SF_S3_ENDPOINT=http://localstack:4566` on the compose network. + ## Prerequisites - Docker @@ -37,14 +52,14 @@ and gives each Python asset an isolated environment built with `uv` from ```bash export LOCALSTACK_AUTH_TOKEN=ls-... -make all # start LocalStack, seed S3, run the pipeline, print the report -make stop # stop and remove the LocalStack container +make all # start containers, seed S3, run the pipeline, print the report +make stop # stop and remove the containers ``` Or step by step: ```bash -make start # docker compose up (localstack/snowflake image: AWS + Snowflake on :4566) +make start # docker compose up, wait for S3 and Snowflake make seed # create buckets, upload data/*.csv to s3://shop-raw make validate # bruin validate make lineage # bruin lineage for the report asset @@ -53,6 +68,10 @@ make report # print s3://shop-reports/revenue_by_country.json make test # assert the report has the expected totals (used in CI) ``` +If the default ports are taken, override them for every `make` call, for example +`export LOCALSTACK_PORT=4578 SNOWFLAKE_PORT=4577`. When running `bruin` directly instead +of through `make`, export `SNOWFLAKE_PORT` too, since `.bruin.yml` reads it. + To see the quality checks stop the pipeline, upload an order for a customer that doesn't exist and run again: @@ -72,5 +91,4 @@ Bruin's native Snowflake connection has no `host` option (its Go driver always c to `.snowflakecomputing.com`), so the Snowflake steps are Python assets. Each one declares the `localstack-snowflake` connection under `secrets`, Bruin injects it as JSON, and [`common.py`](pipeline/assets/common.py) connects with `snowflake-connector-python` -using `host=snowflake.localhost.localstack.cloud`, `port=4566`. The S3 endpoint for boto3 -comes from a `generic` connection in the same file. +using `host=snowflake.localhost.localstack.cloud` and the `SNOWFLAKE_PORT`. diff --git a/localstack-bruin/docker-compose.yml b/localstack-bruin/docker-compose.yml index ea4d6a0..718577e 100644 --- a/localstack-bruin/docker-compose.yml +++ b/localstack-bruin/docker-compose.yml @@ -1,11 +1,23 @@ services: + # LocalStack for the AWS side (S3 data lake + reports bucket) localstack: - container_name: localstack-bruin - image: localstack/snowflake:latest + container_name: localstack-bruin-aws + image: localstack/localstack-pro:latest ports: - - "127.0.0.1:4566:4566" + - "127.0.0.1:${LOCALSTACK_PORT:-4566}:4566" environment: - LOCALSTACK_AUTH_TOKEN=${LOCALSTACK_AUTH_TOKEN:?LOCALSTACK_AUTH_TOKEN is required} - - DEBUG=${DEBUG:-0} - volumes: - - /var/run/docker.sock:/var/run/docker.sock + - SERVICES=s3 + + # LocalStack Snowflake emulator (next generation) + snowflake: + container_name: localstack-bruin-snowflake + image: localstack/snowflake-next:latest + ports: + - "127.0.0.1:${SNOWFLAKE_PORT:-4567}:4566" + environment: + - LOCALSTACK_AUTH_TOKEN=${LOCALSTACK_AUTH_TOKEN:?LOCALSTACK_AUTH_TOKEN is required} + # External S3 stages are read from / unloaded to the LocalStack container above + - SF_S3_ENDPOINT=http://localstack:4566 + depends_on: + - localstack diff --git a/localstack-bruin/pipeline/assets/common.py b/localstack-bruin/pipeline/assets/common.py index bb8c971..cde1889 100644 --- a/localstack-bruin/pipeline/assets/common.py +++ b/localstack-bruin/pipeline/assets/common.py @@ -1,31 +1,22 @@ -"""Shared helpers for connecting the Python assets to LocalStack (Snowflake + AWS).""" +"""Shared helper for connecting the Python assets to the LocalStack Snowflake emulator.""" import json import os -import boto3 import snowflake.connector +DEFAULT_SNOWFLAKE_PORT = 4567 + def snowflake_connect(**kwargs): # Bruin injects the `localstack-snowflake` connection as JSON (see `secrets` in each asset) conn = json.loads(os.environ["SNOWFLAKE_CONN"]) return snowflake.connector.connect( host=os.environ["SNOWFLAKE_HOST"], - port=4566, + port=int(os.environ.get("SNOWFLAKE_PORT") or DEFAULT_SNOWFLAKE_PORT), account=conn["account"], user=conn["username"], password=conn["password"], warehouse=conn.get("warehouse"), **kwargs, ) - - -def aws_client(service): - return boto3.client( - service, - endpoint_url=os.environ["AWS_ENDPOINT_URL"], - aws_access_key_id="test", - aws_secret_access_key="test", - region_name="us-east-1", - ) diff --git a/localstack-bruin/pipeline/assets/orders_enriched.py b/localstack-bruin/pipeline/assets/orders_enriched.py index 27428de..1d1832b 100644 --- a/localstack-bruin/pipeline/assets/orders_enriched.py +++ b/localstack-bruin/pipeline/assets/orders_enriched.py @@ -12,6 +12,7 @@ - key: localstack-snowflake inject_as: SNOWFLAKE_CONN - key: SNOWFLAKE_HOST + - key: SNOWFLAKE_PORT @bruin""" from .common import snowflake_connect diff --git a/localstack-bruin/pipeline/assets/raw_load.py b/localstack-bruin/pipeline/assets/raw_load.py index d18005f..0b26df5 100644 --- a/localstack-bruin/pipeline/assets/raw_load.py +++ b/localstack-bruin/pipeline/assets/raw_load.py @@ -9,12 +9,11 @@ - key: localstack-snowflake inject_as: SNOWFLAKE_CONN - key: SNOWFLAKE_HOST + - key: SNOWFLAKE_PORT @bruin""" from .common import snowflake_connect -CSV_FORMAT = "FILE_FORMAT = (TYPE = CSV SKIP_HEADER = 1)" - def main(): with snowflake_connect() as sf, sf.cursor() as cur: @@ -22,17 +21,19 @@ def main(): "CREATE DATABASE IF NOT EXISTS SHOP", "CREATE SCHEMA IF NOT EXISTS SHOP.RAW", "USE SCHEMA SHOP.RAW", + "CREATE OR REPLACE FILE FORMAT csv_format TYPE = CSV PARSE_HEADER = TRUE", # External stage pointing at the LocalStack S3 bucket """CREATE OR REPLACE STAGE raw_stage URL = 's3://shop-raw/' - CREDENTIALS = (AWS_KEY_ID = 'test' AWS_SECRET_KEY = 'test')""", + CREDENTIALS = (AWS_KEY_ID = 'test' AWS_SECRET_KEY = 'test') + FILE_FORMAT = csv_format""", """CREATE OR REPLACE TABLE CUSTOMERS ( customer_id INT, name VARCHAR, country VARCHAR, signup_date DATE)""", """CREATE OR REPLACE TABLE ORDERS ( order_id INT, customer_id INT, order_date DATE, product VARCHAR, quantity INT, unit_price NUMBER(10, 2))""", - f"COPY INTO CUSTOMERS FROM @raw_stage/customers/ {CSV_FORMAT}", - f"COPY INTO ORDERS FROM @raw_stage/orders/ {CSV_FORMAT}", + "COPY INTO CUSTOMERS FROM @raw_stage/customers/ MATCH_BY_COLUMN_NAME = CASE_INSENSITIVE", + "COPY INTO ORDERS FROM @raw_stage/orders/ MATCH_BY_COLUMN_NAME = CASE_INSENSITIVE", ]: cur.execute(stmt) diff --git a/localstack-bruin/pipeline/assets/requirements.txt b/localstack-bruin/pipeline/assets/requirements.txt index 7a6b1c5..8f5e50f 100644 --- a/localstack-bruin/pipeline/assets/requirements.txt +++ b/localstack-bruin/pipeline/assets/requirements.txt @@ -1,2 +1 @@ snowflake-connector-python>=3.12 -boto3>=1.35 diff --git a/localstack-bruin/pipeline/assets/revenue_by_country.py b/localstack-bruin/pipeline/assets/revenue_by_country.py index 3d4fdbb..d05b288 100644 --- a/localstack-bruin/pipeline/assets/revenue_by_country.py +++ b/localstack-bruin/pipeline/assets/revenue_by_country.py @@ -2,7 +2,7 @@ name: analytics.revenue_by_country type: python description: | - Builds a revenue-by-country mart in LocalStack Snowflake, then publishes it as a + Builds a revenue-by-country mart in LocalStack Snowflake, then unloads it as a JSON report to the LocalStack S3 reports bucket. depends: @@ -12,15 +12,13 @@ - key: localstack-snowflake inject_as: SNOWFLAKE_CONN - key: SNOWFLAKE_HOST - - key: AWS_ENDPOINT_URL + - key: SNOWFLAKE_PORT @bruin""" -import json +from .common import snowflake_connect -from .common import aws_client, snowflake_connect - -REPORT_BUCKET = "shop-reports" -REPORT_KEY = "revenue_by_country.json" +REPORT_URL = "s3://shop-reports/" +REPORT_FILE = "revenue_by_country.json" def main(): @@ -34,18 +32,23 @@ def main(): FROM ORDERS_ENRICHED GROUP BY country""" ) - cur.execute("SELECT country, customers, orders, revenue FROM REVENUE_BY_COUNTRY ORDER BY revenue DESC") - rows = [ - {"country": c, "customers": cu, "orders": o, "revenue": float(r)} - for c, cu, o, r in cur.fetchall() - ] - - for row in rows: - print(f"{row['country']:>4} orders={row['orders']:<3} revenue={row['revenue']:>9.2f}") - - s3 = aws_client("s3") - s3.put_object(Bucket=REPORT_BUCKET, Key=REPORT_KEY, Body=json.dumps(rows, indent=2)) - print(f"Report written to s3://{REPORT_BUCKET}/{REPORT_KEY}") + cur.execute("SELECT country, orders, revenue FROM REVENUE_BY_COUNTRY ORDER BY revenue DESC") + for country, orders, revenue in cur.fetchall(): + print(f"{country:>4} orders={orders:<3} revenue={revenue:>9.2f}") + + # Unload the mart straight from Snowflake to LocalStack S3 (one JSON object per line) + cur.execute( + f"""CREATE OR REPLACE STAGE reports_stage + URL = '{REPORT_URL}' + CREDENTIALS = (AWS_KEY_ID = 'test' AWS_SECRET_KEY = 'test')""" + ) + cur.execute( + f"""COPY INTO @reports_stage/{REPORT_FILE} + FROM (SELECT OBJECT_CONSTRUCT(*) FROM REVENUE_BY_COUNTRY ORDER BY revenue DESC) + FILE_FORMAT = (TYPE = JSON COMPRESSION = NONE) + SINGLE = TRUE OVERWRITE = TRUE""" + ) + print(f"Report unloaded to {REPORT_URL}{REPORT_FILE}") if __name__ == "__main__":