Data pipeline chuyển đổi 2 dataset Olist từ Kaggle lên Cloudflare R2 theo kiến trúc Medallion (Bronze → Silver → Gold), tối ưu cho Data Engineering projects.
Link: closed_deals.seller_id = sellers.seller_id → phân tích funnel-to-revenue.
Cài đặt & Chạy
Prerequisites
pip install duckdb boto3 kagglehub
Chạy pipeline
bash
1# Cách 1: Truyền credentials trực tiếp2python olist_r2_pipeline.py \3 --r2-endpoint "https://<account_id>.r2.cloudflarestorage.com"\4 --r2-access-key "<your_access_key>"\5 --r2-secret "<your_secret_key>"67# Cách 2: Qua environment variables8exportR2_ENDPOINT_URL="https://<account_id>.r2.cloudflarestorage.com"9exportR2_ACCESS_KEY_ID="..."10exportR2_SECRET_ACCESS_KEY="..."11python olist_r2_pipeline.py
1213# Test local trước (không upload lên R2)14python olist_r2_pipeline.py --skip-download --skip-upload
Options
Flag
Mô tả
--r2-endpoint
Cloudflare R2 endpoint URL
--r2-access-key
R2 Access Key ID
--r2-secret
R2 Secret Access Key
--skip-download
Bỏ qua download từ KaggleHub
--skip-upload
Chỉ build local, không upload R2
--skip-validate
Bỏ qua bước validate trên R2
Query dữ liệu trên R2 bằng DuckDB
Sau khi upload, query trực tiếp R2 qua DuckDB (zero egress fee):
sql
1-- Setup một lần2CREATE SECRET (TYPE r2, KEY_ID '...', SECRET '...', ACCOUNT_ID '...');34-- Query Silver layer (toàn bộ dataset ~30MB Parquet)5SELECT c.customer_state,COUNT(DISTINCT o.order_id)AS orders,SUM(oi.price)AS revenue
6FROM read_parquet('r2://olist-data-lake/staging/orders.parquet') o
7JOIN read_parquet('r2://olist-data-lake/staging/customers.parquet') c ON o.customer_id = c.customer_id
8JOIN read_parquet('r2://olist-data-lake/staging/order_items.parquet') oi ON o.order_id = oi.order_id
9WHERE o.order_status ='delivered'10GROUPBY1ORDERBY2DESC;1112-- Query Gold với partition pruning (chỉ scan tháng cần)13SELECT*FROM read_parquet(14'r2://olist-data-lake/analytics/orders/**/*.parquet',15 hive_partitioning =true16)WHEREyear='2017'ANDmonth='06';1718-- Funnel analysis (đã join sẵn, có days_to_close)19SELECT origin,COUNT(*)AS leads,AVG(days_to_close)AS avg_days
20FROM read_parquet('r2://olist-data-lake/analytics/funnel/**/*.parquet', hive_partitioning=true)21GROUPBY1;
Các quyết định thiết kế
Quyết định
Giải thích
Parquet thay vì CSV
Columnar + compressed → DuckDB đọc nhanh hơn 10-100x qua network
Snappy compression
Cân bằng tốc độ đọc/ghi + tỉ lệ nén (~3x so với CSV)
Hive partition year/month
Hầu hết query DE filter theo thời gian → partition pruning
geolocation dedup
1M rows → 19K unique zip codes → giảm 50x
TRY_CAST everywhere
Chịu được cả VARCHAR và TIMESTAMP từ read_csv_auto
Funnel pre-joined ở Gold
Có sẵn days_to_close, không cần join lại mỗi lần query
Bảng nhỏ không partition
Sellers (3K), categories (71) giữ single file ở staging
Schema chi tiết
orders (Silver)
Column
Type
order_id
VARCHAR
customer_id
VARCHAR
order_status
VARCHAR
order_purchase_timestamp
TIMESTAMP
order_approved_at
TIMESTAMP
order_delivered_carrier_date
TIMESTAMP
order_delivered_customer_date
TIMESTAMP
order_estimated_delivery_date
TIMESTAMP
funnel mart (Gold) - đã join sẵn
Column
Type
Mô tả
mql_id
VARCHAR
Marketing Qualified Lead ID
seller_id
VARCHAR
FK → sellers
won_date
TIMESTAMP
Ngày chốt deal
business_segment
VARCHAR
Phân khúc
lead_type
VARCHAR
online_medium, online_small...
lead_behaviour_profile
VARCHAR
cat, wolf, eagle, shark
origin
VARCHAR
Nguồn traffic (organic_search, paid_search...)
first_contact_date
TIMESTAMP
Ngày first touch
days_to_close
INTEGER
Số ngày từ MQL → won
Xem đầy đủ schema của tất cả bảng trong source code.