A recommender system that surfaces relevant content is only half the value. The other half is automating the action: sending the right email, triggering the right push notification, or queuing the right ad creative — without a human deciding on each one. This post covers the engineering for a full recommendation-to-execution pipeline: collaborative filtering for personalization, NLP for content matching, a speech-to-text pipeline for processing audio content, and the integration layer that translates model outputs into campaign actions.
The system has three components:
Most real recommendation systems work with implicit feedback: clicks, opens, time spent, purchases — not explicit ratings. The implicit ALS (Alternating Least Squares) algorithm is the standard approach.
import implicit
import numpy as np
import scipy.sparse as sp
from sklearn.preprocessing import normalize
def build_interaction_matrix(events_df, user_col="user_id", item_col="item_id",
weight_col="weight", n_users: int = None, n_items: int = None):
"""
Build a sparse user-item interaction matrix from event data.
events_df: DataFrame with columns [user_id, item_id, weight]
weight: engagement score (e.g., 1 for click, 3 for purchase, 0.5 for open)
"""
# Aggregate multiple interactions between the same user-item pair
agg = events_df.groupby([user_col, item_col])[weight_col].sum().reset_index()
n_users = n_users or agg[user_col].max() + 1
n_items = n_items or agg[item_col].max() + 1
matrix = sp.csr_matrix(
(agg[weight_col].values, (agg[user_col].values, agg[item_col].values)),
shape=(n_users, n_items),
)
return matrix
def train_als_model(interaction_matrix: sp.csr_matrix, factors: int = 64,
iterations: int = 30, regularization: float = 0.01):
"""Train an Alternating Least Squares model."""
model = implicit.als.AlternatingLeastSquares(
factors=factors,
iterations=iterations,
regularization=regularization,
use_gpu=False, # set True if CUDA is available
calculate_training_loss=True,
)
# Transpose: implicit expects item-user matrix (items as rows)
model.fit(interaction_matrix.T)
return model
def get_user_recommendations(model, user_id: int, interaction_matrix: sp.csr_matrix,
n: int = 20, filter_already_seen: bool = True):
"""Get top-N item recommendations for a user."""
item_ids, scores = model.recommend(
user_id,
interaction_matrix[user_id],
N=n,
filter_already_liked=filter_already_seen,
)
return list(zip(item_ids.tolist(), scores.tolist()))
import numpy as np
from sklearn.metrics.pairwise import cosine_similarity
class HybridRecommender:
"""
Combines collaborative filtering (for known users) with
content-based filtering (for new users with no interaction history).
"""
def __init__(self, als_model, item_features: np.ndarray, popular_items: list):
self.als_model = als_model
self.item_features = normalize(item_features, norm="l2") # L2-normalized content embeddings
self.popular_items = popular_items # fallback for unknown users
def recommend(self, user_id: int, interaction_matrix: sp.csr_matrix,
user_content_profile: np.ndarray = None, n: int = 20) -> list:
"""
user_content_profile: embedding of the user's interests (from profile, signup data, etc.)
"""
n_interactions = interaction_matrix[user_id].nnz if user_id < interaction_matrix.shape[0] else 0
if n_interactions >= 5:
# Enough history — use collaborative filtering
return get_user_recommendations(self.als_model, user_id, interaction_matrix, n=n)
elif user_content_profile is not None:
# New user with profile data — content-based filtering
profile = normalize(user_content_profile.reshape(1, -1), norm="l2")
sims = cosine_similarity(profile, self.item_features)[0]
top_indices = sims.argsort()[::-1][:n]
return [(int(idx), float(sims[idx])) for idx in top_indices]
else:
# Truly cold — return popular items
return [(item_id, 1.0) for item_id in self.popular_items[:n]]
Content-based filtering requires a feature representation of each content item. For text content (emails, product descriptions, articles), sentence transformers produce dense embeddings that capture semantic meaning.
from sentence_transformers import SentenceTransformer
import numpy as np
import pandas as pd
def embed_content_items(items_df: pd.DataFrame, text_col: str = "description",
model_name: str = "all-MiniLM-L6-v2") -> np.ndarray:
"""
Compute dense embeddings for content items.
Returns matrix of shape (n_items, embedding_dim).
"""
model = SentenceTransformer(model_name)
texts = items_df[text_col].fillna("").tolist()
# Batch encoding for efficiency
embeddings = model.encode(
texts,
batch_size=64,
show_progress_bar=True,
normalize_embeddings=True, # L2 normalize for cosine similarity
)
return embeddings
def find_similar_items(query_item_id: int, item_embeddings: np.ndarray,
items_df: pd.DataFrame, top_n: int = 10) -> pd.DataFrame:
"""Find items semantically similar to a given item."""
from sklearn.metrics.pairwise import cosine_similarity
query_embedding = item_embeddings[query_item_id].reshape(1, -1)
similarities = cosine_similarity(query_embedding, item_embeddings)[0]
top_indices = similarities.argsort()[::-1][1:top_n + 1] # exclude self
result = items_df.iloc[top_indices].copy()
result["similarity_score"] = similarities[top_indices]
return result.sort_values("similarity_score", ascending=False)
from sklearn.feature_extraction.text import TfidfVectorizer
from sklearn.decomposition import LatentDirichletAllocation
import numpy as np
def extract_topics(texts: list[str], n_topics: int = 20, n_top_words: int = 10):
"""
Extract topics from a corpus of content descriptions using LDA.
Returns: (topic_distributions per document, top words per topic)
"""
vectorizer = TfidfVectorizer(
max_features=5000,
stop_words="english",
min_df=5, # ignore very rare terms
max_df=0.85, # ignore very common terms
)
X = vectorizer.fit_transform(texts)
vocab = vectorizer.get_feature_names_out()
lda = LatentDirichletAllocation(
n_components=n_topics,
random_state=42,
max_iter=20,
learning_method="online",
)
topic_distributions = lda.fit_transform(X)
# Extract top words per topic
topics = []
for topic_idx, topic in enumerate(lda.components_):
top_words = [vocab[i] for i in topic.argsort()[:-n_top_words - 1:-1]]
topics.append({"topic_id": topic_idx, "top_words": top_words})
return topic_distributions, topics
When content includes audio (podcasts, recorded sales calls, video ads), you need to extract text before applying NLP. OpenAI's Whisper provides state-of-the-art transcription.
import whisper
import subprocess
import tempfile
import os
from pathlib import Path
class AudioContentProcessor:
def __init__(self, model_size: str = "base"):
# Model sizes: tiny, base, small, medium, large
# base: fast, good accuracy; large: best accuracy, ~10x slower
self.model = whisper.load_model(model_size)
def transcribe(self, audio_path: str) -> dict:
"""
Transcribe an audio file to text.
Supports mp3, mp4, wav, flac, ogg, and more.
Returns: {text, language, segments (with timestamps)}
"""
result = self.model.transcribe(
audio_path,
fp16=False, # set True if CUDA is available
language=None, # auto-detect language
word_timestamps=True, # per-word timestamps (large model only)
verbose=False,
)
return {
"text": result["text"].strip(),
"language": result["language"],
"duration_seconds": result.get("duration"),
"segments": [
{
"start": seg["start"],
"end": seg["end"],
"text": seg["text"].strip(),
}
for seg in result.get("segments", [])
],
}
def transcribe_from_url(self, url: str) -> dict:
"""Download and transcribe audio from a URL (e.g., S3, CDN)."""
with tempfile.NamedTemporaryFile(suffix=".mp3", delete=False) as f:
tmp_path = f.name
try:
subprocess.run(["wget", "-q", "-O", tmp_path, url], check=True, timeout=120)
return self.transcribe(tmp_path)
finally:
os.unlink(tmp_path)
def process_audio_content_batch(items_df, audio_url_col: str = "audio_url"):
"""Process all audio content items and return enriched DataFrame."""
processor = AudioContentProcessor(model_size="base")
transcriptions = []
for _, row in items_df.iterrows():
if pd.notna(row.get(audio_url_col)):
try:
result = processor.transcribe_from_url(row[audio_url_col])
transcriptions.append({
"item_id": row["item_id"],
"transcript": result["text"],
"language": result["language"],
"duration": result.get("duration_seconds"),
})
except Exception as e:
transcriptions.append({"item_id": row["item_id"], "transcript": "", "error": str(e)})
return pd.DataFrame(transcriptions)
The recommendation engine outputs a ranked list of items per user. The campaign execution layer converts this into concrete campaign actions: which creative to use, when to send, through which channel.
from dataclasses import dataclass
from typing import Optional
import requests
@dataclass
class CampaignAction:
user_id: str
channel: str # email, push, sms
template_id: str # creative template
personalization: dict # template variables (product name, image URL, etc.)
send_at: str # ISO 8601 datetime
priority: int # 1=high, 3=low
class CampaignOrchestrator:
"""
Translates recommendation outputs into campaign execution actions.
Integrates with a marketing automation platform (e.g., Braze, Iterable, HubSpot).
"""
def __init__(self, api_url: str, api_key: str, items_catalog: pd.DataFrame):
self.api_url = api_url
self.headers = {"Authorization": f"Bearer {api_key}", "Content-Type": "application/json"}
self.catalog = items_catalog.set_index("item_id")
def recommendations_to_actions(
self,
recommendations: dict[str, list[tuple[int, float]]],
channel: str = "email",
send_at: str = "now",
max_items_per_campaign: int = 3,
) -> list[CampaignAction]:
"""
recommendations: {user_id: [(item_id, score), ...]} from the recommender
Returns list of CampaignAction objects ready to execute.
"""
actions = []
for user_id, ranked_items in recommendations.items():
top_items = ranked_items[:max_items_per_campaign]
if not top_items:
continue
# Build personalization dict from top recommended items
personalization = {
"primary_item_name": self.catalog.loc[top_items[0][0], "name"],
"primary_item_image": self.catalog.loc[top_items[0][0], "image_url"],
"primary_item_url": self.catalog.loc[top_items[0][0], "product_url"],
"secondary_items": [
{
"name": self.catalog.loc[iid, "name"],
"url": self.catalog.loc[iid, "product_url"],
}
for iid, _ in top_items[1:]
],
}
# Select template based on channel and number of items
template_id = self._select_template(channel, len(top_items), top_items[0][1])
actions.append(CampaignAction(
user_id=user_id,
channel=channel,
template_id=template_id,
personalization=personalization,
send_at=send_at,
priority=1 if top_items[0][1] > 0.8 else 2,
))
return actions
def _select_template(self, channel: str, n_items: int, top_score: float) -> str:
"""Select campaign template based on context."""
if channel == "email":
return "email-multi-product" if n_items > 1 else "email-single-product"
elif channel == "push":
return "push-urgent" if top_score > 0.9 else "push-standard"
return "generic-template"
def execute_batch(self, actions: list[CampaignAction], batch_size: int = 100) -> dict:
"""Send campaign actions to the marketing platform API."""
results = {"sent": 0, "failed": 0, "errors": []}
for i in range(0, len(actions), batch_size):
batch = actions[i:i + batch_size]
payload = {
"campaigns": [
{
"recipient_id": a.user_id,
"channel": a.channel,
"template_id": a.template_id,
"send_at": a.send_at,
"personalization": a.personalization,
}
for a in batch
]
}
try:
resp = requests.post(
f"{self.api_url}/campaigns/batch",
json=payload,
headers=self.headers,
timeout=30,
)
resp.raise_for_status()
r = resp.json()
results["sent"] += r.get("queued", 0)
results["failed"] += r.get("failed", 0)
except Exception as e:
results["errors"].append(str(e))
results["failed"] += len(batch)
print(f"Campaign execution: {results['sent']} queued, {results['failed']} failed")
return results
import hashlib
def get_experiment_arm(user_id: str, experiment_name: str, arms: list[str]) -> str:
"""
Deterministic assignment: same user always lands in the same arm.
Uses hashing for reproducibility without a database lookup.
"""
key = f"{experiment_name}:{user_id}"
bucket = int(hashlib.md5(key.encode()).hexdigest(), 16) % 100
thresholds = [i * (100 // len(arms)) for i in range(1, len(arms))]
for i, t in enumerate(thresholds):
if bucket < t:
return arms[i]
return arms[-1]
# Usage: 50/50 split between collaborative filtering and popularity baseline
arm = get_experiment_arm(user_id="user_12345", experiment_name="rec-v2-launch", arms=["treatment", "control"])
if arm == "treatment":
recs = hybrid_recommender.recommend(user_id, interaction_matrix, n=20)
else:
recs = [(item_id, 1.0) for item_id in popular_items[:20]]
A recommendation system that drives campaign automation requires three components working together: a collaborative filtering model (ALS for implicit feedback is the standard choice), a content understanding pipeline (sentence transformer embeddings + optional LDA topic modeling), and a campaign execution layer that translates ranked item lists into concrete, personalized campaign actions. The speech-to-text component adds value when content includes audio — it feeds transcript text back into the NLP pipeline to produce content embeddings. The full pipeline, from user interaction data to triggered campaign delivery, can run end-to-end in under an hour for millions of users when properly batched and parallelized.