← Coursework

Movie Recommendation System

A movie recommender run as a production service: viewing events arrive as a stream, are checked and stored once, feed a loop that trains and registers candidate models, and the live traffic itself decides which one serves.

Course
Machine Learning in Production
School
Carnegie Mellon University
Term
Spring 2025
Role
ML Engineer
Stack
  • Python
  • Apache Kafka
  • PostgreSQL
  • Surprise
  • MLflow
  • Flask
  • Prometheus
  • Grafana
  • Docker Compose
  • Jenkins
  • pytest
Code
Offline evaluationTime ordered splitSorted by last watch timeoldest 80% trains the modelsnewest 20% is held out to testsrc/model_training.py · no shuffleCandidate modelsTwo Surprise SVD configurationsv1: 15 epochs, lr 0.0005v2: 25 epochs, lr 0.001svd_model_v1, v2 · reg 0.7 and 0.9train and test setsFeedback loopOnline · A/B testtest RMSEregistered v1, v2fresh snapshotRetrainReruns ingest, then trainingon the fresh snapshotregisters new versionsab_test/periodic_retrain.pyOffline scoresTest RMSE for each candidatePrecision@20 on the held out20%, relevant at 3.5 or moreoffline_eval.pyA/B splitBoth candidates live at onceeven user ids see v1odd user ids see v2recommender_ab_test.py · user_id % 2new eventsserved listsWhat people watch nextServed lists shape laterwatches and ratings, whichland back on the streamKafka · topic movielog6Feedback checkDivisive titles recommended9.15 times on average, vs 6.00Top tenth: 57.31% of ratingsfeedbackloop.ipynb · p 0.0151Online score per armHit: a served movie watchedlater, or rated 2.5 or moreone accuracy gauge per armonline_eval.py · online_eval_ab_test.pyPromote20 movie ids per useraccuracy_a, accuracy_bRedeployA script picks the winnerto go in the service imageand rebuilds the stackcompare_and_deploy.pyCompareReads both scores after 100 sof a 180 s live windowhigher wins, a tie goes to v1pipeline_ab_test.sh · Prometheus querywinner
The loop that decides which model serves. Candidates are scored offline, then live side by side, and what people watch next becomes the data the next round learns from.

Problem

Any service with a large catalogue and many users, from video streaming to retail and news, has to pick a short list for each person the moment they ask, and learn from what they do next. Error on held out ratings says little about whether people watch what they were shown, so the real test happens on live traffic, and it never ends. Tastes drift, new users arrive with no history, and every list served changes the data the next model learns from. This project builds that loop for a simulated streaming service with tens of thousands of users, from the event stream to a model chosen by an online experiment.

Solution

An event stream carries every watch, rating and served list, and an ingest job checks each event and keeps one durable store of ratings and viewing sessions, the dataset every model trains on. Training splits that history in time order, fits candidate models, and registers each one with its parameters, its error, its training data and its code version, so any model that serves can be traced back to the run that made it. A stateless service loads the latest registered version and answers with twenty movie ids per user, falling back to the movies the model scores highest overall for people it has never seen. A separate evaluator reads the same stream to score what was served against what people watched next, and an A/B split by user puts two candidates live side by side and compares their scores to pick the one to redeploy.

IngestViewer event streamWatches, ratings, served listsfrom a simulated streamingservice, over an SSH tunnelKafka · topic movielog6 · :9092Validating ingest jobBatches of 700,000, up to 24Rejects bad ratings and minutesStops at a day boundaryrating 1 to 5 · minute 0 to 5220watch, rate eventsServeWatchStoreserved list, loggedserved lists, later eventschecked rowsRecommendation service20 movie ids per requestLoads the latest versionNew users: overall top moviesFlask · GET /recommend/<user_id> · :8082Online evaluatorServed lists joined to laterevents; a hit is a watch, ora rating of 2.5 or moreonline_eval.py · exporter :8765Event storeRatings and viewing sessionsWatches within 120 min mergeDrift checks over the tablePostgreSQL · ratings, watch_timeaccuracy, status, latencylatest rating per pairMetrics and dashboardsAccuracy, status and latencyScraped every 15 sPrometheus :9090 · Grafana :3000Training snapshotLatest rating per user, movieplus the viewing sessionsThe only input training readsdata/*_pipeline.csvmodels:/svd_model_v1winner, rebuiltquery per armratings + sessions CSVServeA/B split and promotionEven user ids to v1, odd to v2Higher score wins, ties to v1Rebuilds the stack to redeployuser_id % 2 · 180 s live windowLearnv1 and v2LearnRun log and model registryEvery run with its provenanceNumbered versions per modelMLflow · the-sixth-sense · :6001Training jobTime ordered 80/20 splitTwo SVD candidates, v1 and v2Logs RMSE, data and commitsrc/model_training.py · Surprise SVDruns, params, RMSEtests on buildTest pipelineRuns the unit testsMocked database and modelJenkins · pytest · :8080
One store sits between the stream and training, and a registry between training and serving. The stream that feeds the models also judges them.

Learnings

  • Learning

    Evaluation that happens in production

    Offline error picks a candidate, but only what people do after a list is served says whether it works. Scoring that from the event log, in a process of its own, keeps the serving path fast and gives one number that dashboards and experiments share. This is what lets a team change the model often and know whether each change helped.

  • Learning

    A registry as the contract between training and serving

    When a model is a numbered version carrying its data, parameters and commit, the service never needs to know how it was built, and any answer can be traced to the run behind it. Swapping or rolling back becomes a name and a version instead of a file copied by hand. This is what makes retraining routine rather than a release.

  • Learning

    A recommender shapes its own data

    Every list served changes what people watch, and so what the next model learns: the team measured that titles with divided ratings were recommended more often than the rest, and that a tenth of the catalogue drew most of the ratings. New users arrive with no history at all. Planning for both from the start, with a default for strangers and a watch on what gets recommended, is what keeps a recommender useful after launch instead of narrowing over time.

← Back to coursework