A short-video app has an unusual property for a recommender: the user tells you whether you were right within seconds. Every swipe away, rewatch, like or share is a label, and a user can produce a hundred of them in a session. Interests also move fast; a topic that held someone's attention in the morning can bore them by evening. A system that retrains once a day is always learning yesterday's user. This is the design pressure that shapes a TikTok-style feed, and it explains why the interesting part is not the ranking model alone but the loop that keeps it current.

Precision about sources matters here. TikTok has not published a full description of its production For You system. The best public engineering source from ByteDance is the Monolith paper (Liu et al., arXiv 2209.07663, 2022), which describes a real-time recommendation system with a collisionless embedding table and online training, and states that it landed in the BytePlus Recommend product. This article uses Monolith for the learning loop, quoting its numbers, and a widely used industry pattern for the serving stages, marked as such. It is a reference design you can build, not a leak of anyone's internals.

Advertisement

What the product demands

  • Instant start. The next video must begin playing the moment the user swipes, so delivery and ranking happen before the swipe, not after.
  • Fresh interests. Signals from the last few minutes should change the next recommendations.
  • Fresh content. A video uploaded an hour ago by an unknown creator must be able to reach an audience, otherwise the supply side dies.
  • Sparse, huge ID spaces. Users, videos, creators, sounds and hashtags are each IDs with embeddings; new ones appear constantly and most old ones go quiet.
  • Implicit labels. Watch time and completion dominate, because most users never like or comment.

These demands pull in one direction: a short, cheap serving path and a learning path measured in minutes, with memory controlled despite an unbounded ID space.

The architecture end to end

Reference design for a short-video feed: serving path (top) and real-time learning loop (bottom)Clientprefetch next videosRetrievalmany sources, ~thousandsPre-ranklight model, ~hundredsRankmulti-task modelRe-rankdiversity, policyServing parameter serversembeddings, updated each minutelookupCDNvideo segmentsfetchAction logKafka: views, swipesFeature logKafka: features at serveeventsOnline joinerFlink: features + labelsTraining workersonline + batchexamplesTraining parameter serverscollisionless tablesgradientstouched keys / minuteThe learning loop is drawn from ByteDance's Monolith paper; the serving stages are a common industry pattern.
Serving path across the top, from client through four recommendation stages, with video bytes from a CDN. The real-time learning loop across the bottom, from two Kafka logs through a Flink joiner to training and then to the serving parameter servers.

On upload, a video is transcoded into several bitrates and short segments, pushed to a CDN, and analysed by content-understanding models that produce embeddings from frames, audio and text. Those embeddings let a brand-new video be retrieved before anyone has watched it. When a client asks for more videos, the request flows through four stages, each narrowing the set with a more expensive model. Everything the ranker saw is written to a feature log; everything the user then did goes to an action log. A streaming joiner pairs them into training examples, training updates embeddings continuously, and changed embeddings are pushed to the serving parameter servers every minute or so.

Advertisement

Serving: four stages and a prefetch

Retrieval gathers a few thousand candidates from several independent sources: an approximate nearest neighbour search over video embeddings using the user's embedding, as in vector search at scale; videos from followed creators; items similar to recent positive interactions; trending content in the user's region and language; and a fresh-content pool for exploration. Pre-ranking scores them with a cheap model and keeps a few hundred. Ranking runs the large multi-task model that predicts several outcomes per video. Re-ranking applies rules a model should not decide alone: do not show three videos by one creator in a row, mix topics, enforce safety and policy filters, and insert exploration slots.

The client hides latency by asking for a small batch of videos at a time and prefetching the first segment of the next few while the current one plays, so the swipe plays bytes already on the device. The video bytes themselves come from a CDN close to the user, with caching behaviour as described in CDN design. The ranking call is never on the critical path of a swipe; it runs ahead of need, which is why a budget of a few hundred milliseconds per batch is workable.

Labels: joining what was served with what happened

A training example needs the exact features the model saw at serve time and the outcome that followed. Recomputing features later from tables invites training-serving skew, because the tables have changed in the meantime. Monolith's design logs features at serve time to one Kafka topic and user actions to another, and uses a Flink job, its online joiner, to concatenate features with labels and emit training examples to a third topic read by both online and batch training. A simplified version of the joiner logic follows.

# Online joiner: pair what was served (features) with what the user did (labels).
# Keyed by request_id; emits one training example per impression.
JOIN_WINDOW_S = 600          # how long to wait for late actions before giving up

def on_feature_log(rec, state, timers):
    state[rec.request_id] = {"features": rec.features, "labels": {}, "ts": rec.ts}
    timers.register(rec.request_id, rec.ts + JOIN_WINDOW_S)

def on_action(evt, state):
    ex = state.get(evt.request_id)
    if ex is None:            # feature record expired or not yet arrived: count it
        metrics.incr("joiner.unmatched_action")
        return
    if evt.kind == "watch":
        ex["labels"]["watch_ratio"] = max(ex["labels"].get("watch_ratio", 0.0),
                                          evt.watched_ms / evt.duration_ms)
    else:                     # like, share, follow, comment, skip
        ex["labels"][evt.kind] = 1

def on_timer(request_id, state, out):
    ex = state.pop(request_id)
    out.emit(make_example(ex["features"], ex["labels"]))   # missing labels mean 0

The join window is a trade-off. Too short and late likes are lost, too long and memory grows and the model learns later. Videos are short, so most actions arrive within seconds of the impression. Monitor the unmatched-action counter: a rise means the feature log is lagging or being dropped. Most impressions end with no positive action, so training typically down-samples negatives; Monolith applies log odds correction at serving time so that the online model remains an unbiased estimator of the original distribution. The joiner's state handling follows ordinary Flink patterns from Apache Flink.

Embeddings without collisions

Most recommendation features are IDs, and the standard trick of hashing IDs into a fixed-size table means different IDs share a row. Monolith measured the cost on MovieLens: with 7.73% of user IDs and 2.86% of movie IDs colliding, the collisionless model consistently outperformed the hashed ones. Its answer is a key-value table built on a Cuckoo HashMap, which inserts new keys without colliding with existing ones and gives worst-case constant-time lookups. Two policies keep memory bounded: frequency filtering, which admits an ID only after it has occurred a tunable number of times, helped by a probabilistic filter, and expirable embeddings, which drop IDs inactive for a period set per table.

# Sketch of a collisionless embedding table with admission and expiry, after Monolith.
import time, numpy as np

class EmbeddingTable:
    def __init__(self, dim, min_count, ttl_s):
        self.dim, self.min_count, self.ttl_s = dim, min_count, ttl_s
        self.rows = {}           # id -> vector; production uses a cuckoo hash map
        self.seen = {}           # id -> count; production uses a probabilistic filter
        self.last = {}           # id -> last update time
        self.touched = set()     # ids updated since the last push to serving

    def lookup(self, fid):
        v = self.rows.get(fid)
        return v if v is not None else np.zeros(self.dim, np.float32)   # not admitted yet

    def apply_grad(self, fid, grad, lr):
        if fid not in self.rows:
            self.seen[fid] = self.seen.get(fid, 0) + 1
            if self.seen[fid] < self.min_count:
                return           # rare ids never get a row
            self.rows[fid] = np.random.normal(0, 0.01, self.dim).astype(np.float32)
        self.rows[fid] -= lr * grad
        self.last[fid] = time.time()
        self.touched.add(fid)

    def expire(self, now):
        for fid in [f for f, t in self.last.items() if now - t > self.ttl_s]:
            self.rows.pop(fid); self.last.pop(fid); self.touched.discard(fid)

    def drain_touched(self):     # called every minute by the sync job
        batch = {f: self.rows[f] for f in self.touched if f in self.rows}
        self.touched.clear()
        return batch

The sketch uses dictionaries to show the policies; the real structure matters for memory and speed. The expiry period is a product decision as much as a systems one: a user ID silent for a month and a video ID silent for a week have different chances of coming back.

Minute-level sync and what it buys

Training and serving run on separate parameter servers. The training side keeps a set of touched keys, IDs whose embeddings changed since the last push, and every minute only those rows are sent to the serving side. The paper's arithmetic: if 100,000 IDs change in a minute with 1024-dimensional embeddings, about 4 KB each, the push is roughly 400 MB per minute. Dense network weights change slowly, so they are synchronised daily at the lowest-traffic hour.

Sync interval (Criteo, DeepFM)Online-trained AUCBatch-trained AUC
5 hours79.6679.42
1 hour79.7879.44
30 minutes79.8079.43

Shorter intervals help steadily on a public dataset. In a live A/B test on an ads model, online training improved AUC over batch training by between 14.0% and 18.1% on each of seven days. Ads are not short videos, but the direction is the point: on drifting user behaviour, freshness is worth a lot. The general patterns for learning from streams are in online learning.

Reliability you can trade away

Monolith makes an argument worth borrowing: snapshot the parameter servers daily, not continuously. With parameters sharded over 1,000 servers and a 0.01% daily failure rate per server, one server fails roughly every ten days, losing one day of updates for the IDs on it. With 15 million daily users spread evenly, that is one day of feedback from about 15,000 users every ten days, which the paper judged acceptable, because user embeddings retrain quickly and dense weights barely move in a day. The lesson is to price reliability in model quality, not to assume every update is sacred.

Scoring and exploration

The ranking model predicts several outcomes, and the business decides how to combine them. A typical blend multiplies corrected probabilities and an expected watch ratio by tuned weights.

def final_score(p, w, sampling_rate):
    """p: per-task predicted probabilities from the ranking model.
    Negatives were down-sampled in training, so correct each probability first."""
    def correct(q):                      # undo negative down-sampling at rate r
        return q / (q + (1.0 - q) / sampling_rate)
    s = 0.0
    for task in ("finish", "like", "share", "follow", "comment"):
        s += w[task] * correct(p[task])
    s += w["watch"] * p["expected_watch_ratio"]   # regression head, no correction
    return s

Weights are where product values live: raise share and follow and you favour content that builds communities; raise watch ratio alone and you favour compulsive content. New videos need traffic to earn statistics, so a common pattern gives each one a small initial audience and promotes it to larger pools only if completion and engagement beat a threshold for its category. This keeps exploration cheap and makes the system fair to new creators without flooding feeds with untested content.

Worked example: sizing the loop

Take an illustrative service with 100 million daily users, each viewing 100 videos a day. That is 10 billion impressions a day, or 10,000,000,000 / 86,400, about 116,000 per second on average; plan for three times that at the evening peak, about 350,000 per second. Each impression produces one feature-log record and one or more action records, so the joiner holds state for every impression in its window: with a 600-second window at peak, about 210 million open entries. At 2 KB of features each, that is around 420 GB of state across the Flink cluster, which is why the window length is a capacity decision. For the embedding tables, 500 million active video IDs with 64-dimensional float32 vectors take 128 GB before optimiser state, and Adam-style optimisers roughly triple it; expiry and admission are what keep that number flat as uploads pile up.

Failure modes and trade-offs

  • Feedback loops. A model trained only on what it chose to show learns its own biases. Keep exploration traffic and log the serving propensity.
  • Joiner lag. If the joiner falls behind, labels arrive after the window closes and positives become negatives. Alert on unmatched actions and consumer lag.
  • Skew. Features recomputed at training time differ from those served. Log features at serve time, as above, and share feature code between serving and training, as described in feature stores.
  • Bad pushes. A minute-level sync spreads a corrupted update everywhere in a minute. Validate pushed rows for NaNs and norms and keep the previous snapshot ready.
  • Table growth. Without admission and expiry, embedding memory grows with total uploads, not active ones.
  • Hot content. A viral video concentrates CDN and parameter-server load; replicate hot rows and pre-warm edges.

What to do next

  1. Log the exact features served for every impression to a durable stream with a request ID.
  2. Build a streaming joiner with an explicit window, an unmatched-action metric and documented label definitions.
  3. Replace hashed ID embeddings with a keyed table that has admission thresholds and per-table expiry.
  4. Push only changed sparse rows to serving on a short interval; sync dense weights on a slower schedule.
  5. Measure model quality against sync interval offline before investing in minute-level infrastructure.
  6. Reserve exploration traffic for new content and log propensities to correct feedback loops.
  7. Read the Monolith paper itself for details this page summarises.
Key takeaway: A short-video feed is defined by fast labels and fast-moving interests. Serve through cheap-to-expensive stages and prefetch so ranking never sits on the swipe path; log served features and join them with actions in a streaming joiner; keep ID embeddings collision-free but bounded with admission and expiry; and push changed rows to serving every minute. ByteDance's Monolith paper shows the payoff of freshness and argues that daily snapshots are an acceptable reliability trade. TikTok's production system is not public, so treat this as a reference design to build on.