- name
- amee-joshi-data-engineering-portfolio
- description
- Reference portfolio demonstrating Azure data engineering patterns, Medallion architecture, and end-to-end analytics solutions
- triggers
- ["show me data engineering portfolio examples","how to build azure data pipelines","what is medallion architecture implementation","data engineering project structure examples","how to design lakehouse architecture","azure databricks end-to-end patterns","data warehouse modeling examples","ETL pipeline design patterns"]
# Amee Joshi Data Engineering Portfolio
> Skill by [ara.so](https://ara.so) — Data Skills collection.
This portfolio showcases production-grade data engineering patterns and architectures for building scalable, cloud-native data platforms. It demonstrates end-to-end solutions covering data ingestion, transformation, modeling, and analytics using Azure services, Databricks, SQL Server, and BI tools.
## What This Portfolio Demonstrates
This is a reference collection showing:
- **Medallion Architecture (Bronze-Silver-Gold)** implementations
- **Azure cloud data platforms** (ADF, ADLS Gen2, Databricks, Synapse Analytics)
- **Data lakehouse patterns** with Delta Lake
- **Dimensional modeling** (Star Schema, SCD Type 1 & 2)
- **Metadata-driven ingestion frameworks**
- **Analytics-ready datasets** for BI consumption
- **ETL/ELT pipeline design** with incremental loading
- **Power BI and Tableau** reporting solutions
## Key Portfolio Projects
### 1. Azure Databricks Retail Lakehouse
**Repository:** `azure-databricks-end-to-end-retail-lakehouse`
**Pattern:** Enterprise Medallion Architecture with Delta Lake
**Architecture:**
```
Bronze (Raw) → Silver (Cleansed) → Gold (Analytics-Ready)
```
**Key Implementation Concepts:**
```python
# Bronze Layer - Raw Ingestion
from pyspark.sql import SparkSession
from delta.tables import DeltaTable
# Ingest raw data with metadata
df_raw = (spark.read
.format("parquet")
.load(f"{bronze_path}/source_data/")
.withColumn("ingestion_timestamp", current_timestamp())
.withColumn("source_file", input_file_name())
)
# Write to Bronze Delta table
(df_raw.write
.format("delta")
.mode("append")
.option("mergeSchema", "true")
.save(f"{bronze_path}/retail_transactions")
)
```
```python
# Silver Layer - Data Quality & Transformation
from pyspark.sql.functions import col, when, trim, upper
# Cleanse and standardize
df_silver = (df_bronze
.filter(col("transaction_id").isNotNull())
.withColumn("customer_name", trim(upper(col("customer_name"))))
.withColumn("transaction_amount",
when(col("transaction_amount") < 0, 0)
.otherwise(col("transaction_amount")))
.dropDuplicates(["transaction_id"])
.select("transaction_id", "customer_id", "product_id",
"transaction_amount", "transaction_date")
)
# Write with schema enforcement
(df_silver.write
.format("delta")
.mode("overwrite")
.option("overwriteSchema", "false")
.save(f"{silver_path}/transactions")
)
```
```python
# Gold Layer - SCD Type 2 Dimension
def apply_scd_type2(target_table, source_df, key_columns, scd_columns):
"""
Implements Slowly Changing Dimension Type 2
"""
from delta.tables import DeltaTable
from pyspark.sql.functions import lit, current_timestamp
# Prepare source with SCD metadata
source_prepared = (source_df
.withColumn("effective_date", current_timestamp())
.withColumn("end_date", lit(None).cast("timestamp"))
.withColumn("is_current", lit(True))
)
# Read existing target
target_delta = DeltaTable.forPath(spark, target_table)
# Identify changes
merge_condition = " AND ".join([f"target.{k} = source.{k}" for k in key_columns])
# Perform SCD Type 2 merge
(target_delta.alias("target")
.merge(source_prepared.alias("source"), merge_condition)
.whenMatchedUpdate(
condition = "target.is_current = true AND " +
" OR ".join([f"target.{c} != source.{c}" for c in scd_columns]),
set = {
"is_current": "false",
"end_date": "current_timestamp()"
}
)
.whenNotMatchedInsertAll()
.execute()
)
```
### 2. Metadata-Driven Ingestion Framework
**Pattern:** Dynamic, configuration-based pipeline generation
**Configuration Schema:**
```json
{
"pipeline_config": {
"source_system": "SQL_SERVER",
"target_layer": "bronze",
"ingestion_type": "incremental",
"watermark_column": "modified_date",
"tables": [
{
"schema_name": "sales",
"table_name": "orders",
"partition_column": "order_date",
"primary_key": ["order_id"],
"target_path": "/bronze/sales/orders"
}
]
}
}
```
**Azure Data Factory Pattern:**
```python
# Dynamic pipeline parameter processing
# This represents the logic implemented in ADF
def generate_copy_activity(table_config):
"""
Generates ADF copy activity from metadata
"""
return {
"name": f"Copy_{table_config['table_name']}",
"type": "Copy",
"inputs": [{
"referenceName": "SourceDataset",
"type": "DatasetReference",
"parameters": {
"schemaName": table_config['schema_name'],
"tableName": table_config['table_name']
}
}],
"outputs": [{
"referenceName": "SinkDataset",
"type": "DatasetReference",
"parameters": {
"targetPath": table_config['target_path']
}
}],
"typeProperties": {
"source": {
"type": "SqlServerSource",
"sqlReaderQuery": f"""
SELECT * FROM {table_config['schema_name']}.{table_config['table_name']}
WHERE {table_config['watermark_column']} > '@{{pipeline().parameters.watermarkValue}}'
"""
},
"sink": {
"type": "ParquetSink",
"storeSettings": {
"type": "AzureBlobFSWriteSettings",
"copyBehavior": "PreserveHierarchy"
}
}
}
}
```
### 3. Star Schema Data Warehouse
**Pattern:** Dimensional Modeling with SQL Server
**Dimension Table (SCD Type 1):**
```sql
-- Dimension: Product (SCD Type 1)
CREATE TABLE dim_product (
product_key INT IDENTITY(1,1) PRIMARY KEY,
product_id INT NOT NULL,
product_name NVARCHAR(100),
category NVARCHAR(50),
subcategory NVARCHAR(50),
unit_price DECIMAL(10,2),
modified_date DATETIME DEFAULT GETDATE(),
CONSTRAINT uk_product UNIQUE (product_id)
);
-- ETL Merge (SCD Type 1 - Overwrite)
MERGE INTO dim_product AS target
USING (
SELECT
product_id,
product_name,
category,
subcategory,
unit_price
FROM staging.products
) AS source
ON target.product_id = source.product_id
WHEN MATCHED AND (
target.product_name != source.product_name OR
target.category != source.category OR
target.unit_price != source.unit_price
)
THEN UPDATE SET
target.product_name = source.product_name,
target.category = source.category,
target.subcategory = source.subcategory,
target.unit_price = source.unit_price,
target.modified_date = GETDATE()
WHEN NOT MATCHED BY TARGET
THEN INSERT (product_id, product_name, category, subcategory, unit_price)
VALUES (source.product_id, source.product_name, source.category,
source.subcategory, source.unit_price);
```
**Dimension Table (SCD Type 2):**
```sql
-- Dimension: Customer (SCD Type 2)
CREATE TABLE dim_customer (
customer_key INT IDENTITY(1,1) PRIMARY KEY,
customer_id INT NOT NULL,
customer_name NVARCHAR(100),
email NVARCHAR(100),
city NVARCHAR(50),
state NVARCHAR(50),
effective_date DATETIME NOT NULL,
end_date DATETIME NULL,
is_current BIT DEFAULT 1,
CONSTRAINT uk_customer_current UNIQUE (customer_id, is_current)
);
-- ETL for SCD Type 2
-- Step 1: Expire changed records
UPDATE dim_customer
SET
end_date = GETDATE(),
is_current = 0
WHERE customer_id IN (
SELECT s.customer_id
FROM staging.customers s
INNER JOIN dim_customer d ON s.customer_id = d.customer_id
WHERE d.is_current = 1
AND (s.city != d.city OR s.state != d.state)
);
-- Step 2: Insert new versions
INSERT INTO dim_customer (
customer_id, customer_name, email, city, state,
effective_date, end_date, is_current
)
SELECT
s.customer_id,
s.customer_name,
s.email,
s.city,
s.state,
GETDATE() AS effective_date,
NULL AS end_date,
1 AS is_current
FROM staging.customers s
LEFT JOIN dim_customer d ON s.customer_id = d.customer_id AND d.is_current = 1
WHERE d.customer_key IS NULL
OR s.city != d.city
OR s.state != d.state;
```
**Fact Table:**
```sql
-- Fact: Sales Transactions
CREATE TABLE fact_sales (
sales_key BIGINT IDENTITY(1,1) PRIMARY KEY,
date_key INT NOT NULL,
customer_key INT NOT NULL,
product_key INT NOT NULL,
store_key INT NOT NULL,
quantity INT NOT NULL,
unit_price DECIMAL(10,2) NOT NULL,
discount_amount DECIMAL(10,2) DEFAULT 0,
tax_amount DECIMAL(10,2) DEFAULT 0,
total_amount DECIMAL(10,2) NOT NULL,
CONSTRAINT fk_date FOREIGN KEY (date_key) REFERENCES dim_date(date_key),
CONSTRAINT fk_customer FOREIGN KEY (customer_key) REFERENCES dim_customer(customer_key),
CONSTRAINT fk_product FOREIGN KEY (product_key) REFERENCES dim_product(product_key),
CONSTRAINT fk_store FOREIGN KEY (store_key) REFERENCES dim_store(store_key)
);
-- Create columnstore index for analytics
CREATE NONCLUSTERED COLUMNSTORE INDEX idx_fact_sales_cs
ON fact_sales (date_key, customer_key, product_key, store_key,
quantity, unit_price, total_amount);
-- ETL Load
INSERT INTO fact_sales (
date_key, customer_key, product_key, store_key,
quantity, unit_price, discount_amount, tax_amount, total_amount
)
SELECT
dd.date_key,
dc.customer_key,
dp.product_key,
ds.store_key,
st.quantity,
st.unit_price,
st.discount_amount,
st.tax_amount,
st.total_amount
FROM staging.transactions st
INNER JOIN dim_date dd ON CAST(st.transaction_date AS DATE) = dd.date
INNER JOIN dim_customer dc ON st.customer_id = dc.customer_id AND dc.is_current = 1
INNER JOIN dim_product dp ON st.product_id = dp.product_id
INNER JOIN dim_store ds ON st.store_id = ds.store_id;
```
### 4. Incremental Data Loading Pattern
**Watermark-Based Incremental Load:**
```python
# Databricks notebook - Incremental load with watermark
from pyspark.sql.functions import col, max as spark_max
from delta.tables import DeltaTable
# Configuration
source_table = "source_database.transactions"
target_path = "/mnt/silver/transactions"
watermark_table = "control.watermark"
watermark_column = "modified_date"
# Get last watermark
last_watermark = (spark.table(watermark_table)
.filter(col("table_name") == source_table)
.select("watermark_value")
.first()[0]
)
# Read incremental data
df_incremental = (spark.table(source_table)
.filter(col(watermark_column) > last_watermark)
)
# Check if target exists
if DeltaTable.isDeltaTable(spark, target_path):
# Merge into existing table
target_table = DeltaTable.forPath(spark, target_path)
(target_table.alias("target")
.merge(
df_incremental.alias("source"),
"target.transaction_id = source.transaction_id"
)
.whenMatchedUpdateAll()
.whenNotMatchedInsertAll()
.execute()
)
else:
# Initial load
(df_incremental.write
.format("delta")
.mode("overwrite")
.save(target_path)
)
# Update watermark
new_watermark = df_incremental.agg(spark_max(watermark_column)).first()[0]
spark.sql(f"""
UPDATE {watermark_table}
SET watermark_value = '{new_watermark}',
last_updated = current_timestamp()
WHERE table_name = '{source_table}'
""")
```
### 5. Data Quality Framework
**Quality Checks Pattern:**
```python
from pyspark.sql.functions import col, count, sum as spark_sum, when
class DataQualityChecker:
"""
Data quality validation framework
"""
def __init__(self, dataframe, table_name):
self.df = dataframe
GitHubで見る