| name | zomato-ai-data-engineering-pipeline |
| description | End-to-end batch data pipeline with Snowflake, dbt, Airflow, and OpenAI for food delivery analytics |
| triggers | ["build a zomato data pipeline","set up snowflake medallion architecture","create dbt incremental models for zomato","orchestrate data pipeline with airflow","enrich reviews with openai llm","implement rag for text data","build text to sql with openai","configure s3 snowflake integration"] |
zomato-ai-data-engineering-pipeline
Skill by ara.so — Data Skills collection.
Complete batch data engineering pipeline that processes food delivery data through a medallion architecture (Bronze → Silver → Gold) using Amazon S3, Snowflake, dbt, Airflow orchestration, and OpenAI-powered AI capabilities (LLM enrichment, RAG, text-to-SQL).
Project Overview
Pipeline Flow:
CSVs → S3 Data Lake → Snowflake RAW (Bronze) → dbt STAGING (Silver) → dbt MARTS (Gold) → AI Layer
Architecture Layers:
- Bronze (RAW): Direct
COPY INTO from S3 via storage integration
- Silver (STAGING): dbt views for cleaning, typing, renaming
- Gold (MARTS): Dimensions, incremental facts (MERGE), business aggregates, SCD2 snapshots
- AI: LLM enrichment, RAG chat, text-to-SQL queries
Data Scale:
- 10M orders
- 23M order items
- 300K text reviews
- 7 source tables (restaurants, users, food, menu, orders, order_items, reviews)
Installation & Setup
Prerequisites
git clone https://github.com/darshilparmar/zomato-ai-data-engineering-end-to-end-project
cd zomato-ai-data-engineering-end-to-end-project
AWS S3 Setup
aws s3 mb s3://your-zomato-bucket
aws s3 sync data/ s3://your-zomato-bucket/raw/ --exclude "*" \
--include "restaurants/*" \
--include "users/*" \
--include "food/*" \
--include "menu/*" \
--include "orders/*" \
--include "order_items/*" \
--include "reviews/*"
aws iam create-policy \
--policy-name zomato-s3-read \
--policy-document file://aws/iam/s3-read-policy.json
aws iam create-role \
--role-name snowflake-s3-role \
--assume-role-policy-document file://aws/iam/snowflake-role-trust-policy-initial.json
aws iam attach-role-policy \
--role-name snowflake-s3-role \
--policy-arn arn:aws:iam::YOUR_ACCOUNT:policy/zomato-s3-read
Snowflake Setup
CREATE WAREHOUSE ZOMATO_WH
WITH WAREHOUSE_SIZE = 'MEDIUM'
AUTO_SUSPEND = 60
AUTO_RESUME = TRUE;
CREATE DATABASE ZOMATO;
USE DATABASE ZOMATO;
CREATE SCHEMA RAW;
CREATE SCHEMA STAGING;
CREATE SCHEMA MARTS;
CREATE SCHEMA SNAPSHOTS;
CREATE SCHEMA AI;
CREATE ROLE DBT_ROLE;
GRANT USAGE ON WAREHOUSE ZOMATO_WH TO ROLE DBT_ROLE;
GRANT ALL ON DATABASE ZOMATO TO ROLE DBT_ROLE;
GRANT ALL ON ALL SCHEMAS IN DATABASE ZOMATO TO ROLE DBT_ROLE;
GRANT ROLE DBT_ROLE TO USER YOUR_USER;
CREATE STORAGE INTEGRATION s3_zomato_integration
TYPE = EXTERNAL_STAGE
STORAGE_PROVIDER = 'S3'
ENABLED = TRUE
STORAGE_AWS_ROLE_ARN = 'arn:aws:iam::YOUR_ACCOUNT:role/snowflake-s3-role'
STORAGE_ALLOWED_LOCATIONS = ('s3://your-zomato-bucket/raw/');
DESC STORAGE INTEGRATION s3_zomato_integration;
STAGE s3_stage
STORAGE_INTEGRATION s3_zomato_integration
URL ;
RAW.RESTAURANTS (
restaurant_id NUMBER,
name ,
city ,
rating ,
rating_count NUMBER,
cost ,
cuisine ,
lic_no ,
link ,
address ,
menu
);
RAW.USERS (
user_id NUMBER,
name ,
email ,
password ,
age NUMBER,
gender ,
marital_status ,
occupation ,
monthly_income NUMBER,
educational_qualifications ,
family_size NUMBER
);
RAW.FOOD (
food_id NUMBER,
item ,
veg_or_non_veg
);
RAW.MENU (
menu_id NUMBER,
restaurant_id NUMBER,
food_id NUMBER,
cuisine ,
price NUMBER
);
RAW.ORDERS (
order_id NUMBER,
user_id NUMBER,
restaurant_id NUMBER,
order_date ,
order_time ,
order_status ,
order_value NUMBER
);
RAW.ORDER_ITEMS (
order_item_id NUMBER,
order_id NUMBER,
food_id NUMBER,
quantity NUMBER,
price NUMBER
);
RAW.REVIEWS (
review_id NUMBER,
order_id NUMBER,
restaurant_id NUMBER,
user_id NUMBER,
rating NUMBER,
review_text ,
review_date
);
dbt Configuration
cd zomato
cat > profiles.yml <<EOF
zomato:
target: dev
outputs:
dev:
type: snowflake
account: "{{ env_var('SNOWFLAKE_ACCOUNT') }}"
user: "{{ env_var('SNOWFLAKE_USER') }}"
password: "{{ env_var('SNOWFLAKE_PASSWORD') }}"
role: DBT_ROLE
database: ZOMATO
warehouse: ZOMATO_WH
schema: STAGING
threads: 4
EOF
export SNOWFLAKE_ACCOUNT=your_account.region
export SNOWFLAKE_USER=your_user
export SNOWFLAKE_PASSWORD=your_password
dbt debug
dbt deps
Airflow Setup
cd airflow
cp example.env .env
docker compose build
docker compose up -d
Key dbt Models
Staging (Silver Layer)
version: 2
sources:
- name: raw
database: ZOMATO
schema: RAW
tables:
- name: restaurants
- name: users
- name: food
- name: menu
- name: orders
- name: order_items
- name: reviews
WITH source AS (
SELECT * FROM {{ source('raw', 'restaurants') }}
),
cleaned AS (
SELECT
restaurant_id,
TRIM(name) AS restaurant_name,
LOWER(TRIM(city)) AS city,
rating,
rating_count,
TRY_CAST(
REPLACE(REPLACE(cost, '₹', ''), ' ', '')
AS NUMBER
) AS avg_cost_for_two,
TRIM(cuisine) AS cuisine,
NULLIF(TRIM(lic_no), '--') AS license_number,
link AS restaurant_url,
address,
menu AS menu_url
FROM source
)
SELECT * FROM cleaned
WITH source AS (
SELECT * FROM {{ source('raw', 'orders') }}
),
cleaned AS (
SELECT
order_id,
user_id,
restaurant_id,
order_date,
order_time,
LOWER(TRIM(order_status)) AS order_status,
order_value,
CASE
WHEN order_status = 'delivered' THEN TRUE
ELSE FALSE
END AS is_delivered,
CASE
WHEN order_status IN ('cancelled', 'canceled') THEN TRUE
ELSE FALSE
END AS is_cancelled
FROM source
)
SELECT * FROM cleaned
Marts (Gold Layer)
{{ config(
materialized='table'
) }}
SELECT
restaurant_id,
restaurant_name,
city,
rating,
rating_count,
avg_cost_for_two,
cuisine,
license_number,
restaurant_url,
address
FROM {{ ref('stg_restaurants') }}
{{ config(
materialized='table'
) }}
WITH customers AS (
SELECT
user_id,
name AS customer_name,
LOWER(email) AS email,
age,
gender,
marital_status,
occupation,
monthly_income,
educational_qualifications,
family_size,
CASE
WHEN age < 25 THEN '18-24'
WHEN age BETWEEN 25 AND 34 THEN '25-34'
WHEN age BETWEEN 35 AND 44 THEN '35-44'
WHEN age BETWEEN 45 AND 54 THEN '45-54'
WHEN age >= 55 THEN '55+'
ELSE 'Unknown'
END AS age_segment
FROM {{ ref('stg_users') }}
)
SELECT * FROM customers
{{ config(
materialized='incremental',
unique_key='order_id',
on_schema_change='append_new_columns'
) }}
WITH orders AS (
SELECT
order_id,
user_id,
restaurant_id,
order_date,
order_time,
order_status,
order_value,
is_delivered,
is_cancelled
FROM {{ ref('stg_orders') }}
{% if is_incremental() %}
WHERE order_date > (SELECT MAX(order_date) FROM {{ this }})
{% endif %}
)
SELECT * FROM orders
{{ config(
materialized='incremental',
unique_key='order_item_id',
on_schema_change='append_new_columns'
) }}
WITH order_items AS (
SELECT
oi.order_item_id,
oi.order_id,
oi.food_id,
oi.quantity,
oi.price,
oi.quantity * oi.price AS line_total,
o.order_date
FROM {{ ref('stg_order_items') }} oi
JOIN {{ ref('stg_orders') }} o ON oi.order_id = o.order_id
{% if is_incremental() %}
WHERE o.order_date > (SELECT MAX(order_date) FROM {{ this }})
{% endif %}
)
SELECT * FROM order_items
{{ config(
materialized='table'
) }}
WITH daily_metrics AS (
SELECT
o.order_date,
r.city,
COUNT(DISTINCT o.order_id) AS total_orders,
COUNT(DISTINCT CASE WHEN o.is_delivered THEN o.order_id END) AS delivered_orders,
COUNT(DISTINCT CASE WHEN o.is_cancelled THEN o.order_id END) AS cancelled_orders,
SUM(CASE WHEN o.is_delivered THEN o.order_value ELSE 0 END) AS gmv,
AVG(CASE WHEN o.is_delivered THEN o.order_value END) AS aov,
COUNT(DISTINCT o.user_id) AS active_customers,
COUNT(DISTINCT o.restaurant_id) AS active_restaurants
FROM {{ ref('fct_orders') }} o
JOIN {{ ref('dim_restaurants') }} r ON o.restaurant_id = r.restaurant_id
o.order_date, r.city
)
order_date,
city,
total_orders,
delivered_orders,
cancelled_orders,
ROUND(cancelled_orders:: (total_orders, ) , ) cancellation_rate_pct,
gmv,
aov,
active_customers,
active_restaurants,
ROUND(gmv (active_restaurants, ), ) revenue_per_restaurant
daily_metrics
dbt Testing
version: 2
models:
- name: dim_restaurants
description: Restaurant dimension
columns:
- name: restaurant_id
description: Primary key
tests:
- unique
- not_null
- name: city
tests:
- not_null
- name: fct_orders
description: Orders fact table (incremental)
columns:
- name: order_id
description: Primary key
tests:
- unique
- not_null
- name: user_id
tests:
- not_null
- relationships:
[, , , ]
Airflow DAG
from datetime import datetime, timedelta
from airflow import DAG
from airflow.providers.snowflake.operators.snowflake import SnowflakeOperator
from airflow.operators.bash import BashOperator
from airflow.operators.python import PythonOperator
import os
SNOWFLAKE_CONN_ID = 'snowflake_default'
S3_BUCKET = os.getenv('S3_BUCKET')
default_args = {
'owner': 'data-eng',
'depends_on_past': False,
'email_on_failure': False,
'email_on_retry': False,
'retries': 1,
'retry_delay': timedelta(minutes=5),
}
with DAG(
'zomato_batch',
default_args=default_args,
description='Zomato end-to-end batch pipeline',
schedule_interval='@daily',
start_date=datetime(2024, 1, 1),
catchup=False,
tags=['zomato', 'batch', 'ai'],
) as dag:
reload_raw = SnowflakeOperator(
task_id='reload_raw',
snowflake_conn_id=SNOWFLAKE_CONN_ID,
sql=f"""
USE SCHEMA ZOMATO.RAW;
COPY INTO RESTAURANTS FROM @s3_stage/restaurants/
FILE_FORMAT = (TYPE = 'CSV' FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1)
FORCE = TRUE;
COPY INTO USERS FROM @s3_stage/users/
FILE_FORMAT = (TYPE = 'CSV' FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1)
FORCE = TRUE;
COPY INTO FOOD FROM @s3_stage/food/
FILE_FORMAT = (TYPE = 'CSV' FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1)
FORCE = TRUE;
COPY INTO MENU FROM @s3_stage/menu/
FILE_FORMAT = (TYPE = 'CSV' FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1)
FORCE = TRUE;
COPY INTO ORDERS FROM @s3_stage/orders/
FILE_FORMAT = (TYPE = 'CSV' FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1)
FORCE = TRUE;
COPY INTO ORDER_ITEMS FROM @s3_stage/order_items/
FILE_FORMAT = (TYPE = 'CSV' FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1)
FORCE = TRUE;
COPY INTO REVIEWS FROM @s3_stage/reviews/
FILE_FORMAT = (TYPE = 'CSV' FIELD_OPTIONALLY_ENCLOSED_BY = '"' SKIP_HEADER = 1)
FORCE = TRUE;
"""
)
dbt_build_core = BashOperator(
task_id=,
bash_command=,
env={
: os.getenv(),
: os.getenv(),
: os.getenv(),
}
)
():
sys
sys.path.append()
enrich_reviews enrich_reviews
enrich_reviews()
enrich_reviews_task = PythonOperator(
task_id=,
python_callable=run_enrich_reviews,
env={
: os.getenv(),
: os.getenv(),
: os.getenv(),
: os.getenv(),
: os.getenv(, ),
}
)
dbt_build_ai = BashOperator(
task_id=,
bash_command=,
env={
: os.getenv(),
: os.getenv(),
: os.getenv(),
}
)
reload_raw >> dbt_build_core >> enrich_reviews_task >> dbt_build_ai
AI Layer
LLM Enrichment
import os
import json
from openai import OpenAI
import snowflake.connector
SNOWFLAKE_ACCOUNT = os.getenv('SNOWFLAKE_ACCOUNT')
SNOWFLAKE_USER = os.getenv('SNOWFLAKE_USER')
SNOWFLAKE_PASSWORD = os.getenv('SNOWFLAKE_PASSWORD')
OPENAI_API_KEY = os.getenv('OPENAI_API_KEY')
SAMPLE_N = int(os.getenv('SAMPLE_N', '1000'))
client = OpenAI(api_key=OPENAI_API_KEY)
def get_snowflake_connection():
return snowflake.connector.connect(
account=SNOWFLAKE_ACCOUNT,
user=SNOWFLAKE_USER,
password=SNOWFLAKE_PASSWORD,
warehouse='ZOMATO_WH',
database='ZOMATO',
schema='RAW',
role='DBT_ROLE'
)
def enrich_review(review_text):
"""Use LLM to extract sentiment and topic from review text."""
prompt = f"""
Analyze this restaurant review and return JSON with:
- sentiment: "positive", "negative", or "neutral"
- topic: main topic like "food_quality", "service", "delivery", "price", "ambiance"
Review: {review_text}
Return only valid JSON, no explanation.
"""
response = client.chat.completions.create(
model='gpt-4o-mini',
messages=[
{'role': 'system', 'content': 'You are a sentiment and topic extraction expert. Return only JSON.'},
{'role': 'user', 'content': prompt}
],
temperature=0
)
result = json.loads(response.choices[].message.content)
result[], result[]
():
conn = get_snowflake_connection()
cursor = conn.cursor()
cursor.execute()
cursor.execute()
reviews = cursor.fetchall()
()
review reviews:
review_id, order_id, restaurant_id, user_id, rating, review_text, review_date = review
:
sentiment, topic = enrich_review(review_text)
cursor.execute(, (review_id, order_id, restaurant_id, user_id, rating, review_text, review_date, sentiment, topic))
()
Exception e:
()
conn.commit()
cursor.close()
conn.close()
()
__name__ == :
enrich_reviews()
RAG Chat
import os
import streamlit as st
from openai import OpenAI
import snowflake.connector
import numpy as np
OPENAI_API_KEY = os.getenv('OPENAI_API_KEY')
client = OpenAI(api_key=OPENAI_API_KEY)
def get_snowflake_connection():
return snowflake.connector.connect(
account=os.getenv('SNOWFLAKE_ACCOUNT'),
user=os.getenv('SNOWFLAKE_USER'),
password=os.getenv('SNOWFLAKE_PASSWORD'),
warehouse='ZOMATO_WH',
database='ZOMATO',
schema='AI',
role='DBT_ROLE'
)
def get_embedding(text):
"""Generate embedding for text."""
response = client.embeddings.create(
model='text-embedding-3-small',
input=text
)
return response.data[0].embedding
def cosine_similarity(a, b):
"""Calculate cosine similarity between two vectors."""
return np.dot(a, b) / (np.linalg.norm(a) * np.linalg.norm(b))
def retrieve_reviews(question, top_k=5):
"""Retrieve most relevant reviews for a question."""
question_embedding = get_embedding(question)
conn = get_snowflake_connection()
cursor = conn.cursor()
cursor.execute("""
SELECT review_id, review_text, sentiment, topic, rating
FROM REVIEW_ENRICHED
""")
reviews = cursor.fetchall()
cursor.close()
conn.close()
scored_reviews = []
review reviews:
review_id, review_text, sentiment, topic, rating = review
review_embedding = get_embedding(review_text)
similarity = cosine_similarity(question_embedding, review_embedding)
scored_reviews.append((similarity, review_id, review_text, sentiment, topic, rating))
scored_reviews.sort(reverse=, key= x: x[])
scored_reviews[:top_k]
():
context = .join([
i, r (context_reviews)
])
prompt =
response = client.chat.completions.create(
model=,
messages=[
{: , : },
{: , : prompt}
]
)
response.choices[].message.content
st.title()
st.write()
question = st.text_input(, placeholder=)
st.button():
question:
st.spinner():
relevant_reviews = retrieve_reviews(question, top_k=)
st.spinner():
answer = generate_answer(question, relevant_reviews)
st.success()
st.write(answer)
st.subheader()
i, (score, review_id, review_text, sentiment, topic, rating) (relevant_reviews):
st.write()
st.write()
st.write()
st.write()
Text-to-SQL
import os
import streamlit as st
from openai import OpenAI
import snowflake.connector
import pandas as pd
import re
OPENAI_API_KEY = os.getenv('OPENAI_API_KEY')
client = OpenAI(api_key=OPENAI_API_KEY)
def get_snowflake_connection():
return snowflake.connector.connect(
account=os.getenv('SNOWFLAKE_ACCOUNT'),
user=os.getenv('SNOWFLAKE_USER'),
password=os.getenv('SNOWFLAKE_PASSWORD'),
warehouse='ZOMATO_WH',
database='ZOMATO',
schema='MARTS',
role='DBT_ROLE'
)
def get_schema_info():
"""Get schema information for marts."""
conn = get_snowflake_connection()
cursor = conn.cursor()
cursor.execute("""
SELECT table_name, column_name, data_type
FROM ZOMATO.INFORMATION_SCHEMA.COLUMNS
WHERE table_schema = 'MARTS'
ORDER BY table_name, ordinal_position
""")
schema = {}