| 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 — 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:
from pyspark.sql import SparkSession
from delta.tables import DeltaTable
df_raw = (spark.read
.format("parquet")
.load(f"{bronze_path}/source_data/")
.withColumn("ingestion_timestamp", current_timestamp())
.withColumn("source_file", input_file_name())
)
(df_raw.write
.format("delta")
.mode("append")
.option("mergeSchema", "true")
.save(f"{bronze_path}/retail_transactions")
)
from pyspark.sql.functions import col, when, trim, upper
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")
)
(df_silver.write
.format("delta")
.mode("overwrite")
.option("overwriteSchema", "false")
.save(f"{silver_path}/transactions")
)
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
source_prepared = (source_df
.withColumn("effective_date", current_timestamp())
.withColumn("end_date", lit(None).cast("timestamp"))
.withColumn("is_current", lit(True))
)
target_delta = DeltaTable.forPath(spark, target_table)
merge_condition = " AND ".join([f"target.{k} = source.{k}" for k in key_columns])
(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:
{
"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:
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}}'
"""
},
: {
: ,
: {
: ,
:
}
}
}
}
3. Star Schema Data Warehouse
Pattern: Dimensional Modeling with SQL Server
Dimension Table (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)
);
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 TARGET
(product_id, product_name, category, subcategory, unit_price)
(source.product_id, source.product_name, source.category,
source.subcategory, source.unit_price);
Dimension Table (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)
);
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)
);
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,
end_date,
is_current
staging.customers s
dim_customer d s.customer_id d.customer_id d.is_current
d.customer_key
s.city d.city
s.state d.state;
Fact Table:
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 NONCLUSTERED COLUMNSTORE INDEX idx_fact_sales_cs
fact_sales (date_key, customer_key, product_key, store_key,
quantity, unit_price, total_amount);
fact_sales (
date_key, customer_key, product_key, store_key,
quantity, unit_price, discount_amount, tax_amount, total_amount
)
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
staging.transactions st
dim_date dd (st.transaction_date ) dd.date
dim_customer dc st.customer_id dc.customer_id dc.is_current
dim_product dp st.product_id dp.product_id
dim_store ds st.store_id ds.store_id;
4. Incremental Data Loading Pattern
Watermark-Based Incremental Load:
from pyspark.sql.functions import col, max as spark_max
from delta.tables import DeltaTable
source_table = "source_database.transactions"
target_path = "/mnt/silver/transactions"
watermark_table = "control.watermark"
watermark_column = "modified_date"
last_watermark = (spark.table(watermark_table)
.filter(col("table_name") == source_table)
.select("watermark_value")
.first()[0]
)
df_incremental = (spark.table(source_table)
.filter(col(watermark_column) > last_watermark)
)
if DeltaTable.isDeltaTable(spark, target_path):
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:
(df_incremental.write
.format("delta")
.mode("overwrite")
.save(target_path)
)
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:
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
self.table_name = table_name
self.quality_results = []
def check_null_values(self, columns):
"""Check for null values in critical columns"""
for column in columns:
null_count = self.df.filter(col(column).isNull()).count()
total_count = self.df.count()
self.quality_results.append({
"check_type": "null_check",
"column": column,
"null_count": null_count,
"total_count": total_count,
"null_percentage": (null_count / total_count * 100) if total_count > 0 else 0,
"passed": null_count == 0
})
return self
def check_duplicates(self, key_columns):
"""Check for duplicate records"""
duplicate_count = (.df
.groupBy(key_columns)
.count()
.(col() > )
.count()
)
.quality_results.append({
: ,
: key_columns,
: duplicate_count,
: duplicate_count ==
})
():
missing_references = (.df
.select(foreign_key)
.distinct()
.join(reference_df.select(reference_key),
col(foreign_key) == col(reference_key),
)
.count()
)
.quality_results.append({
: ,
: foreign_key,
: missing_references,
: missing_references ==
})
():
out_of_range = .df.(
(col(column) < min_value min_value ) |
(col(column) > max_value max_value )
).count()
.quality_results.append({
: ,
: column,
: min_value,
: max_value,
: out_of_range,
: out_of_range ==
})
():
.quality_results
df_transactions = spark.read.().load()
df_customers = spark.read.().load()
quality_checker = DataQualityChecker(df_transactions, )
results = (quality_checker
.check_null_values([, , ])
.check_duplicates([])
.check_referential_integrity(, df_customers, )
.check_value_range(, min_value=, max_value=)
.get_results()
)
result results:
()
Power BI Analytics Patterns
DAX Measures for KPIs:
// Total Sales
Total Sales = SUM(fact_sales[total_amount])
// Year-over-Year Growth
Sales YoY Growth =
VAR CurrentYearSales = [Total Sales]
VAR PreviousYearSales =
CALCULATE(
[Total Sales],
DATEADD(dim_date[Date], -1, YEAR)
)
RETURN
DIVIDE(
CurrentYearSales - PreviousYearSales,
PreviousYearSales,
0
)
// Customer Lifetime Value
Customer LTV =
CALCULATE(
[Total Sales],
ALLEXCEPT(dim_customer, dim_customer[customer_id])
)
// Moving Average (3 months)
Sales 3M MA =
CALCULATE(
[Total Sales],
DATESINPERIOD(
dim_date[Date],
LASTDATE(dim_date[Date]),
-3,
MONTH
)
) / 3
// Rank by Sales
Product Sales Rank =
RANKX(
ALL(dim_product[product_name]),
[Total Sales],
,
DESC,
DENSE
)
Common Architectural Patterns
Medallion Architecture Best Practices
Bronze Layer:
- Raw data ingestion with minimal transformation
- Add audit columns (ingestion_timestamp, source_file)
- Preserve source schema with schema evolution enabled
- Partition by ingestion date for performance
Silver Layer:
- Data cleansing and standardization
- Deduplication based on business keys
- Data type conversions and validations
- Enforce schema constraints
- Join related datasets
Gold Layer:
- Business-aggregated datasets
- Dimensional models (Star/Snowflake schema)
- Pre-calculated metrics and KPIs
- Optimized for BI tool consumption
Delta Lake Optimization
from delta.tables import DeltaTable
deltaTable = DeltaTable.forPath(spark, "/mnt/gold/fact_sales")
deltaTable.optimize().executeZOrderBy("date_key", "customer_key")
deltaTable.vacuum(168)
spark.sql("ANALYZE TABLE gold.fact_sales COMPUTE STATISTICS FOR ALL COLUMNS")
Unity Catalog Security
CREATE CATALOG IF NOT EXISTS retail_analytics;
CREATE SCHEMA IF NOT EXISTS retail_analytics.gold;
GRANT USE CATALOG ON CATALOG retail_analytics TO `data_analysts`;
GRANT USE SCHEMA ON SCHEMA retail_analytics.gold TO `data_analysts`;
GRANT SELECT ON TABLE retail_analytics.gold.fact_sales TO `data_analysts`;
CREATE FUNCTION retail_analytics.gold.customer_filter(customer_region STRING)
RETURN customer_region = current_user_region();
ALTER TABLE retail_analytics.gold.fact_sales
SET ROW FILTER retail_analytics.gold.customer_filter ON (region);
Environment Setup
Azure Configuration:
export AZURE_SUBSCRIPTION_ID=your_subscription_id
export AZURE_RESOURCE_GROUP=rg-data-platform
export AZURE_STORAGE_ACCOUNT=datalakestorage
export AZURE_DATABRICKS_WORKSPACE=databricks-workspace
export ADF_FACTORY_NAME=adf-data-ingestion
export ADF_LINKED_SERVICE_NAME=ls-sqlserver-source
Databricks Configuration:
configs = {
"fs.azure.account.auth.type": "OAuth",
"fs.azure.account.oauth.provider.type": "org.apache.hadoop.fs.azurebfs.oauth2.ClientCredsTokenProvider",
"fs.azure.account.oauth2.client.id": dbutils.secrets.get(scope="keyvault", key="client-id"),
"fs.azure.account.oauth2.client.secret": dbutils.secrets.get(scope="keyvault", key="client-secret"),
"fs.azure.account.oauth2.client.endpoint": f"https://login.microsoftonline.com/{dbutils.secrets.get(scope='keyvault', key='tenant-id')}/oauth2/token"
}
dbutils.fs.mount(
source = "abfss://bronze@datalakestorage.dfs.core.windows.net/",
mount_point = "/mnt/bronze",
extra_configs = configs
)
Troubleshooting
Issue: Delta Lake merge taking too long
from delta.tables import DeltaTable
target_table = DeltaTable.forPath(spark, target_path)
target_table.optimize().executeCompaction()
spark.sql(f"""
ALTER TABLE delta.`{target_path}`
SET TBLPROPERTIES (
delta.autoOptimize.optimizeWrite = true,
delta.autoOptimize.autoCompact = true
)
""")
Issue: ADF pipeline timeout
{
"typeProperties": {
"timeout": "0.12:00:00"
},
"policy": {
"timeout": "7.00:00:00",
"retry": 2,
"retryIntervalInSeconds": 30
}
}
Issue: Power BI slow refresh
// Use incremental refresh configuration
// In Power BI Desktop: Table Tools > Incremental Refresh
// Or optimize DAX measures
Optimized Total Sales =
CALCULATE(
SUM(fact_sales[total_amount]),
KEEPFILTERS(dim_date[Date]) // Reduce context transition overhead
)
Issue: Schema evolution conflicts
(df.write
.format("delta")
.mode("append")
.option("mergeSchema", "true")
.save(target_path)
)
(df.write
.format("delta")
.mode("overwrite")
.option("overwriteSchema", "true")
.save(target_path)
)
Reference Architecture
This portfolio demonstrates a typical enterprise data platform architecture:
┌─────────────────┐
│ Source Systems │
│ (SQL Server, │
│ APIs, Files) │
└────────┬────────┘
│
▼
┌─────────────────┐
│ Azure Data │
│ Factory (ADF) │ ◄──── Metadata-driven ingestion
└────────┬────────┘
│
▼
┌─────────────────┐
│ ADLS Gen2 │
│ Bronze Layer │ ◄──── Raw data landing
└────────┬────────┘
│
▼
┌─────────────────┐
│ Databricks │
│ Silver Layer │ ◄──── Cleansing & transformation
└────────┬────────┘
│
▼
┌─────────────────┐
│ Databricks │
│ Gold Layer │ ◄──── Analytics-ready datasets
└────────┬────────┘
│
├──────────────────┐
▼ ▼
┌─────────────────┐ ┌─────────────────┐
│ Power BI │ │ Synapse │
│ Reporting │ │ Analytics │
└─────────────────┘ └─────────────────┘
This skill provides patterns and code examples for building production-grade data platforms following industry best practices demonstrated across the portfolio projects.