
Key Features & Highlights
Technologies & Architecture
FX Rate Analytics Pipeline
About Project
An automated, enterprise-grade Foreign Exchange Rate Analytics platform built on a Medallion Data Warehouse Architecture using Microsoft SQL Server, Apache Spark, dbt, Apache Airflow, Scikit-Learn, FastAPI, and Next.js.
The platform ingests live spot exchange rates from dual independent APIs into a Bronze raw audit layer, transforms and flattens heterogeneous JSON schemas via PySpark into a standardized Silver layer with idempotent T-SQL MERGE upserts, and generates business-ready Gold analytics marts for rolling volatility, Month-over-Month (MoM) performance, and Isolation Forest ML anomaly detection.
System Architecture
The system ingests live foreign exchange data from two independent REST APIs (Frankfurter and Open Exchange Rates). Python handles API ingestion and raw JSON preservation, while PySpark cleans, flattens, and standardizes disparate schema formats into a staging table.
Microsoft SQL Server hosts the Bronze, Silver, and Gold medallion layers. dbt transforms Silver data into business-ready analytical marts, while an unsupervised Isolation Forest model detects unusual currency fluctuations. Apache Airflow 2.10 orchestrates the entire 7-step pipeline on a scheduled cron, and a high-performance FastAPI backend serves analytics to a real-time Next.js dashboard.
Data Architecture — Medallion Pattern
The data platform adheres strictly to the Medallion Data Architecture:
• Bronze Layer (`bronze.raw_api_responses`): Stores raw, unmodified API response payloads with ingestion timestamps for full auditability and replayability. • Silver Layer (`silver.rates`): Cleaned, deduplicated, and unified currency pairs with standardized column definitions across both API sources. • Gold Layer: Houses business-ready analytics including `gold.daily_volatility` (7-day rolling moving averages & volatility %), `gold.currency_performance` (monthly aggregated rates & MoM % growth), and `gold.anomaly_alerts` (ML-scored outlier events).
Pipeline Orchestration & Airflow DAG
Apache Airflow 2.10 orchestrates the end-to-end workflow through an automated 7-step DAG (`fx_ingestion_dag`):
1. `fetch_fx_rates`: Extracts rates from both external APIs. 2. `load_bronze`: Inserts raw JSON payloads into MSSQL Bronze. 3. `spark_transform_bronze_to_staging`: PySpark flattens currency arrays and writes to a staging table. 4. `merge_staging_to_silver`: Executes an atomic T-SQL `MERGE` to update Silver. 5. `dbt_run`: Materializes Gold volatility and monthly performance models. 6. `dbt_test`: Validates schema integrity, non-null, and unique constraints. 7. `detect_anomalies`: Runs Isolation Forest anomaly scoring and updates the Gold alerts table.
Data Processing & Spark Staging-to-Merge Pattern
Because the two FX APIs produce completely different JSON response structures (one nested with date objects, one flat with epoch timestamps), PySpark is utilized to parse and standardize them into a unified tabular schema.
Since Spark's native JDBC writer only supports append and overwrite modes (lacking native UPSERT capabilities), an architectural pattern was designed where Spark writes batch records to `staging.stg_fx_rates`, followed by an atomic T-SQL `MERGE` into `silver.rates`. This delivers distributed transformation speed combined with relational transactional guarantees.
Unsupervised ML Anomaly Detection
The pipeline employs Scikit-Learn's Isolation Forest algorithm to detect anomalous exchange rate swings without requiring historical labeled training data.
Rather than comparing raw rates across different pairs (which would bias high-value pairs like USD/JPY against low-value pairs like USD/EUR), the model is trained per currency pair using percentage deviation from its own 7-day rolling moving average and rolling volatility. Flagged anomalies are scored and categorized into High, Medium, and Low severity tiers.
Interactive Financial Cockpit & Request Flow
The frontend is built with Next.js 14 App Router, TypeScript, and Tailwind CSS. It connects to FastAPI using Server-Side Rendering (SSR) for initial loads and dynamic client-side fetching for real-time currency exploration.
Key UI features include a 6×2 Country Currency Selector driven by ISO country metadata, interactive Recharts glowing gradient area charts with moving average overlays, a dbt Gold MoM performance mart table, an instant Live FX conversion tool, and one-click CSV export.
Key Engineering Decisions & Learnings
• End-to-End Idempotency: Engineered every pipeline stage so that repeated executions on the same day never produce duplicate records or corrupted stats. • Cross-Platform Containerization: Standardized the complete database, Airflow, backend, and dashboard microservices on Docker Compose with SQL Server Authentication, overcoming Linux/Windows networking boundaries. • Automated CI/CD: Implemented GitHub Actions workflows executing automated Flake8 linting and dbt model validation tests on every push.