Clickstream pipeline
How a page view travels from the app to the warehouse and onto the dashboards, and how fresh it is at each stop.
A click leaves the app as a JSON event and reaches a Looker dashboard about 40 minutes later. The path is stream first and batch last. The stream stages exist so that fraud and session features see events within seconds. The batch stages exist so that the marts are cheap and reproducible. The two halves meet at S3, which is the only place a day can be replayed from.
Assumptions
Where the data goes
The enricher is the one stage that changes an event. It adds geo, device class, and the session id, and it drops events that fail schema validation into a dead-letter topic. Nothing after the enricher rejects an event, so a bad field that passes the enricher reaches the dashboards. The collector writes nothing to disk; if Kafka is unreachable, the app retries with a 24-hour local buffer.
Page view from app to dashboard
The event on the wire
page_viewed is the highest-volume event and the one every mart depends on. Consumers key on anonymous_id, so all events from one browser land on one partition and the session builder sees them in order. A consumer that cannot handle a duplicate is wrong; Kafka delivers at least once and the app retries on timeout.
A user rendered a page or a screen. One event per navigation, sent from the client.
| Field | Type | Description | |
|---|---|---|---|
| event_id | uuid | Client-generated; consumers dedupe on it | |
| # | anonymous_id | uuid | Browser or device identity; the partition key |
| ? | user_id | string | Set after login; absent for guests |
| url | string | Full URL without query string; screen name on mobile | |
| ? | referrer | string | Previous URL or the deep link source |
| ts_client | timestamp | Client clock at render; may be skewed | |
| app_version | string | Semantic version of the app build |
| Field | Type | Description | |
|---|---|---|---|
| schema_version | string | Matches the version above | |
| sent_at | timestamp | Client clock at send; the enricher replaces it with the broker time |
{ "event_id": "5f0c…", "anonymous_id": "a1b2…", "url": "https://shop.example.com/products/oak-desk", "ts_client": "2026-09-12T10:41:07Z", "app_version": "4.18.0" }
| Error | When |
|---|---|
| SchemaMismatch | a required field is missing or has the wrong type; the event goes to clickstream.dlq |
| ClockSkew | ts_client is more than 24 hours from the broker time; the enricher keeps the event and flags it |
The client batches 50 events or 5 seconds, whichever comes first. A batch is retried for 24 hours from a local buffer.
Replaying a day
S3 is the replay source, never Kafka. The raw topic keeps seven days and the Parquet files keep 400, so a mart bug found in week three is fixed from S3. The procedure is idempotent: the loader appends, the dbt models merge on event_id, and a second run of the same day changes nothing.
Rebuild one day in the marts
- 1Confirm the landing files are complete
The day has 288 files, one per five minutes. A gap means the sink was down. Re-sink the day from Kafka first; that is only possible inside the seven-day retention.
bashaws s3 ls s3://landing/clickstream/dt=2026-09-10/ | wc -l - 2Reload the day into RAW
Snowpipe skips files it already loaded. Force the day through the manual COPY to complete a partial load.
bashsnowsql -q "COPY INTO RAW.CLICKSTREAM FROM @landing/clickstream/dt=2026-09-10/ FORCE=TRUE" - 3Rebuild the marts for the day
The incremental models merge on event_id. The run takes about 12 minutes for one day.
bashdbt run --select tag:clickstream --vars '{"replay_date": "2026-09-10"}'Post in
- 4Compare the row counts
The count in MARTS.DAILY_PAGES for the day must equal the count in RAW within 0.1%. A larger gap means the enricher dropped events and the day needs a look at clickstream.dlq.
bashdbt test --select daily_pages --vars '{"replay_date": "2026-09-10"}'
How fresh is fresh
Each stop has a freshness target and an owner. The stream targets are seconds because fraud scoring acts on them. The warehouse targets are minutes, because a dashboard that is 40 minutes old is good enough. A warehouse that loads every minute costs four times as much.
Freshness per stop · 30 days to 12 Sep 2026
The marts objective is over budget. Two dbt runs in the window took 70 minutes because the sessions model scanned a full month after a partition change. The fix shipped on 9 Sep.
Events in clickstream.raw within 5 s of sent_at
Events in clickstream.enriched within 30 s of sent_at
Five-minute files closed within 8 min of the window end
Hourly dbt run finished within 45 min of the hour