Skip to main content

amee-joshi-data-engineering-portfolio

Reference portfolio demonstrating Azure data engineering patterns, Medallion architecture, and end-to-end analytics solutions

معلومات المصدر

المستودع
reason-machines/data-skills
آخر نشاط في المصدر
٢٢ مايو ٢٠٢٦ في ٢٢:٥٨
لغة SKILL.md المكتشفة
الإنجليزية
النجوم
٥
التفرعات
١

خيارات التثبيت

يُحدَّد Prompt الذي يراجع المصدر أولًا بشكل افتراضي. يمكنك التبديل إلى أمر مباشر أو تنزيل نسخة محلية.

مراجعة ملفات المصدر

اقرأ SKILL.md وأي ملفات مرافقة يعرضها SkillsMP قبل أن تقرر التثبيت.

عرض SKILL.md

SKILL.md
تعليمات المصدر · معاينة للقراءة فقط
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
ملف SKILL.md هذا كبير جدا، لذلك يعرض SkillsMP القسم الاول فقط هنا. عرض على GitHub