Adding Semantic IDs to your stack with Zipline
This post covers how teams building personalization, search or recommender systems can use Zipline to incorporate Semantic ID based generative recommendations into their stack.
What are Semantic IDs
Semantic IDs replace arbitrary item identifiers (item_123) with short sequences of discrete tokens derived from an item's content embedding, typically by passing the embedding through a residual quantizer so that similar items share token prefixes.
Because the tokens carry meaning, a transformer can treat your catalog like a language: it can generate the next items a user will engage with, and it can generalize to brand-new items that share tokens with familiar ones.
The idea was introduced in Google's TIGER paper (Recommender Systems with Generative Retrieval), and Eugene Yan's survey of LLM-era recsys techniques is another great deep dive.
When and why to use Semantic IDs
The strongest evidence comes from teams with large, fast-moving catalogs with cold-start and embedding-table bloat problems.
YouTube found that replacing hash-based item IDs with Semantic IDs in their production ranker improved CTR on new and tail items while shrinking embedding tables (Better Generalization with Semantic IDs, RecSys '24).
Kuaishou's OneRec replaced its cascaded retrieval/ranking stack with a single generative model over Semantic IDs, reporting gains in app stay time alongside roughly an order-of-magnitude reduction in serving compute, and >20% GMV lift when applied to their local-services surface (OneRec technical report).
Snap's practitioner's handbook (GRID) is a candid account of what it takes to make these systems work in production. If your catalog churns daily, your embedding tables are in the hundreds of GBs, or you want one model that serves both search and recs, Semantic IDs are worth an experiment.
Two ways to use Semantic IDs
There are two distinct architectures in the literature, and it's worth being explicit about which one you're building:
- SIDs as features (this post). Keep your existing discriminative ranker and feed it the user's recent activity as a sequence of Semantic IDs, alongside your production features. This is the YouTube approach: minimal architectural risk, a clean A/B against production, and it directly attacks cold-start and embedding-table bloat. Your model and serving stack stay exactly where they are — only the inputs change.
- Fully generative retrieval. Replace the retrieval stage (or the whole cascade) with a seq2seq transformer that takes the user's SID sequence as input and autoregressively decodes the Semantic IDs of the next items — TIGER and OneRec style. This is the bigger prize but also the bigger lift: it needs constrained decoding so the model only emits valid catalog SIDs, a reverse SID→item lookup, and a transformer inference service with its own latency profile.
The data infrastructure is the same for both: an item→SID mapping, a real-time enriched activity sequence, and point-in-time correct training data. Path 1 is the natural first experiment and everything built below carries over unchanged if you later go generative, where the join output already is the sequence-to-next-SID training corpus.
Challenges
Most of the difficulty is not the model, but rather the data infrastructure around it:
- Real-time sequence construction. The model's input is the user's recent activity, tokenized into Semantic IDs. That sequence must update within seconds of a click or view, which usually means a bespoke Flink/Kafka pipeline writing scalably to a KV store, and keeping it consistent with training.
- Point-in-time correct training data. Every training example needs the user's SID sequence as it existed at that moment. Getting this wrong (training on leaked future events) silently inflates offline metrics and craters online. Backfills with bespoke SQL pipelines can easily go out of sync with online serving pipelines, causing training/serving skew.
- Two codepaths, one feature. Batch backfill logic and streaming serving logic tend to be written twice, in two engines, by two teams, and training/serving skew creeps in between them.
- Catalog churn. New items need embeddings and SID assignments continuously, not in a weekly batch.
Zipline solves these challenges by offering a single API that can power:
- Scalable online pipelines: Zipline is built on the Chronon engine which powers search ranking at companies like Airbnb and Netflix.
- Fast and cheap backfills: See batch computation for more details.
- Consistency guaranteed: This is the core guarantee of the Chronon engine, a single API to define features and embeddings, and the ability to power training data generation and production serving pipelines with consistent results.
Building on Zipline
In the sections below we'll go step by step and illustrate how a user would build this pipeline on Zipline. The result will be:
- ~100 lines of Python config, no hand-written Spark or Flink jobs
- ~1s freshness on the user's Semantic ID sequence (streaming updates on every event)
- ~10ms fetch latency for the full model input at serving time
- Point-in-time correct training data backfilled over 90 days of history, consistent with online serving by construction
Users generally interact with Zipline using their preferred coding agent + the Zipline skill, so we'll start there and then go step by step through what the system produced.
Here's the prompt a data scientist might give Claude Code:
We run homepage recommendations with a two-tower retriever and a DNN ranker. Item features live in
item-catalog, we already produce content embeddings for every item intoitem-embeddings(updated hourly), and user clicks/views/add-to-carts land in theuser-activitiesstream. I want to run a Semantic ID experiment: quantize our item embeddings into 3-level Semantic IDs with RQ-VAE, maintain each user's last 50 activities as a real-time SID sequence, and feed that sequence into the ranker as a new feature alongside our existing ones. Backfill 90 days of point-in-time correct training data keyed offranking-requestsso I can retrain and A/B the new feature against production. Use the Zipline skill.
Let's walk through what gets built.
Step 1: Tokenize the catalog
First we need an item → Semantic ID mapping. The RQ-VAE tokenizer is defined with Zipline's Model API, taking the existing content embeddings as input:
# models/semantic_id_tokenizer.py
tokenizer = Model(
source="item-embeddings", # hourly content embeddings per item
model_type=ModelType.RQ_VAE,
params={
"num_levels": 3,
"codebook_size": 1024,
"embedding_dim": 256,
},
outputs=["semantic_id"], # e.g. [412, 87, 903]
)The output is materialized as an item-keyed feature, so any downstream pipeline can look up an item's SID. Because it's fed by the hourly embedding source, newly listed items get tokenized within the hour — no weekly batch job:
# group_bys/item_semantic_id.py
item_sid = GroupBy(
sources=["item-embeddings"],
keys=["item_id"],
aggregations=None, # simple lookup: latest SID per item
derivations=[tokenizer.semantic_id],
online=True,
)Step 2: Enrich the activity stream with Semantic IDs
Raw events in user-activities carry an item_id but not a Semantic ID, so before we can build sequences, each event needs to be enriched with the SID of the item it touched. In Zipline this is a chaining Join: the activity stream on the left, the item→SID mapping from Step 1 on the right. It runs as a streaming enrichment online and as a point-in-time correct join over history offline:
# joins/enriched_activities.py
enriched_activities = Join(
left="user-activities", # click / view / add-to-cart stream
right_parts=[
JoinPart(group_by=item_sid), # attaches semantic_id by item_id
],
online=True,
)Each event flows through with its SID attached — (user_id, action_type, ts, semantic_id) — within milliseconds of landing on the stream. Crucially, the offline version of this enrichment attaches the SID an item had at event time, so if the tokenizer is retrained and an item's SID changes, historical training data still reflects what the model would have seen.
Step 3: Build the real-time user sequence
Now the user's activity history as a sequence of Semantic IDs. This is a single GroupBy that consumes the enriched stream from Step 2 as its source (a JoinSource, in Chronon terms), with LAST_K maintaining the ordered sequence:
# group_bys/user_sid_sequence.py
user_sid_sequence = GroupBy(
sources=[JoinSource(join=enriched_activities)],
keys=["user_id"],
aggregations=[
Aggregation(
input_column="semantic_id",
operation=Operation.LAST_K(50),
),
Aggregation(
input_column="action_type",
operation=Operation.LAST_K(50),
),
],
online=True, # streaming updates to the KV store
)That's the whole real-time pipeline: event lands → enriched with its SID → appended to the user's sequence, end to end in about a second. From these two definitions Zipline generates the Flink jobs, the batch jobs that compute the same logic over warehouse history, and the serving endpoint — one set of definitions, no skew.
Step 4: Combine with existing features and backfill point-in-time training data
Training data is defined with a Join. The left side is the events we want to train on (historical ranking requests) and the right side is the new sequence feature plus the existing production feature groups:
# joins/homepage_ranker_sid_v1.py
v1 = Join(
left="ranking-requests", # one row per (user_id, ts) request
right_parts=[
JoinPart(group_by=user_sid_sequence), # the new experiment feature
JoinPart(group_by=user_engagement), # existing production features
JoinPart(group_by=item_stats),
],
online=True,
)Running the backfill produces 90 days of training rows where every SID sequence is accurate as of the millisecond of each request — the model never sees an event from the future. This is the property that makes offline results trustworthy, and it comes for free from the Join semantics rather than from a hand-audited backfill script.
Step 5: Train, serve, compare
This is where the pipeline hands off to your existing ML stack. The join output is a warehouse table with one row per historical ranking request, with the SID sequence sitting alongside your production features. The data scientist retrains the ranker exactly as they would for any new feature: point the existing training job (PyTorch, TensorFlow, SageMaker, Vertex, etc.) at the new table, add the sequence as an input to the model, and deploy through your normal model serving path.
Zipline's job at inference time is to feed that model. When a ranking request arrives, one fetcher call returns the full feature vector including fresh SID sequence in ~10ms, and your ranker consumes it like any other input:
features = fetcher.fetch_join(
join="homepage_ranker_sid_v1",
keys={"user_id": "u_1834"},
)Since the v1 join functions additively by bootstrapping existing data without recomputation, it facilitates a rigorous A/B comparison: you maintain your production baseline while introducing a single sequence feature. This infrastructure is purpose-built for the long term. These very same definitions transition directly into your production stack, and the enriched activity sequences you build here serve as the foundation for training a fully generative retrieval system down the road.
Bonus
In this experiment, we tried:
- 3-level Semantic IDs
- Last-50 user activities
- 90 days of training data
Once we see some results for this as a baseline, we might want to run a broader set of experiments, i.e.
- 3 vs. 4 vs. 5-level Semantic IDs
- last-50 vs. last-100 vs. last-150 user activities
- 90 vs. 180 vs. 365 days of training data
Zipline makes running these parallel experiments easy and cheap. See a demo here: Agentic Research with Zipline
Wrapping up
Semantic IDs are one of the clearest signals of where recommender systems are heading: smaller embedding tables, better cold-start behavior, and a path to unifying retrieval and ranking in a single generative model. The modeling ideas are public — the hard part has been the real-time, point-in-time correct data infrastructure underneath them. With Zipline, that part is ~100 lines of Python. If you'd like to run this experiment on your own stack, book a demo or reach out to the founders.