🎬 Recommendation Engine with Cassandra
Cassandra powers real-time recommendations at massive scale - fast reads, writes, and personalization for millions!
Why Cassandra for Recommendations?
- ⚡ Fast Reads: Sub-millisecond personalized recommendations
- 📊 User-Item Matrix: Store billions of user-item interactions
- 🔄 Real-time Updates: Immediate preference learning
- 👥 Scalability: Handle millions of concurrent users
- 💾 Denormalization: Pre-compute recommendations efficiently
- 🌍 Global Distribution: Low-latency recommendations worldwide
Real-World Examples
Production recommendation systems:
• Netflix: 80% of watched content from recommendations
• Spotify: Discover Weekly uses Cassandra for user history
• Amazon: 35% of revenue from recommendations
• YouTube: 70% of watch time from recommendations
🏗️ Recommendation Data Modeling
Typical Recommendation Queries
Q1: Get personalized recommendations for user
Q2: Get similar items (people who liked X also liked Y)
Q3: Get user's interaction history (ratings, views, purchases)
Q4: Find similar users (collaborative filtering)
Q5: Get trending/popular items
Q6: Update user preferences in real-time
Recommendation Approaches
👥
Collaborative Filtering
User-based similarity:
- Find similar users
- Recommend what they liked
- Works without item metadata
- Cold start problem
📊
Content-Based
Item features similarity:
- Analyze item attributes
- Match user preferences
- No cold start for users
- Needs item metadata
🔀
Hybrid
Best of both worlds:
- Combine multiple signals
- More accurate
- Handles cold start
- More complex
Data Modeling Strategy!
Key Principle: Denormalize and pre-compute recommendations
- Store pre-computed recommendations per user
- Update recommendations asynchronously (batch jobs)
- Keep user interaction history for real-time updates
- Cache popular items and trending content
📋 Schema Design
1
User Interactions (Ratings, Views, Clicks)
CREATE TABLE user_interactions (
user_id text,
item_id text,
interaction_type text,
timestamp timestamp,
rating decimal,
duration int,
metadata map<text, text>,
PRIMARY KEY (user_id, timestamp, item_id)
) WITH CLUSTERING ORDER BY (timestamp DESC);
2
Pre-computed Recommendations
CREATE TABLE user_recommendations (
user_id text,
recommendation_type text,
item_id text,
score decimal,
reason text,
generated_at timestamp,
PRIMARY KEY ((user_id, recommendation_type), score, item_id)
) WITH CLUSTERING ORDER BY (score DESC);
3
Item Similarities (Item-to-Item)
CREATE TABLE item_similarities (
item_id text,
similar_item_id text,
similarity_score decimal,
co_occurrence_count int,
PRIMARY KEY (item_id, similarity_score, similar_item_id)
) WITH CLUSTERING ORDER BY (similarity_score DESC);
4
User Similarity (User-to-User)
CREATE TABLE user_similarities (
user_id text,
similar_user_id text,
similarity_score decimal,
common_items set<text>,
PRIMARY KEY (user_id, similarity_score, similar_user_id)
) WITH CLUSTERING ORDER BY (similarity_score DESC);
5
Item Catalog (Metadata)
CREATE TABLE items (
item_id text PRIMARY KEY,
title text,
category text,
tags set<text>,
genres set<text>,
popularity_score decimal,
release_date date,
avg_rating decimal,
total_ratings counter,
metadata map<text, text>
);
6
Trending Items
CREATE TABLE trending_items (
category text,
time_window text,
item_id text,
trend_score decimal,
view_count counter,
PRIMARY KEY ((category, time_window), trend_score, item_id)
) WITH CLUSTERING ORDER BY (trend_score DESC);
👥 Collaborative Filtering
Computing Item Similarities
from cassandra.cluster import Cluster
from collections import defaultdict
import math
cluster = Cluster(['127.0.0.1'])
session = cluster.connect('recommendations')
def compute_item_similarities():
query = "SELECT user_id, item_id, rating FROM user_interactions WHERE rating IS NOT NULL"
interactions = session.execute(query)
user_items = defaultdict(lambda: {})
item_users = defaultdict(set)
for row in interactions:
user_items[row.user_id][row.item_id] = row.rating
item_users[row.item_id].add(row.user_id)
items = list(item_users.keys())
for i, item1 in enumerate(items):
for item2 in items[i+1:]:
common_users = item_users[item1] & item_users[item2]
if len(common_users) >= 5:
dot_product = 0
norm1 = 0
norm2 = 0
for user in common_users:
r1 = user_items[user][item1]
r2 = user_items[user][item2]
dot_product += r1 * r2
norm1 += r1 ** 2
norm2 += r2 ** 2
similarity = dot_product / (math.sqrt(norm1) * math.sqrt(norm2))
insert = """
INSERT INTO item_similarities
(item_id, similar_item_id, similarity_score, co_occurrence_count)
VALUES (?, ?, ?, ?)
"""
session.execute(insert, (item1, item2, similarity, len(common_users)))
session.execute(insert, (item2, item1, similarity, len(common_users)))
print("Item similarities computed!")
compute_item_similarities()
Generate Recommendations from Similar Items
def generate_recommendations_for_user(user_id):
query = """
SELECT item_id, rating FROM user_interactions
WHERE user_id = ? AND rating >= 4.0
LIMIT 50
"""
liked_items = session.execute(query, (user_id,))
recommendations = {}
for interaction in liked_items:
similar_query = """
SELECT similar_item_id, similarity_score
FROM item_similarities
WHERE item_id = ?
LIMIT 10
"""
similar_items = session.execute(similar_query, (interaction.item_id,))
for item in similar_items:
if item.similar_item_id not in recommendations:
score = item.similarity_score * interaction.rating
recommendations[item.similar_item_id] = score
else:
recommendations[item.similar_item_id] += item.similarity_score * interaction.rating
top_recommendations = sorted(
recommendations.items(),
key=lambda x: x[1],
reverse=True
)[:20]
insert = """
INSERT INTO user_recommendations
(user_id, recommendation_type, item_id, score, generated_at)
VALUES (?, ?, ?, ?, ?)
"""
from datetime import datetime
now = datetime.now()
for item_id, score in top_recommendations:
session.execute(insert, (user_id, 'personalized', item_id, score, now))
print(f"Generated {len(top_recommendations)} recommendations for {user_id}")
📊 Content-Based Filtering
Building User Profiles
def build_user_profile(user_id):
query = """
SELECT item_id, rating FROM user_interactions
WHERE user_id = ? AND rating >= 4.0
"""
interactions = session.execute(query, (user_id,))
genre_scores = defaultdict(lambda: 0)
tag_scores = defaultdict(lambda: 0)
for interaction in interactions:
item = session.execute(
"SELECT genres, tags FROM items WHERE item_id = ?",
(interaction.item_id,)
).one()
weight = interaction.rating / 5.0
for genre in item.genres:
genre_scores[genre] += weight
for tag in item.tags:
tag_scores[tag] += weight
return {
'genres': dict(genre_scores),
'tags': dict(tag_scores)
}
def recommend_by_content(user_id):
profile = build_user_profile(user_id)
top_genres = sorted(profile['genres'].items(), key=lambda x: x[1], reverse=True)[:3]
recommendations = []
for genre, score in top_genres:
pass
return recommendations
⚡ Real-time Updates
Track User Interactions
from datetime import datetime
def track_interaction(user_id, item_id, interaction_type, rating=None, duration=None):
insert = """
INSERT INTO user_interactions
(user_id, item_id, interaction_type, timestamp, rating, duration)
VALUES (?, ?, ?, ?, ?, ?)
"""
session.execute(insert, (
user_id, item_id, interaction_type,
datetime.now(), rating, duration
))
if interaction_type == 'view':
session.execute("""
UPDATE items
SET total_ratings = total_ratings + 1
WHERE item_id = ?
""", (item_id,))
if rating and rating >= 4.0:
update_recommendations_realtime(user_id, item_id, rating)
def update_recommendations_realtime(user_id, liked_item_id, rating):
query = """
SELECT similar_item_id, similarity_score
FROM item_similarities
WHERE item_id = ?
LIMIT 5
"""
similar_items = session.execute(query, (liked_item_id,))
insert = """
INSERT INTO user_recommendations
(user_id, recommendation_type, item_id, score, reason, generated_at)
VALUES (?, ?, ?, ?, ?, ?)
"""
for item in similar_items:
score = item.similarity_score * rating
reason = f"Because you rated {liked_item_id} highly"
session.execute(insert, (
user_id, 'realtime', item.similar_item_id,
score, reason, datetime.now()
))
print(f"Updated real-time recommendations for {user_id}")
API Endpoint Example
from flask import Flask, request, jsonify
app = Flask(__name__)
@app.route('/recommendations/<user_id>')
def get_recommendations(user_id):
query = """
SELECT item_id, score, reason
FROM user_recommendations
WHERE user_id = ? AND recommendation_type = 'personalized'
LIMIT 20
"""
results = session.execute(query, (user_id,))
recommendations = []
for row in results:
item = session.execute(
"SELECT title, category FROM items WHERE item_id = ?",
(row.item_id,)
).one()
recommendations.append({
'item_id': row.item_id,
'title': item.title,
'category': item.category,
'score': float(row.score),
'reason': row.reason
})
return jsonify(recommendations)
@app.route('/track', methods=['POST'])
def track():
data = request.json
track_interaction(
data['user_id'],
data['item_id'],
data['type'],
data.get('rating'),
data.get('duration')
)
return jsonify({'status': 'success'})
🔀 Hybrid Recommendation System
Combining Multiple Signals
def generate_hybrid_recommendations(user_id):
recommendations = {}
collab_recs = get_collaborative_recommendations(user_id)
for item_id, score in collab_recs:
recommendations[item_id] = score * 0.4
content_recs = get_content_based_recommendations(user_id)
for item_id, score in content_recs:
if item_id in recommendations:
recommendations[item_id] += score * 0.3
else:
recommendations[item_id] = score * 0.3
trending = get_trending_items()
for item_id, score in trending:
if item_id in recommendations:
recommendations[item_id] += score * 0.2
else:
recommendations[item_id] = score * 0.2
popular = get_popular_items()
for item_id, score in popular:
if item_id in recommendations:
recommendations[item_id] += score * 0.1
else:
recommendations[item_id] = score * 0.1
top_recs = sorted(recommendations.items(), key=lambda x: x[1], reverse=True)[:50]
return top_recs
def get_trending_items():
query = """
SELECT item_id, trend_score
FROM trending_items
WHERE category = 'all' AND time_window = 'daily'
LIMIT 20
"""
results = session.execute(query)
return [(r.item_id, r.trend_score) for r in results]
💡 Best Practices
Recommendation Best Practices
- ✅ Pre-compute: Generate recommendations in batch jobs
- ✅ Denormalize: Store recommendations per user for fast reads
- ✅ Multiple signals: Combine collaborative + content + trending
- ✅ Real-time updates: Quick updates for immediate interactions
- ✅ Diversity: Don't recommend only similar items
- ✅ Explain: Show why items are recommended
Common Mistakes
- ❌ Online computation: Computing recommendations on the fly
- ❌ Cold start: No fallback for new users/items
- ❌ Filter bubble: Only recommending similar content
- ❌ No freshness: Stale recommendations
- ❌ Over-personalization: Missing popular/trending items
- ❌ No A/B testing: Not measuring recommendation quality
Production Architecture
🔄
Batch Processing
Offline computation:
- Spark jobs (daily/weekly)
- Compute similarities
- Generate recommendations
- Store in Cassandra
- High accuracy
⚡
Real-time Layer
Online updates:
- Track interactions instantly
- Quick similarity lookups
- Add recent interactions
- Low latency (<10ms)
- Fresh recommendations
📊
Serving Layer
API responses:
- Read from Cassandra
- Cache popular recs (Redis)
- Personalize ordering
- A/B test variants
- Sub-100ms response
Production Checklist
Before Going Live
- ☑️ Batch job pipeline (Spark/Airflow)
- ☑️ Multiple recommendation tables (personalized, trending, similar)
- ☑️ Real-time tracking implemented
- ☑️ Cold start handling (new users → trending, new items → content-based)
- ☑️ Diversity algorithm to avoid filter bubbles
- ☑️ Explanation/reasoning for recommendations
- ☑️ A/B testing framework
- ☑️ Caching layer (Redis) for hot recommendations
- ☑️ Monitoring: click-through rate, engagement
- ☑️ Fallback to popular items
🎉 Build Your Recommendation Engine!
You're now ready to build Netflix-style personalized recommendations at scale!
🚀 Implementation Roadmap:
- ✅ Design schema for interactions and recommendations
- ✅ Build batch pipeline (Spark) for collaborative filtering
- ✅ Compute item-item and user-user similarities
- ✅ Pre-compute recommendations per user
- ✅ Add real-time tracking layer
- ✅ Implement hybrid approach (multiple signals)
- ✅ Build fast serving API (<100ms)
- ✅ Add caching layer for popular recommendations
- ✅ Set up A/B testing framework
- ✅ Monitor and iterate! 📈
💡 Key Takeaways:
- 📊 Denormalize: Pre-compute and store recommendations
- ⚡ Batch + Real-time: Combine offline accuracy with online freshness
- 🔀 Hybrid signals: Collaborative + content + trending
- 👥 Cold start: Always have fallback recommendations
- 📈 Diversity: Avoid filter bubbles with variety
- 🎯 Explain: Show users why they see each recommendation
🎬 Real-world Examples:
- 🎥 Video streaming (Netflix, YouTube, Disney+)
- 🎵 Music discovery (Spotify, Apple Music)
- 🛍️ E-commerce (Amazon, eBay)
- 📰 News personalization (Google News, Flipboard)
- 📱 Social media feeds (Instagram, TikTok)
- 🎮 Game recommendations (Steam, Epic Games)
🎬 Delight users with personalized recommendations! 🚀