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
ソースの最終更新活動
2026年5月22日 22:58
検出された SKILL.md の言語
英語
スター
5
フォーク
1

インストール方法

デフォルトでは、最初にソースを確認する 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で見る