Muhammad Mudassir Raza

Muhammad Mudassir Raza

Senior Data Engineer

Designing reliable data platforms, production ETL/ELT pipelines, cloud-native systems, and AI-powered workflows across AWS, GCP, and Azure.

Python SQL ETL / ELT BigQuery AWS Azure GCP LLM / RAG Snowflake AI / LLM Automation
Core capabilities

Engineering that moves data into action.

A practical mix of data engineering, cloud infrastructure, automation, and AI application development.

๐Ÿ›  Skills

Languages & Backend

Python, SQL, Flask, FastAPI, Django

Data Engineering

ETL/ELT pipelines, data warehousing, data modeling, SCD2

Big Data & Streaming

Apache Spark (PySpark, SparkSQL), Kafka (batch & streaming)

Orchestration

Apache Airflow, Cloud Scheduler, Cron

Cloud โ€” AWS

S3, Glue, EC2, Lambda, Athena, IAM, DataBrew, QuickSight

Cloud โ€” GCP

BigQuery, Cloud Functions, Cloud Storage, Vertex AI, Compute Engine, Looker Studio

Cloud โ€” Azure

Azure Functions, Logic Apps

Data Platforms

Snowflake, Databricks, dbt Cloud, PostgreSQL, MySQL, MongoDB

AI / Automation

GPT-4, LLM workflows, n8n automation, REST API integration

DevOps

Docker, Docker Compose, CI/CD, microservices, Git, Linux

Web Scraping

Requests, BeautifulSoup, Scrapy, Selenium

Visualization

Pandas, Matplotlib, Seaborn, Streamlit

๐ŸŽ“ Certifications & Credentials

๐Ÿ… Recognition

๐Ÿงฉ Selected Projects

Architecture-led projects spanning cloud data pipelines, analytics engineering, AI automation, streaming, and data modeling.

๐Ÿง  Mindify

AI-powered mental wellness companion โ€” Backend Developer, Team DataGrains

An empathetic chatbot that takes a user's text or voice input, detects the underlying emotion, retrieves grounded CBT/psychology insights from a vector database, and uses Gemini 2.0 Flash to craft a caring, practical response โ€” complete with coping tips, next steps, and a cited book insight.

StreamlitGemini 2.0 FlashQdrant Sentence-Transformers (MiniLM-L6-v2)SpeechRecognitionPython
User text or voice Streamlit UI app.py chat + audio upload backend.py detect_emotion() Gemini 2.0 Flash retrieve_insights() MiniLM embed โ†’ Qdrant search build prompt + context Qdrant CBT insight vectors Gemini 2.0 Flash generates empathetic reply + CBT tips response rendered in chat

โšก Real-Time Crypto Data Pipeline

End-to-end data engineering project โ€” Kafka, S3, Snowflake, fully dockerized

A real-time pipeline that scrapes live crypto prices, streams them through a Kafka producer/consumer pair, cleans and lands them as JSON in an S3 bucket, then auto-ingests into Snowflake via Snowpipe the moment new files land โ€” no manual loading step. Every service (Zookeeper, Kafka broker, producer, consumer) runs as its own Docker container.

PythonBeautifulSoupApache Kafka Docker ComposeAWS S3Snowflake + SnowpipeEC2 (optional)
crypto.com live price pages docker-compose services producer.py scrapes + sends to Kafka topic Kafka broker + Zookeeper topic: demo_testing2 consumer.py cleans fields, casts types, writes JSON via boto3 AWS S3 real-time JSON files event Snowflake via Snowpipe auto-ingest, no manual load

๐Ÿ“Š Crypto Market Analytics Pipeline

End-to-end analytics engineering โ€” BigQuery star schema, SCD2, LLM commentary

Pulls live coin data from the CoinGecko API into BigQuery, models it into a proper star schema with a SCD Type 2 table that tracks each coin's market-cap rank history over time, then feeds a daily summary view into two places at once: an LLM that writes newsletter-style market commentary, and a Streamlit dashboard that shows both the charts and that commentary side by side.

PythonCoinGecko APIGoogle BigQuery SQL (window functions)SCD Type 2OpenAI GPT-4o-mini Streamlit + PlotlyLooker Studio
CoinGecko API public, no auth crypto_ingest.py requests + service acct BigQuery raw_crypto_prices SQL transformation layer ROW_NUMBER() dedup โ†’ latest snapshot/day Star schema: FCT_PRICES ยท DIM_COIN ยท DIM_DATE LAG/LEAD โ†’ SCD2 rank history (DIM_COIN_HISTORY) VW_DAILY_SUMMARY latest price + rank movement ๐Ÿค– llm_insights.py GPT-4o-mini โ†’ AI_COMMENTARY ๐Ÿ“ˆ Looker Studio trends, gainers/losers ๐Ÿ–ฅ๏ธ Streamlit + Plotly charts + AI commentary in one app

๐ŸŒ€ Data Orchestration with Airflow

Batch pipeline orchestration โ€” Airflow, GCS, BigQuery, Docker

A daily-scheduled Airflow DAG that downloads a public NYC taxi dataset, converts it from CSV to Parquet for efficient storage and querying, uploads it to Google Cloud Storage, then loads it straight into BigQuery โ€” each step a separate, retry-able task chained by dependencies. Airflow itself (webserver, scheduler, Postgres metadata DB) runs entirely through Docker.

Apache AirflowPythonDocker Compose Google Cloud StorageBigQueryPyArrow (Parquet)PostgreSQL
Airflow DAG ยท data_ingestion_gcs_dag ยท runs @daily (Dockerized) download_dataset_task BashOperator curl NYC taxi zone lookup CSV format_to_parquet_task PythonOperator CSV โ†’ Parquet via PyArrow local_to_gcs_task PythonOperator uploads Parquet to GCS bucket (raw/) bigquery_external _table_task GCSToBigQueryOperator loads into BigQuery Each task retries independently ยท dependencies chained with >> operators

๐Ÿ“ˆ Stock Market Real-Time Pipeline (Kafka + Glue/Athena + Snowflake)

Real-time streaming pipeline with dual query layers

Streams historical stock index data row-by-row through Kafka to simulate a live feed, lands it as JSON in S3, then makes it queryable two different ways: an AWS Glue crawler builds a schema catalog so Athena can query S3 directly with no load step, while Snowflake ingests the same files automatically through Snowpipe for warehouse-side analysis โ€” comparing a serverless query path against a managed warehouse path on the same data.

PythonApache KafkaPandas AWS S3AWS Glue CrawlerAmazon AthenaSnowflake + Snowpipe
Stock CSV indexProcessed.csv Producer samples rows โ†’ Kafka Kafka topic demo_testing2 Consumer writes JSON to S3 AWS S3 JSON files Glue Crawler โ†’ Catalog then queried via Athena Snowflake via Snowpipe Same S3 data queried two ways: serverless (Glue+Athena) vs. warehouse (Snowflake)

๐Ÿค– Personal Assistant โ€” n8n AI Agent

No-code AI automation โ€” Gemini agent orchestrating 6 Google/web tools

An AI agent built as an n8n workflow: a webhook receives a request, hands it to a Gemini-powered agent that keeps conversational memory across turns, and the agent decides which of six connected tool groups to call โ€” Gmail, web search, Calendar, Docs, Tasks, or Sheets โ€” before the result is sent back through the same webhook.

n8nGoogle GeminiGmail APIGoogle Calendar Google DocsGoogle TasksGoogle SheetsSerpAPI
Webhook POST trigger AI Agent Gemini Chat Model Simple Memory picks a tool ๐Ÿ“ง Gmail ๐Ÿ” SerpAPI search ๐Ÿ“… Calendar ๐Ÿ“ Google Docs โœ… Google Tasks ๐Ÿ’ฐ Sheets + Calc Respond to Webhook

๐Ÿ’ต Nealytics โ€” Multi-Source Revenue Data Pipeline

Cloud Function orchestrating 4 payment platforms into one warehouse

A Google Cloud Function that pulls revenue data from ThriveCart, Stripe (across multiple accounts, fetched in parallel), Shopify, and PayPal, and loads it all into BigQuery โ€” with incremental loading for ThriveCart (only new rows past the last synced date) and parallel fetching for Stripe to keep runtime down. Every run posts a pass/fail summary to Slack, so failures get flagged immediately instead of going unnoticed.

PythonGoogle Cloud FunctionsBigQuery Stripe APIShopify APIPayPal APIThriveCart API Slack Webhooksjoblib (parallel fetch)
HTTP request {pipelines: "..."} main.py Cloud Function runs requested pipelines ThriveCart incremental load Stripe parallel across accounts Shopify last 60 days PayPal last 60 days BigQuery raw dataset per source (one table each) ๐Ÿ’ฌ Slack pass/fail summary

๐Ÿ›ก๏ธ Insightly CRM โ†’ OneDrive Export Pipeline

Azure-native CRM data export โ€” Functions, Logic Apps, Microsoft Graph

A scheduled Azure Logic App triggers five separate Azure Functions, each responsible for a slice of Insightly CRM data โ€” quotes, organisations, opportunities, equipment/invoices/users, tasks, and opportunity stages. Every function pulls its entity via the Insightly API, exports it to CSV, then authenticates with Microsoft Graph (MSAL client-credentials flow) to upload the file straight into a shared OneDrive/SharePoint folder โ€” splitting the work across functions keeps each run well under Azure's execution time limits.

Azure FunctionsAzure Logic AppsPython Insightly CRM APIMicrosoft Graph APIMSAL (OAuth) OneDrive / SharePointGitHub Actions (CI/CD)
Azure Logic App scheduled trigger 5 HTTP-triggered Azure Functions Trigger1 โ†’ Quote + Organisation Trigger2 โ†’ Opportunity Trigger3 โ†’ Equipment/Invoice/Users Trigger4 โ†’ Task Trigger5 โ†’ Opportunity Stage Insightly CRM fetch entity via API Export to CSV per entity Microsoft Graph MSAL client-credentials token auth โ†’ OneDrive / SharePoint

๐Ÿ“‰ GA4 โ†’ BigQuery Reporting Pipeline

Google Cloud Function โ€” self-healing analytics sync with retry logic

A scheduled Cloud Function that pulls the last 10 days of sessions, events, and key events from a Google Analytics 4 property (by landing page and date), then reconciles it into BigQuery using a delete-then-append pattern โ€” deleting the overlapping date range first so re-runs never create duplicates. Every GA4 call and BigQuery write retries with exponential backoff, and a Slack alert fires the moment a step exhausts its retries.

Google Cloud FunctionsPythonGA4 Reporting API BigQuerySlack WebhooksOAuth2 (service account + refresh token)
Scheduled Cloud Function call fetch_data() GA4 Reporting API last 10 days, retries + backoff delete_table() clears overlapping date range upload_to_bq() WRITE_APPEND, autodetect BigQuery ga4_reporting_api_data ๐Ÿ’ฌ Slack alert on failure

๐Ÿ”Ž Search Console โ†’ BigQuery Pipeline

Google Cloud Function โ€” dual-path GSC sync (by query & by page)

A single Cloud Function that runs two independent sync paths against the same Search Console property: one pulls performance grouped by search query, the other by landing page. The query-level path uses the same delete-then-append pattern as the GA4 pipeline to stay idempotent on re-runs; the page-level path currently appends only. Both land in their own BigQuery table, validated afterwards with clicks/impressions/position roll-up queries.

Google Cloud FunctionsPythonSearch Console API BigQuerySlack Webhooksfunctions-framework
main() HTTP Cloud Function fetch_data_by_query() grouped by search query delete overlap โ†’ append (last 10 days) BigQuery ...search_console_data fetch_data_by_page() grouped by landing page append-only BigQuery ...data_byPage validated after load with clicks/impressions/position roll-up SQL queries

๐Ÿ“ž CallRail โ†’ BigQuery Pipeline

Incremental sync with parallel fetch and ID-level deduplication

Syncs call tracking data from CallRail into BigQuery without ever re-processing what's already there: it reads the latest timestamp already in BigQuery as a watermark, splits the gap up to today into monthly chunks, fetches all of them in parallel (7 threads via joblib), then diffs the fetched IDs against what's already in BigQuery before uploading โ€” so only genuinely new calls, users, and form submissions ever get written.

Google Cloud FunctionsPythonCallRail API BigQueryjoblib (parallel fetch)Pandas
Read watermark MAX(start_time) from BigQuery Monthly chunks watermark โ†’ today Parallel fetch CallRail API calls ยท users ยท forms joblib, 7 threads one chunk per thread Dedupe drop IDs already in BigQuery BigQuery Call_rails dataset ยท new rows only

๐Ÿ“ Google Business Profile โ†’ BigQuery Pipeline

Multi-location metrics sync โ€” daily + monthly, fetched in parallel

Pulls performance metrics for six Google Business Profile store locations at once, running two sub-pipelines in the same function call: daily metrics for the last 8 days, and monthly search-keyword data for the last month. Each pipeline fetches all six locations in parallel (7 workers), then reconciles into BigQuery with the same delete-then-append pattern used across the other Google-data pipelines, so re-runs stay duplicate-free.

Google Cloud FunctionsPythonGoogle Business Profile API BigQueryjoblib (parallel fetch)Pandas
main() 6 GBP locations Michigan Auto Law monthly_metrics() search keywords, last month 7 parallel workers daily_metrics_data() views, calls, last 8 days 7 parallel workers BigQuery monthly_metrics_data BigQuery daily_metrics_data delete overlapping date range, then append (no dupes)

๐Ÿ† Wincher SEO Rank-Tracking Pipeline

Multi-client keyword & competitor tracking โ€” 15+ scheduled Cloud Functions

Tracks keyword rankings, keyword groups, and competitor positions from Wincher for over a dozen client sites (Michigan Auto Law, IHOP, Outback, Qdoba, Ferguson Roofing, and more), each run as its own scheduled Cloud Function built on one shared Wincher client module. Every run covers a rolling 3-month window split into single-day intervals, fetched in parallel (5 threads) with automatic retry on transient API errors, then truncates and reloads that date range in BigQuery so the numbers stay accurate as Wincher's own data gets revised.

Google Cloud FunctionsPythonWincher API BigQueryThreadPoolExecutorSlack Webhooks
Michigan Auto Law IHOP ยท Outback ยท Qdoba Ferguson Roofing + 10 more clients each its own scheduled fn Wincher API client keywords ยท groups ยท competitors Daily intervals rolling 3-month window 5 parallel threads retry on 500 errors BigQuery truncate window โ†’ reload Same shared client module powers every client's pipeline โ€” only the website ID and destination table change per script

๐Ÿ“ฑ Facebook & Instagram Ads Data Pipeline

Meta Marketing API โ†’ BigQuery, multi-endpoint social data sync

Pulls marketing and organic data from Meta's platform through five independent modules โ€” Page insights, post performance, ad performance, ad creative details, and Instagram user insights โ€” each loading into its own BigQuery table. A same-day-run guard checks whether a table already has today's data before pulling again, avoiding wasted API calls when the function gets triggered more than once in a day.

Google Cloud FunctionsPythonMeta Marketing API Instagram Graph APIBigQuery
main() Cloud Function facebook_page_insights() facebook_posts() facebook_ads_creative() facebook_ads() instagram_user_insights() BigQuery one table per endpoint facebook_data_raw dataset should_run_today() checks today's data skips duplicate runs

๐Ÿ›๏ธ CFPB Complaints Pipeline โ€” Airflow + Streamlit

Recursive multi-state API scrape โ†’ Google Sheets โ†’ live dashboard

A daily Airflow DAG recurses through every U.S. state to pull consumer financial complaint data from the CFPB's public API, transforms it, and pushes it into Google Sheets โ€” which a Streamlit dashboard then reads directly to give stakeholders a state-by-state, always-current view of complaint trends without anyone touching a spreadsheet by hand.

Apache AirflowPythonCFPB Public API Google Sheets APIStreamlit
extract_task CFPB API, all 50 states database_task stage raw records transform_task clean + reshape google_sheet task writes results Streamlit live dashboard

๐Ÿ™ Scraping 1M+ GitHub Repositories on GCP

Large-scale API harvesting with multithreading + EDA

Paginates through GitHub's public repository listing and fetches over a million repos, enriching each with follower counts, languages, and stargazers through concurrent worker threads with automatic retry on failed requests. The resulting dataset then feeds an exploratory analysis notebook that surfaces trends across languages, audiences, and user attributes โ€” turning raw API output into a dataset someone can actually draw conclusions from.

PythonGitHub REST APIconcurrent.futures Retry/backoffPandasMatplotlib
GitHub REST API paginated repo list Threaded fetch + retry enrich: followers, languages, stargazers per repo CSV dataset 1M+ rows EDA notebook Pandas + Matplotlib language & audience trends

โญ Northwind OLTP โ†’ OLAP Star Schema (PySpark SCD2)

Data warehouse modeling + Slowly Changing Dimension implementation

Redesigns the classic Northwind OLTP database into a proper star schema: a fact table grained at order/order-detail level with foreign keys to denormalized Employee, Customer, and Product dimensions, plus a Date dimension. MySQL procedures migrate the data from OLTP to OLAP, and PySpark implements a working Slowly Changing Dimension (Type 2) on the Employee dimension โ€” validated against both an insert case and an update case to prove history is tracked correctly.

PySparkMySQLStar Schema Modeling SCD Type 2SQLAlchemy
MySQL OLTP northwind database Migration procedures denormalize dims + build fact grain MySQL OLAP star schema (starflow) Fact: Order + OrderDetail PySpark SCD2 Employee dim, insert + update tested

๐Ÿšข Nested JSON Streaming โ€” Kafka + PySpark

Exam solution โ€” structured streaming, schema flattening, data quality rules

Reads nested Titanic-dataset JSON messages off a Dockerized Kafka topic, applies an explicit schema to work with the nested fields through PySpark's API, then flattens everything to a single level. From there it drops duplicate rows, enforces that key columns are never null, removes fields that aren't needed downstream, fixes numeric types, and writes the cleaned result out as line-delimited JSON.

Apache KafkaPySparkDocker ComposeStructured Schema
titanic.csv reshaped to nested JSON Kafka topic titanic_topic (batch) PySpark processing Apply schema โ†’ flatten nested cols Drop duplicates ยท enforce not-null Drop unused cols ยท fix types (Age, Fare) JSON lines output via DataFrameWriter