Skip to content

About

An end-to-end ELT lakehouse pipeline that detects inventory shrinkage, protects retail margins, and delivers audit-ready analytics using Airflow, BigQuery, Apache Iceberg, and OpenMetadata.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Latest commit

 

History

8 Commits

Folders and files

NameName
Last commit message
Last commit date
 
 
 
 
 
 
 
 

Repository files navigation

ShrinkGuard — Lakehouse Pipeline for Retail Inventory Shrinkage

Docker Airflow Python GCS Apache Iceberg Lakekeeper BigQuery OpenMetadata

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

Table of Contents

Architecture at a Glance

ShrinkGuard Lakehouse Architecture


The Business Problem

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:

  1. Which retail store locations are experiencing unusual inventory shrinkage?
  2. What is the financial impact of inventory losses caused by vendor short-shipments?
  3. 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.


Personal Motivation

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.


Solution Overview

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.


What the Pipeline Builds

DAG 0 — One-time Dimension Load (manual)

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.

Setup demension data in GCS Bucket and push to BigQuery

DAG 1 — Generate Landing Data (daily)

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.

Generate Inventory data and load into GCS Bucket and then push to BigQuery

DAG 2 — Full ELT Pipeline (daily)

The core of the solution:

  1. Load landing files into BigQuery.
  2. Transform into a run-scoped staging table (CREATE OR REPLACE so retries are self-healing).
  3. Merge/load into a partitioned fact table (with optional reject/failed-records handling).
  4. Move processed CSVs from landing/ to archival/.
  5. Log every step (status, timestamps, row counts) to a central log table.
  6. Refresh three business views that answer the key retail questions listed above.

Landing and staging tables are temporary and auto-expire after 7 days.

Core ELT pipeline

DAG 3 — Archive to Iceberg (monthly)

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.

Moving Old records to Lakekeeper


Getting Started

Prerequisites

  • Docker & Docker Compose
  • A Google Cloud project with billing enabled
  • A GCS bucket
  • A service-account JSON key with the permissions listed below
  • Git

1. Clone the repository

git clone <your-repo-url>
cd shrinkguard

2. Configure application settings

cd 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.

3. Place your GCP credentials

# Copy your service-account key into the plugins folder
cp /path/to/your-service-account.json docker-compose/plugins/gcp-key.json

Required 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

4. Start OpenMetadata first

docker compose -f docker-compose/docker-compose-omd.yaml up -d
# Wait ~60 seconds for all services to become healthy

Open 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>"}}

5. Start the Airflow + Lakekeeper stack

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_HERE

Save .env and then run

docker compose -f docker-compose/docker-compose-airflow-lakekeeper.yaml up -d
# Wait ~60 seconds

Open 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).

6. Connect Airflow to OpenMetadata

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

7. Connect BigQuery to OpenMetadata

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

Useful monitoring commands

# 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=50

Tear Down

To 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 down

If 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 --volumes

If 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

Tool Stack

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

Suggested Improvements

  1. 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.

  2. 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.

  3. Alerting
    Add Airflow or OpenMetadata alerts when shrinkage exceeds a configurable threshold for any store or category.

  4. 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.

About

An end-to-end ELT lakehouse pipeline that detects inventory shrinkage, protects retail margins, and delivers audit-ready analytics using Airflow, BigQuery, Apache Iceberg, and OpenMetadata.

Resources

Stars

0 stars

Watchers

0 watching

Forks

Releases

Packages

Contributors

Languages