An end-to-end governed data lakehouse that turns raw inventory events into actionable insights on shrinkage, spoilage, and vendor short-shipments. Raw CSV files land in GCS → move through a controlled ELT process in BigQuery → are served via business views → and are eventually archived to Iceberg. OpenMetadata captures lineage across the entire flow
- Architecture at a Glance
- The Business Problem
- Personal Motivation
- Solution Overview
- What the Pipeline Builds
- Getting Started
- Tool Stack
- Suggested Improvements
Retail businesses lose a significant portion of revenue every year to inventory shrinkage (the gap between recorded stock and actual physical stock). Shrinkage comes from several sources:
- Theft (internal and external)
- Administrative and process errors
- Vendor short-shipments
- Spoilage and damage
- Unrecorded returns or adjustments
Even a 1–2 % shrinkage rate can translate into millions of naira (or dollars) in lost profit. Traditional approaches rely on periodic physical counts and spreadsheet reconciliations. By the time discrepancies surface, the root causes have already compounded across stores, regions, and product categories.
Finance, loss-prevention, and operations teams need timely, trustworthy answers to questions such as:
- Which retail store locations are experiencing unusual inventory shrinkage?
- What is the financial impact of inventory losses caused by vendor short-shipments?
- What total monetary value of inventory losses needs to be categorized and reported?
Without a reliable pipeline that captures, validates, transforms, and archives inventory events — while preserving full lineage and data-quality rules — these questions remain difficult or impossible to answer with confidence.
This project was inspired by my experience in audit. While examining inventory records I repeatedly saw how shrinkage, spoilage, and simple recording errors can quietly erode a retailer’s financial health. Small discrepancies that look insignificant in isolation become material when aggregated across hundreds of stores and thousands of SKUs.
I wanted to design a modern data solution that gives retail organisations continuous visibility into these losses, rather than discovering them only at the end of a quarter or during a physical stocktake. ShrinkGuard is the result of that motivation — a production-style lakehouse pipeline that turns inventory events into governed, queryable insights.
ShrinkGuard implements a complete medallion-style lakehouse on Google Cloud, orchestrated by Apache Airflow and governed by OpenMetadata:
| Layer | Technology | Purpose |
|---|---|---|
| Orchestration | Apache Airflow 2.11.2 | Three idempotent DAGs that generate, transform, and archive data |
| Storage | Google Cloud Storage | Landing, temporary, archival, and Iceberg zones |
| Warehouse | BigQuery | Landing → Staging → Fact tables + dimension tables + business views |
| Archive | Apache Iceberg + Lakekeeper | Long-term, retention-policy-driven archival of aged fact data |
| Governance | OpenMetadata 1.13.1 | End-to-end lineage, schema catalog, and data-quality rules |
Every load is logged. Every transformation is retry-safe. Aged data is moved to Iceberg on a retention policy you control. The entire system is catalogued so lineage from raw CSV to business view is always visible.
Uploads slowly-changing dimension data (stores, products, vendors, etc.) to GCS and loads it into BigQuery dim_* tables. Run this once before the daily pipelines.
Simulates the business producing inventory events.
Writes exactly three deterministically-named CSV files per day into the landing/ zone (e.g. ind_YYYYMMDD_1.csv).
Intentionally includes some bad records so data-validation logic can be exercised.
Idempotent: a re-run overwrites the same three files.
The core of the solution:
- Load landing files into BigQuery.
- Transform into a run-scoped staging table (
CREATE OR REPLACEso retries are self-healing). - Merge/load into a partitioned fact table (with optional reject/failed-records handling).
- Move processed CSVs from
landing/toarchival/. - Log every step (status, timestamps, row counts) to a central log table.
- Refresh three business views that answer the key retail questions listed above.
Landing and staging tables are temporary and auto-expire after 7 days.
Selects aged fact data (configurable retention window — e.g. rows older than 90–180 days) and moves it into Iceberg tables managed by Lakekeeper.
Re-archiving a period overwrites rather than duplicates, keeping the archive clean.
- Docker & Docker Compose
- A Google Cloud project with billing enabled
- A GCS bucket
- A service-account JSON key with the permissions listed below
- Git
git clone <your-repo-url>
cd shrinkguardcd docker-compose/config
cp settings_template.py settings.py
# Edit settings.py and fill in your GCS bucket, BigQuery project/dataset,
# retention window, and any other environment-specific values.# Copy your service-account key into the plugins folder
cp /path/to/your-service-account.json docker-compose/plugins/gcp-key.jsonRequired permissions for the service account (minimum for metadata extraction + usage):
bigquery.datasets.get
bigquery.tables.get
bigquery.tables.getData
bigquery.tables.list
resourcemanager.projects.get
bigquery.jobs.create
bigquery.jobs.listAll
Optional (policy tags / full usage & lineage):
datacatalog.taxonomies.get
datacatalog.taxonomies.list
bigquery.readsessions.create
bigquery.readsessions.getData
docker compose -f docker-compose/docker-compose-omd.yaml up -d
# Wait ~60 seconds for all services to become healthyOpen http://localhost:8585
Login: admin@open-metadata.org / admin
- Go to Settings → Services → Bots → LineageBot
- Copy the JWT token
- Paste it into
docker-compose/.env:
AIRFLOW_OPENLINEAGE_TRANSPORT={"type":"http","url":"http://openmetadata-server:8585/api/v1/openlineage/","endpoint":"lineage","auth":{"type":"api_key","api_key":"<PASTE-TOKEN-HERE>"}}Before starting the Second docker compose, run this in your terminal to generate secret key:
python -c "import secrets; print(secrets.token_urlsafe(32))"- Copy the generated key
- Paste it into
docker-compose/.env:
AIRFLOW__WEBSERVER__SECRET_KEY=INSERT_YOUR_SECRET_KEY_HERESave .env and then run
docker compose -f docker-compose/docker-compose-airflow-lakekeeper.yaml up -d
# Wait ~60 secondsOpen http://localhost:8080 (Airflow UI)
Default credentials: admin / admin
Enable all DAGs.
Trigger DAG 0 manually once. The remaining DAGs run on their schedules after you have made them active (@daily / @monthly).
In OpenMetadata UI → Settings → Services → Pipelines → Add New Service
- Name:
local_airflow - Host Port:
http://airflow-webserver:8080 - Auth: Basic Auth (
admin/admin) - Test connection → Save
Settings → Services → Databases → Add New Service → BigQuery
- Credentials Type: Service Account
- Path:
/opt/openmetadata/plugins/gcp-key.json - Project ID: your GCP project
- Test connection (warnings on optional permissions are acceptable) → Save
# Watch OpenLineage events leaving Airflow
docker logs -f airflow-scheduler | grep -i openlineage
# Watch OpenMetadata receiving lineage
docker logs -f openmetadata_server | grep -i -E "openlineage|lineage"
# Troubleshoot Airflow webserver
docker logs airflow-webserver --tail=50To stop and clean up the local demo environment started above:
# Stop the Airflow + Lakekeeper stack
docker compose -f docker-compose/docker-compose-airflow-lakekeeper.yaml down
# Stop the OpenMetadata stack
docker compose -f docker-compose/docker-compose-omd.yaml downIf you want to remove all containers, networks, and volumes created by compose:
docker compose -f docker-compose/docker-compose-airflow-lakekeeper.yaml down --volumes
docker compose -f docker-compose/docker-compose-omd.yaml down --volumesIf you also want to remove the local images built by compose:
docker compose -f docker-compose/docker-compose-airflow-lakekeeper.yaml down --rmi all --volumes
docker compose -f docker-compose/docker-compose-omd.yaml down --rmi all --volumes| Component | Version / Notes |
|---|---|
| Apache Airflow | 2.11.2 |
| Lakekeeper | 0.13.1 (Iceberg REST catalog) |
| OpenMetadata | 1.13.1 |
| Google Cloud Storage | Landing / temp / archival / iceberg zones |
| BigQuery | Landing → Staging → (Fact + dimensions) → views |
| Docker / Docker Compose | Local development & demo environment |
-
Business Intelligence Layer
Connect the three business views directly to Power BI, Tableau, or Looker Studio. This turns the governed lakehouse into an executive dashboard that loss-prevention and finance teams can use daily. -
Infrastructure as Code
Write a Terraform module that provisions:- GCS bucket + lifecycle policies
- BigQuery datasets and tables
- Service accounts and the exact IAM roles required
- Any networking or Secret Manager entries
This removes manual setup friction and makes the project reproducible across environments.
-
Alerting
Add Airflow or OpenMetadata alerts when shrinkage exceeds a configurable threshold for any store or category. -
Data Quality Expansion
Formalise additional Great Expectations / OpenMetadata test suites for completeness, referential integrity between fact and dimensions, and freshness of the landing zone.
ShrinkGuard — because every unit of inventory that disappears without a trace is a unit of profit that could have been protected.




