Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
58 changes: 58 additions & 0 deletions .github/workflows/test-localstack-bruin.yml
Original file line number Diff line number Diff line change
@@ -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_SNOWFLAKE_AUTH_TOKEN }}

jobs:
test-localstack-bruin:
name: Bruin pipeline on LocalStack S3 + Snowflake (snowflake-next)
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
4 changes: 4 additions & 0 deletions .gitignore
Original file line number Diff line number Diff line change
Expand Up @@ -104,3 +104,7 @@ venv.bak/
.mypy_cache/

.idea/

logs/runs
logs/*.log
logs/queries
21 changes: 21 additions & 0 deletions localstack-bruin/.bruin.yml
Original file line number Diff line number Diff line change
@@ -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/port 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: SNOWFLAKE_PORT
value: ${SNOWFLAKE_PORT}
7 changes: 7 additions & 0 deletions localstack-bruin/.gitignore
Original file line number Diff line number Diff line change
@@ -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
64 changes: 64 additions & 0 deletions localstack-bruin/Makefile
Original file line number Diff line number Diff line change
@@ -0,0 +1,64 @@
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:$(LOCALSTACK_PORT)
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 (S3) and the Snowflake emulator
docker compose up -d
@$(MAKE) --no-print-directory wait

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:$(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

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 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 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 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 everything, seed data, run the pipeline, show the report

logs: ## Show the container logs
docker compose logs

stop: ## Stop and remove the containers
docker compose down
94 changes: 94 additions & 0 deletions localstack-bruin/README.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,94 @@
# 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 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 ─┐ 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 ... 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, 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
- 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 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, 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
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)
```

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:

```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 `<account>.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` and the `SNOWFLAKE_PORT`.
9 changes: 9 additions & 0 deletions localstack-bruin/data/customers.csv
Original file line number Diff line number Diff line change
@@ -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
17 changes: 17 additions & 0 deletions localstack-bruin/data/orders.csv
Original file line number Diff line number Diff line change
@@ -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
23 changes: 23 additions & 0 deletions localstack-bruin/docker-compose.yml
Original file line number Diff line number Diff line change
@@ -0,0 +1,23 @@
services:
# LocalStack for the AWS side (S3 data lake + reports bucket)
localstack:
container_name: localstack-bruin-aws
image: localstack/localstack-pro:latest
ports:
- "127.0.0.1:${LOCALSTACK_PORT:-4566}:4566"
environment:
- LOCALSTACK_AUTH_TOKEN=${LOCALSTACK_AUTH_TOKEN:?LOCALSTACK_AUTH_TOKEN is required}
- 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
22 changes: 22 additions & 0 deletions localstack-bruin/pipeline/assets/common.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,22 @@
"""Shared helper for connecting the Python assets to the LocalStack Snowflake emulator."""

import json
import os

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=int(os.environ.get("SNOWFLAKE_PORT") or DEFAULT_SNOWFLAKE_PORT),
account=conn["account"],
user=conn["username"],
password=conn["password"],
warehouse=conn.get("warehouse"),
**kwargs,
)
55 changes: 55 additions & 0 deletions localstack-bruin/pipeline/assets/orders_enriched.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,55 @@
"""@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
- key: SNOWFLAKE_PORT
@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()
Loading
Loading