Sh.
0%
FX Rate Analytics Pipeline preview 1

Key Features & Highlights

Dual API Ingestion — Ingests and normalizes live FX spot rates from Frankfurter and Open Exchange Rates APIs.
Medallion Architecture — Full separation of concerns across Bronze (raw audit), Silver (standardized rates), and Gold (analytical marts).
Distributed PySpark Transformations — Batch normalizes heterogeneous JSON schemas and executes bulk staging writes via JDBC.
Idempotent T-SQL MERGE — Guarantees exactly-once daily rate upserts without duplicating historical records.
dbt Gold Analytics Marts — Automates rolling 7-day volatility calculations and Month-over-Month (MoM) currency performance with dbt data tests.
Isolation Forest ML Anomaly Radar — Unsupervised outlier detection identifying abnormal FX price spikes using rolling deviation baselines.
Interactive Financial Cockpit — Next.js 14 + Recharts dashboard with a 6×2 Country Currency Selector, Live FX Calculator, glowing area charts, and CSV export.

Technologies & Architecture

Python 3.11Apache Spark (PySpark 3.5)Microsoft SQL Server 2022T-SQLdbt Core (dbt-sqlserver)Apache Airflow 2.10scikit-learn (Isolation Forest)FastAPINext.js 14 (App Router)TypeScriptTailwind CSSRechartsDocker & Docker ComposeGitHub Actions (CI/CD)TerraformAWS S3

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.