MCP Data Integration Server
by RJrohan47
README.md
# Multi-Source Lakehouse Pipeline
An end-to-end, orchestrated data pipeline that uses **FastMCP** to ingest mocked PostgreSQL, CSV, Excel, and JSON sales data into an **AWS S3** landing zone, standardizes it into a common model as **Parquet**, builds it up through a **Databricks (Unity Catalog + Delta Lake)** lakehouse into business-ready Gold tables, and serves those tables two ways: a 5-page **Databricks AI/BI dashboard** and a **Genie Space** for natural-language Q&A. A **Databricks Job** chains the full lakehouse build and dashboard refresh into a single, repeatable, on-demand or scheduled run — from raw files to a query-ready agent, following each step automatically.
> [!NOTE]
> All files in `sample_data/` are mocked demo data and are intentionally kept in the repository. Credentials are read only from a local `.env` file, which is ignored by Git.
## Contents
- [Architecture](#architecture)
- [Features](#features)
- [Repository structure](#repository-structure)
- [Quick start](#quick-start)
- [Configuration](#configuration)
- [Run the pipeline](#run-the-pipeline)
- [MCP tools](#mcp-tools)
- [Data layers and outputs](#data-layers-and-outputs)
- [AI/BI dashboard](#aibi-dashboard)
- [Orchestration](#orchestration)
- [Genie Space](#genie-space)
- [Troubleshooting](#troubleshooting)
- [Security](#security)
- [Roadmap](#roadmap)
## Architecture
```text
BRONZE — RAW INGESTION
PostgreSQL table ──┐
CSV file ──────────┼──> FastMCP tools ───> Amazon S3 landing/
Excel workbook ────┤ ├── sql/
JSON API payload ──┘ ├── csv/
├── excel/
└── json/
│
▼
SILVER — STANDARDIZATION AND PARQUET
Python transformers normalize, combine,
and deduplicate mocked source data
│
▼
Amazon S3 silver/
├── fact_sales_transactions/
├── dim_customers/
└── dim_products/
│
▼
┌─────────────────────────────────────────────────────────────────┐
│ DATABRICKS JOB: MCP_Data_Pipeline_Run (orchestration) │
│ │
│ Setup_CatalogSchema → Load_StagingData → Clean_StagingData → │
│ Creating_GoldLayer → Customer_Sales_Report │
│ │
│ DATABRICKS LAKEHOUSE — UNITY CATALOG (mcp_lakehouse) │
│ │
│ silver_staging (external Delta tables over S3 silver/ Parquet) │
│ │ │
│ ▼ │
│ silver_clean (managed Delta tables — deduped, │
│ surrogate keys, standardized types & currency) │
│ │ │
│ ▼ │
│ gold (business-ready Delta marts) │
│ ├── customer_segmentation │
│ ├── product_performance │
│ ├── regional_sales │
│ ├── regional_product_performance │
│ └── monthly_sales │
└─────────────────────────┬─────────────────────────┬─────────────┘
│ │
▼ ▼
DATABRICKS AI/BI DASHBOARD GENIE SPACE (natural-
(Customer Sales Report) language Q&A over
Executive Summary · Regional Analysis the Gold tables)
Product Performance · Customer Analytics
· Global Filters
```
## Features
- **Multi-source ingestion:** PostgreSQL, CSV, Excel, and JSON inputs.
- **MCP tool interface:** FastMCP exposes ingestion and transformation routines as callable tools.
- **Schema normalization:** Source-specific column names are mapped into common sales, customer, and product entities.
- **In-memory S3 processing:** Landing files are read from S3 and Parquet is written through memory buffers.
- **Lakehouse layer on Databricks:** Unity Catalog (`mcp_lakehouse`) with `silver_staging`, `silver_clean`, and `gold` schemas built on Delta Lake.
- **Business-ready Gold marts:** PySpark aggregations for customer segmentation (RFM-style recency/frequency/monetary with Active/Risky/Inactive status), product performance, and regional/monthly sales.
- **Databricks AI/BI dashboard:** A 5-page Lakeview dashboard (`Customer_Sales_Report.lvdash.json`) built directly on the Gold tables, with KPI counters, trend and regional charts, and an interactively sortable top-customers table.
- **End-to-end orchestration:** A Databricks Job (`MCP_Data_Pipeline_Run`) chains catalog setup, staging, cleaning, Gold-layer build, and dashboard refresh into a single on-demand or scheduled run, with Unity Catalog lineage tracked automatically.
- **Genie Space:** A natural-language analytics agent over the Gold tables, letting non-technical users ask sales, customer, and regional questions directly without writing SQL.
- **Mocked data included:** The repository can be inspected without sharing real business data.
## Repository structure
```text
.
├── sample_data/ # Mocked CSV, Excel, and JSON source files
├── src/
│ ├── config.py # Reads AWS and PostgreSQL variables from .env
│ ├── database.py # PostgreSQL table-to-CSV extraction
│ ├── storage.py # Amazon S3 client and upload helper
│ ├── parquet_utils.py # S3 reads and Snappy Parquet writes
│ ├── server.py # FastMCP server and public tools
│ └── transformers/
│ ├── sql_source.py # PostgreSQL landing-data transformer
│ ├── spreadsheet_source.py # CSV and Excel transformer
│ ├── json_source.py # JSON order transformer
│ └── base.py # Shared name-cleaning helper
├── DDL/
│ ├── all_schemas_access.sql # Catalog + schema bootstrap (silver_staging, silver_clean, gold)
│ ├── silver_staging/ # External Delta table definitions over S3 silver/ Parquet
│ ├── silver_clean/ # Managed Delta table definitions (cleaned, keyed)
│ └── gold/ # Business mart table definitions
├── Notebooks/
│ ├── 1_Creating mcp pipeline catalog and schema.ipynb # Unity Catalog + external table setup
│ ├── 2_Loading and viewing staging silver schema data.ipynb # Sanity-check silver_staging
│ ├── 3_Cleaning and Transforming silver staging and loading in silver clean.ipynb
│ └── 4_Business Ready Gold Layer.ipynb # Builds all 5 Gold tables
├── Customer_Analytical_Report/
│ └── Customer_Sales_Report.lvdash.json # Databricks AI/BI (Lakeview) dashboard definition
├── assets/
│ └── screenshots/ # README screenshots (orchestration Job, Genie Space)
├── test_connections.py # Optional AWS and PostgreSQL connectivity test
├── requirements.txt
├── .gitignore
└── README.md
```
## Quick start
### Prerequisites
- Python 3.10 or later
- An AWS account, S3 bucket, and IAM credentials with bucket access
- A reachable PostgreSQL instance containing the configured source table
- A Databricks workspace with Unity Catalog enabled, and access to the same S3 bucket (for the lakehouse layer, dashboard, orchestration Job, and Genie Space)
### Install
<details>
<summary>Windows PowerShell</summary>
```powershell
python -m venv .venv
.\.venv\Scripts\Activate.ps1
python -m pip install --upgrade pip
python -m pip install -r requirements.txt
```
</details>
<details>
<summary>macOS / Linux</summary>
```bash
python3 -m venv .venv
source .venv/bin/activate
python -m pip install --upgrade pip
python -m pip install -r requirements.txt
```
</details>
## Configuration
Create a `.env` file in the project root. Never commit this file.
```env
# AWS / S3
AWS_ACCESS_KEY_ID=your_access_key_id
AWS_SECRET_ACCESS_KEY=your_secret_access_key
AWS_REGION=eu-north-1
S3_BUCKET_NAME=your_s3_bucket_name
# PostgreSQL
POSTGRES_HOST=localhost
POSTGRES_PORT=5433
POSTGRES_DB=etl_pipeline
POSTGRES_USER=postgres
POSTGRES_PASSWORD=your_postgres_password
```
> [!TIP]
> `.env` and `.env.*` are ignored by Git. You can safely keep local credentials there; use `.env.example` for a commit-safe template if you add one later.
>
> The S3 location referenced by `DDL/silver_staging/*.sql` must point at the same bucket configured in `S3_BUCKET_NAME` — update the `LOCATION` clause in those files if you use a different bucket name.
### Verify connectivity
```powershell
python test_connections.py
```
This tests the configured PostgreSQL connection and lists objects in the configured S3 bucket.
## Run the pipeline
### 1. Stage every source in the Bronze layer
```powershell
python -c "from src.server import batch_ingest_all_sources; print(batch_ingest_all_sources())"
```
This stages the mocked local CSV, Excel, and JSON files plus the configured PostgreSQL table under `landing/` in S3.
### 2. Build the Silver layer (S3 Parquet)
```powershell
python -c "from src.server import build_complete_silver_layer; print(build_complete_silver_layer())"
```
This reads the landing files, normalizes them, removes duplicates, and writes three Parquet datasets to `silver/` in S3.
### 3. Start the FastMCP server
```powershell
python src/server.py
```
An MCP-compatible client can then discover and call the project tools.
### 4. Build the Databricks lakehouse layer
Run the following notebooks, in order, against a Databricks cluster or SQL warehouse with Unity Catalog access:
| Notebook | Purpose |
| --- | --- |
| `1_Creating mcp pipeline catalog and schema.ipynb` | Creates the `mcp_lakehouse` catalog, the `silver_staging` / `silver_clean` / `gold` schemas, and external tables over the S3 `silver/` Parquet output. |
| `2_Loading and viewing staging silver schema data.ipynb` | Sanity-checks that `silver_staging` tables read correctly from S3. |
| `3_Cleaning and Transforming silver staging and loading in silver clean.ipynb` | Cleans, dedupes, and adds surrogate keys, writing managed Delta tables to `silver_clean`. |
| `4_Business Ready Gold Layer.ipynb` | Builds the 5 Gold Delta tables from `silver_clean`. |
The DDL under `DDL/` mirrors what these notebooks create and can be run directly (e.g. `DDL/all_schemas_access.sql`) to bootstrap or inspect the catalog independently of the notebooks.
### 5. Import the AI/BI dashboard
In a Databricks workspace with access to the `mcp_lakehouse.gold` schema, import `Customer_Analytical_Report/Customer_Sales_Report.lvdash.json` as a new Lakeview dashboard (Dashboards → Create → Import). See [AI/BI dashboard](#aibi-dashboard) below for what it contains.
### 6. Automate steps 4–5 with the orchestration Job
Steps 4 and 5 don't have to be run by hand every time — the `MCP_Data_Pipeline_Run` Databricks Job chains catalog setup, staging, cleaning, the Gold build, and the dashboard refresh into one on-demand or scheduled run. See [Orchestration](#orchestration) below.
### 7. Ask questions in the Genie Space
Once the Gold tables exist, the `MCP Data Pipeline Sales Analytics Agent` Genie Space can answer sales, customer, and regional questions directly in natural language. See [Genie Space](#genie-space) below.
## MCP tools
| Tool | Purpose | Default behavior |
| --- | --- | --- |
| `ingest_postgres_table_to_s3` | Extracts a PostgreSQL table as CSV and uploads it to S3. | Reads `store.customer_transactions`. |
| `ingest_file_to_s3` | Uploads one local raw file to a category in the landing layer. | Caller supplies a path and category. |
| `batch_ingest_all_sources` | Runs the PostgreSQL, CSV, Excel, and JSON Bronze ingestion steps. | Uses the files in `sample_data/`. |
| `build_complete_silver_layer` | Builds and uploads the normalized Silver Parquet datasets. | Writes sales, customer, and product tables. |
## Data layers and outputs
### Bronze: raw S3 landing files
```text
landing/sql/customer_transactions.csv
landing/csv/source_regional_sales.csv
landing/excel/raw_regional_sales.xlsx
landing/json/source_api_orders.json
```
### Silver: standardized Parquet datasets (S3)
```text
silver/fact_sales_transactions/data.parquet
silver/dim_customers/data.parquet
silver/dim_products/data.parquet
```
| Dataset | Description | Key fields |
| --- | --- | --- |
| `fact_sales_transactions` | Individual sales transactions from every source. | `transaction_id`, `customer_id`, `product_id`, `transaction_date`, `amount`, `quantity`, `sales_region` |
| `dim_customers` | Deduplicated customer records. | `customer_id`, `full_name`, `email`, `customer_region` |
| `dim_products` | Deduplicated product and category records. | `product_id`, `product_name`, `category` |
### silver_staging: external Delta tables (Databricks, `mcp_lakehouse.silver_staging`)
Read directly off the S3 `silver/` Parquet output — same shape as above, with the addition of `data_source` for lineage back to the originating system.
### silver_clean: cleaned, keyed Delta tables (Databricks, `mcp_lakehouse.silver_clean`)
| Table | Key fields | Notes |
| --- | --- | --- |
| `dim_customers` | `customer_key`, `customer_id`, `full_name`, `email` | Surrogate key added; missing emails backfilled from name; `data_source` dropped. |
| `dim_products` | `product_key`, `product_type`, `category` | `product_id` prefix (`PROD_`) stripped into `product_type`. |
| `fact_sales_transactions` | `transaction_key`, `transaction_id`, `customer_id`, `product_type`, `transaction_date`, `Year`, `Month`, `quantity`, `amount_in_USD`, `sales_region` | Invalid rows (non-positive amount/quantity, missing customer ID) filtered out; only `postgresql_legacy` amounts are converted from INR to USD, other sources pass through unchanged; `Year`/`Month` derived from `transaction_date`. |
### gold: business-ready Delta marts (Databricks, `mcp_lakehouse.gold`)
| Table | Description | Key fields |
| --- | --- | --- |
| `customer_segmentation` | Per-customer recency/frequency/monetary summary with a lifecycle status. | `customer_id`, `days_since_last_purchase`, `purchase_frequency`, `total_monetary_value`, `customer_status` (`Active` ≤30 days, `Risky` ≤100 days, else `Inactive`) |
| `product_performance` | Revenue and volume by product type. | `product_type`, `total_customers`, `total_transactions`, `total_quantity_sold`, `total_revenue`, `avg_transaction_value` |
| `regional_sales` | Revenue and volume by sales region. | `sales_region`, `Total_Customers`, `Total_Transactions`, `Total_Sales`, `Total_Quantity` |
| `regional_product_performance` | Revenue and volume by region × product type. | `sales_region`, `product_type`, `transactions`, `quantity_sold`, `revenue` |
| `monthly_sales` | Revenue and volume by calendar month. | `MonthYear`, `Total_Customers`, `Total_Transactions`, `Total_Sales`, `Total_Quantity` |
<details>
<summary>How data is standardized</summary>
- PostgreSQL column variants such as `id` / `transaction_id` and `total` / `amount` are handled by the SQL transformer.
- CSV and Excel headers are lowercased, trimmed, and converted to underscore-separated names before their records are combined.
- JSON orders are flattened with `pandas.json_normalize()`.
- Customer names are cleaned into first name, last name, and full name fields.
- Sales are deduplicated by `transaction_id` plus `data_source`; customers and products are deduplicated by their IDs.
- In `silver_clean`, only `postgresql_legacy` amounts are converted from INR to USD — other sources are already in USD and pass through unchanged. Surrogate keys (`*_key`) are added, and invalid rows (non-positive amount/quantity, missing customer ID) are dropped before the Gold layer is built.
</details>
## AI/BI dashboard
`Customer_Analytical_Report/Customer_Sales_Report.lvdash.json` is a Databricks AI/BI (Lakeview) dashboard built directly on the 5 Gold tables. It ships as 5 pages:
| Page | Contents |
| --- | --- |
| **Executive Summary** | KPI counters (total revenue, customers, transactions, quantity), a monthly revenue trend chart, and a regional sales bar chart. |
| **Regional Analysis** | Regional quantity trend, revenue/average-revenue counters, a region-by-metric heatmap, and average-transactions counter. |
| **Product Performance** | Top-5-products table, a product trend chart, revenue/transaction counters, and a product-revenue-by-customer-type bar chart. |
| **Customer Analytics** | Customer status breakdown (Active/Risky/Inactive) as a pie chart, purchase-frequency distribution, customer/value/frequency counters, and a top-10-customers table. |
| **Global Filters** | Shared filter definitions applied across pages. |
The Customer Analytics page's top-customers table is parameter-driven rather than statically sorted: `sort_by` (Monetary Value / Purchase Frequency) and `sort_order` (Highest to Lowest / Lowest to Highest) filter widgets are bound to SQL `CASE` expressions in the underlying query, so the ranking updates live as a viewer changes the filters — no static "top N" assumption baked into the query.
## Orchestration

`MCP_Data_Pipeline_Run` is a Databricks Job that chains the entire lakehouse build into a single run, in dependency order:
```text
Setup_CatalogSchema → Load_StagingData → Clean_StagingData → Creating_GoldLayer → Customer_Sales_Report
```
| Task | What it does |
| --- | --- |
| `Setup_CatalogSchema` | Runs Notebook 1 — creates the `mcp_lakehouse` catalog and schemas. |
| `Load_StagingData` | Runs Notebook 2 — verifies `silver_staging` reads correctly from S3. |
| `Clean_StagingData` | Runs Notebook 3 — cleans and loads `silver_clean`. |
| `Creating_GoldLayer` | Runs Notebook 4 — builds all 5 Gold tables. |
| `Customer_Sales_Report` | Refreshes the AI/BI dashboard so it reflects the latest Gold data. |
Unity Catalog tracks lineage across the whole run automatically — the job currently reports **13 upstream tables and 11 downstream tables** — and Photon performance optimization is enabled on the job cluster. The job can be triggered on demand (**Run now**) or put on a recurring schedule from the same Jobs & Pipelines UI.
## Genie Space

**MCP Data Pipeline Sales Analytics Agent** is a Databricks Genie Space built directly on the Gold tables, letting anyone — not just SQL users — ask sales, customer, and regional questions in plain English and get answered from governed, business-ready data. Capabilities include:
- Identifying high-value customers via RFM segmentation, lifetime value, and activity status.
- Tracking monthly sales trends, revenue, customer counts, and product quantities for forecasting and year-over-year comparisons.
- Analyzing product-level performance (customer reach, units sold, revenue) to inform portfolio decisions.
- Comparing product performance and sales across regions for distribution and marketing planning.
Example questions it answers directly: *"Give me Top Customers data who has higher monetary value based on their status,"* *"In which region there was maximum number of transactions?"*, and *"What is the monthly trend of total sales revenue?"*
Because it sits on top of the same Gold tables the dashboard and orchestration Job maintain, the Genie Space always reflects the latest successful pipeline run — no separate data prep required.
## Troubleshooting
| Symptom | Likely cause | Resolution |
| --- | --- | --- |
| `NoCredentialsError` or S3 authentication failure | AWS values are missing or invalid. | Check the `.env` values and the IAM policy for the bucket. |
| Connection refused on PostgreSQL | PostgreSQL is unavailable or using another port. | Check `POSTGRES_HOST`, `POSTGRES_PORT`, and that the database service is running. |
| Missing landing object during Silver build | Bronze ingestion has not run successfully. | Run `batch_ingest_all_sources()` before building Silver. |
| Parquet engine error | `pyarrow` was not installed in the active environment. | Run `python -m pip install -r requirements.txt`. |
| `TABLE_OR_VIEW_NOT_FOUND` on `silver_staging` in Databricks | External table created before S3 `silver/` data existed, or bucket/path mismatch. | Run the Python Silver step first, confirm the `LOCATION` in `DDL/silver_staging/*.sql` matches your bucket, then re-run Notebook 1. |
| `CATALOG_NOT_FOUND` / `SCHEMA_NOT_FOUND` in Databricks | Catalog/schemas not yet created in this workspace. | Run `DDL/all_schemas_access.sql` or Notebook 1 before the cleaning/Gold notebooks. |
| Dashboard imports with broken/empty visuals | Gold tables haven't been built yet, or the importing workspace can't reach `mcp_lakehouse.gold`. | Run Notebook 4 first, and confirm Unity Catalog permissions on the `gold` schema for the importing user. |
| Orchestration Job task fails partway through | An upstream task (e.g. `Clean_StagingData`) errored, so downstream tasks didn't run. | Check the failed task's run logs in Jobs & Pipelines; task dependencies mean a fix only requires re-running from the failed step, not the whole job. |
| Genie Space gives an outdated or empty answer | The Gold tables haven't been refreshed since the last data change. | Trigger `MCP_Data_Pipeline_Run` (or wait for its schedule) before querying the Genie Space again. |
## Security
- `.env`, `.env.*`, private keys, certificates, and common credential-file names are excluded in `.gitignore`.
- Do not paste access keys or passwords into source code, issues, commits, or README examples.
- The included `sample_data/` files are mocked; replace them with governed data sources for production use.
- For production, use short-lived IAM roles or a managed secret store instead of long-lived local access keys where possible.
This server cannot be deployed
Maintenance
ActivityMaintained
ResponsivenessNo issues