Skip to content
chiltepin

Generated from: “How clickstream data gets from the app to the warehouse to the dashboards.”

Clickstream pipeline data pipeline diagram

AI-generated example · 2026-09-14. The source passes chiltepin check. System details and measurements are illustrative; review them before adapting this document.

DOCUMENTDRAFT

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.

SECTION 01 · Note

Assumptions

Note
The request named no tools, so this document picks a common stack. Web and mobile apps send events over HTTPS to a collector service. Kafka carries the raw and enriched streams. S3 is the landing zone, Snowflake is the warehouse, dbt builds the marts, and Looker reads them. Volume is about 9,000 events per second at peak and 350 million per day. Change the names and the shape below still holds.

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.

SECTION 02 · Data flow

Page view from app to dashboard

DFD
Data-flow diagram: 15 nodes, 14 flowsEXTWeb appEXTiOS and Android apps1CollectorDBKafka ·clickstream.raw2EnricherDBKafka ·clickstream.dlqDBKafka ·clickstream.enrichedEXTFraud scoring3S3 sinkDBS3 ·landing/clickstream/dt=4SnowpipeDBSnowflake ·RAW.CLICKSTREAM5dbt hourlyDBSnowflake · MARTSEXTLooker1234567891011121314
1POST /v1/events, batches of 502POST /v1/events, batches of 503publishes, key anonymous_id4consumes5adds geo, device, session_id6schema failures7consumes, sub-second8consumes9Parquet, 5-minute files10S3 event notification11COPY, append only12reads new partitions13sessions, funnels, daily_pages14reads
Legendprocessexternal entitydata storedata flow

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.

SECTION 03 · Event contract
EVENT·v3page_viewedchannelclickstream.raw

A user rendered a page or a screen. One event per navigation, sent from the client.

Producers (3)
web appiOS appAndroid app
Consumers (3)
enrichersession builderfraud scoring
deliveryat-least-onceorderingper-keykeyanonymous_idretention7d
Payload
FieldTypeDescription
event_iduuidClient-generated; consumers dedupe on it
#anonymous_iduuidBrowser or device identity; the partition key
?user_idstringSet after login; absent for guests
urlstringFull URL without query string; screen name on mobile
?referrerstringPrevious URL or the deep link source
ts_clienttimestampClient clock at render; may be skewed
app_versionstringSemantic version of the app build
# partition key · ? optional
Headers
FieldTypeDescription
schema_versionstringMatches the version above
sent_attimestampClient clock at send; the enricher replaces it with the broker time
Example
{ "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" }
Errors
ErrorWhen
SchemaMismatcha required field is missing or has the wrong type; the event goes to clickstream.dlq
ClockSkewts_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.

SECTION 04 · Steps

Rebuild one day in the marts

  1. Confirm 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.

    bash
    aws s3 ls s3://landing/clickstream/dt=2026-09-10/ | wc -l
  2. Reload the day into RAW

    Snowpipe skips files it already loaded. Force the day through the manual COPY to complete a partial load.

    bash
    snowsql -q "COPY INTO RAW.CLICKSTREAM FROM @landing/clickstream/dt=2026-09-10/ FORCE=TRUE"
  3. Rebuild the marts for the day

    The incremental models merge on event_id. The run takes about 12 minutes for one day.

    bash
    dbt run --select tag:clickstream --vars '{"replay_date": "2026-09-10"}'

    Post in

  4. Compare 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.

    bash
    dbt 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.

SECTION 05 · Service objectives

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.

Raw in Kafka30d

Events in clickstream.raw within 5 s of sent_at

Target99.9%
Current99.97%
30% of error budget used
Enriched30d

Events in clickstream.enriched within 30 s of sent_at

Target99.5%
Current99.8%
40% of error budget used
Landed in S330d

Five-minute files closed within 8 min of the window end

Target99%
Current99.6%
40% of error budget used
Marts refreshed30d

Hourly dbt run finished within 45 min of the hour

Target99%
Current97.2%
error budget exhausted
View the Markdown
```meta
title: Clickstream pipeline
subtitle: How a page view travels from the app to the warehouse and onto the dashboards, and how fresh it is at each stop.
tag: DRAFT
```

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.

```callout
tone: note
title: Assumptions
body: "The request named no tools, so this document picks a common stack. Web and mobile apps send events over HTTPS to a collector service. Kafka carries the raw and enriched streams. S3 is the landing zone, Snowflake is the warehouse, dbt builds the marts, and Looker reads them. Volume is about 9,000 events per second at peak and 350 million per day. Change the names and the shape below still holds."
```

## 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.

```dfd
id: clickstream-flow
title: Page view from app to dashboard
dir: LR
nodes:
  - { id: web, col: 1, row: 1, kind: external, name: Web app }
  - { id: mobile, col: 1, row: 2, kind: external, name: iOS and Android apps }
  - { id: collector, col: 2, row: 1, kind: process, name: Collector, num: 1 }
  - { id: raw, col: 3, row: 1, kind: store, name: "Kafka · clickstream.raw" }
  - { id: enricher, col: 4, row: 1, kind: process, name: Enricher, num: 2 }
  - { id: dlq, col: 4, row: 2, kind: store, name: "Kafka · clickstream.dlq" }
  - { id: enriched, col: 5, row: 1, kind: store, name: "Kafka · clickstream.enriched" }
  - { id: fraud, col: 5, row: 2, kind: external, name: Fraud scoring }
  - { id: sink, col: 6, row: 1, kind: process, name: S3 sink, num: 3 }
  - { id: s3, col: 7, row: 1, kind: store, name: "S3 · landing/clickstream/dt=" }
  - { id: loader, col: 8, row: 1, kind: process, name: Snowpipe, num: 4 }
  - { id: rawtbl, col: 9, row: 1, kind: store, name: "Snowflake · RAW.CLICKSTREAM" }
  - { id: dbt, col: 10, row: 1, kind: process, name: dbt hourly, num: 5 }
  - { id: marts, col: 11, row: 1, kind: store, name: "Snowflake · MARTS" }
  - { id: looker, col: 12, row: 1, kind: external, name: Looker }
edges:
  - web -> collector: "POST /v1/events, batches of 50"
  - mobile -> collector: "POST /v1/events, batches of 50"
  - collector -> raw: "publishes, key anonymous_id"
  - raw -> enricher: "consumes"
  - enricher -> enriched: "adds geo, device, session_id"
  - enricher -> dlq: "schema failures"
  - enriched -> fraud: "consumes, sub-second"
  - enriched -> sink: "consumes"
  - sink -> s3: "Parquet, 5-minute files"
  - s3 -> loader: "S3 event notification"
  - loader -> rawtbl: "COPY, append only"
  - rawtbl -> dbt: "reads new partitions"
  - dbt -> marts: "sessions, funnels, daily_pages"
  - marts -> looker: "reads"
```

## 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.

```eventcontract
id: page-viewed
name: page_viewed
version: v3
channel: clickstream.raw
summary: A user rendered a page or a screen. One event per navigation, sent from the client.
producers: [web app, iOS app, Android app]
consumers: [enricher, session builder, fraud scoring]
delivery: at-least-once
ordering: per-key
key: anonymous_id
retention: 7d
schema:
  - event_id uuid required — Client-generated; consumers dedupe on it
  - anonymous_id uuid required — Browser or device identity; the partition key
  - user_id string — Set after login; absent for guests
  - url string required — Full URL without query string; screen name on mobile
  - referrer string — Previous URL or the deep link source
  - ts_client timestamp required — Client clock at render; may be skewed
  - app_version string required — Semantic version of the app build
headers:
  - schema_version string required — Matches the version above
  - sent_at timestamp required — Client clock at send; the enricher replaces it with the broker time
example: |
  { "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" }
errors:
  - 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
note: "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.

```steps
id: clickstream-replay
title: Rebuild one day in the marts
items:
  - title: Confirm the landing files are complete
    body: 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.
    code: aws s3 ls s3://landing/clickstream/dt=2026-09-10/ | wc -l
    lang: bash
  - title: Reload the day into RAW
    body: Snowpipe skips files it already loaded. Force the day through the manual COPY to complete a partial load.
    code: "snowsql -q \"COPY INTO RAW.CLICKSTREAM FROM @landing/clickstream/dt=2026-09-10/ FORCE=TRUE\""
    lang: bash
  - title: Rebuild the marts for the day
    body: The incremental models merge on event_id. The run takes about 12 minutes for one day.
    code: "dbt run --select tag:clickstream --vars '{\"replay_date\": \"2026-09-10\"}'"
    lang: bash
    note: Post in #data-platform before you start. Looker shows the day as partial while the run is in progress.
  - title: Compare the row counts
    body: 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.
    code: "dbt test --select daily_pages --vars '{\"replay_date\": \"2026-09-10\"}'"
    lang: bash
```

## 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.

```slo
id: clickstream-freshness
title: Freshness per stop · 30 days to 12 Sep 2026
items:
  - { name: Raw in Kafka, sli: "Events in clickstream.raw within 5 s of sent_at", target: 99.9%, current: 99.97%, window: 30d, budget: 0.3 }
  - { name: Enriched, sli: "Events in clickstream.enriched within 30 s of sent_at", target: 99.5%, current: 99.8%, window: 30d, budget: 0.4 }
  - { name: Landed in S3, sli: "Five-minute files closed within 8 min of the window end", target: 99%, current: 99.6%, window: 30d, budget: 0.4 }
  - { name: Marts refreshed, sli: "Hourly dbt run finished within 45 min of the hour", target: 99%, current: 97.2%, window: 30d, budget: 1 }
description: "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."
```