# How Session Replay Recording and Storage Works in PostHog: A Technical Deep Dive

> Discover how PostHog handles session replay recording and storage. Learn about event capture, Kafka streaming, S3 batching, and ClickHouse metadata persistence for efficient querying.

- Repository: [PostHog/posthog](https://github.com/PostHog/posthog)
- Tags: deep-dive
- Published: 2026-04-25

---

**PostHog captures client-side session events, streams them through Kafka, batches them into compressed files on S3, and persists aggregated metadata in ClickHouse for fast querying.**

The **session replay recording and storage** architecture in the `PostHog/posthog` repository separates real-time ingestion from long-term storage. Client recordings flow through a Node.js ingestion pipeline that manages batching, compression, and dual-writing to object storage and ClickHouse. This design enables efficient session playback while maintaining performant analytics queries.

## Session Replay Ingestion Pipeline

The ingestion pipeline is defined in [`nodejs/src/ingestion/session_replay/session-replay-pipeline.ts`](https://github.com/PostHog/posthog/blob/main/nodejs/src/ingestion/session_replay/session-replay-pipeline.ts). It orchestrates the flow from Kafka consumption to batch recording using a functional step-based architecture.

### Pipeline Steps

The `createSessionReplayPipeline` function composes six distinct processing steps:

1. **Parse headers and apply restrictions** — `createParseHeadersStep()` validates message headers, then `createApplyEventRestrictionsStep()` drops or overflows messages violating ingestion limits.
2. **Team filter** — `createTeamFilterStep(teamService)` ensures events belong to known teams and enriches messages with team context.
3. **Parse message** — `createParseMessageStep()` transforms raw Kafka payloads into typed `ParsedMessageData` objects.
4. **Library version monitor** — `createLibVersionMonitorStep()` emits warnings for outdated recording libraries.
5. **Record to batch** — `createRecordSessionEventStep()` writes parsed events into the current session batch via `SessionBatchManager`.

The pipeline uses the `newBatchPipelineBuilder` and returns a `BatchPipelineUnwrapper` that processes batches of Kafka `Message` objects.

## Recording Events with SessionBatchManager

The `createRecordSessionEventStep` function in [`nodejs/src/ingestion/session_replay/record-session-event-step.ts`](https://github.com/PostHog/posthog/blob/main/nodejs/src/ingestion/session_replay/record-session-event-step.ts) serves as the bridge between the pipeline and persistent storage.

```typescript
export function createRecordSessionEventStep<T extends RecordSessionEventStepInput>(
    config: RecordSessionEventStepConfig,
): ProcessingStep<T, T> {
    const { sessionBatchManager, isDebugLoggingEnabled } = config

    return async function recordSessionEventStep(input) {
        const { team, parsedMessage } = input

        // Reset revoked-session counters for this run
        SessionRecordingIngesterMetrics.resetSessionsRevoked()

        // Optional debug logging per Kafka partition
        if (isDebugLoggingEnabled(parsedMessage.metadata.partition)) {
            logger.debug('processing_session_recording', { /* metadata */ })
        }

        // Record the event in the current batch
        const batch = sessionBatchManager.getCurrentBatch()
        const messageWithTeam: MessageWithTeam = { team, message: parsedMessage }
        await batch.record(messageWithTeam)

        return ok(input)
    }
}

```

This step operates as a **side-effect**—it does not transform downstream data but hands the parsed event to `SessionBatchManager` for accumulation.

## S3 Storage and Batch Flushing

Batch management and flushing logic reside in [`nodejs/src/session-recording/sessions/session-batch-manager.ts`](https://github.com/PostHog/posthog/blob/main/nodejs/src/session-recording/sessions/session-batch-manager.ts). The `SessionBatchManager` class maintains a `currentBatch` and enforces flush boundaries based on configurable limits.

Key configuration properties include:

- **maxBatchSizeBytes** — Triggers flush when raw size exceeds limit
- **maxBatchAgeMs** — Triggers flush based on batch age
- **maxEventsPerSessionPerBatch** — Rate-limits per-session event density
- **offsetManager** — Tracks and commits Kafka offsets post-flush
- **fileStorage** — Handles compressed batch file writes to S3

When `shouldFlush()` returns true, the manager calls `flush()`:

```typescript
public async flush(): Promise<void> {
    logger.info('session_batch_manager_flushing', { batchSize: this.currentBatch.size })
    await this.currentBatch.flush()               // writes file, uploads to S3, updates ClickHouse
    this.currentBatch = new SessionBatchRecorder(/* same config */)
    this.lastFlushTime = Date.now()
}

```

The `SessionBatchRecorder` writes **compressed JSONL blocks** containing tuples of `[windowId, event]` and uploads them via `SessionBatchFileStorage`.

S3 configuration is controlled in [`posthog/settings/session_replay_v2.py`](https://github.com/PostHog/posthog/blob/main/posthog/settings/session_replay_v2.py):

```python
SESSION_RECORDING_V2_S3_ENDPOINT = os.getenv("SESSION_RECORDING_V2_S3_ENDPOINT", "")
SESSION_RECORDING_V2_S3_ACCESS_KEY_ID = os.getenv("SESSION_RECORDING_V2_S3_ACCESS_KEY_ID", "")
SESSION_RECORDING_V2_S3_SECRET_ACCESS_KEY = os.getenv("SESSION_RECORDING_V2_S3_SECRET_ACCESS_KEY", "")
SESSION_RECORDING_V2_S3_BUCKET = os.getenv("SESSION_RECORDING_V2_S3_BUCKET", "posthog")
SESSION_RECORDING_V2_S3_PREFIX = os.getenv("SESSION_RECORDING_V2_S3_PREFIX", "session_recordings")
SESSION_RECORDING_V2_S3_ENABLED = get_from_env("SESSION_RECORDING_V2_S3_ENABLED", True if DEBUG else False, type_cast=str_to_bool)

```

The bucket typically uses the prefix `session_recordings/` and supports S3-compatible endpoints like SeaweedFS for development environments.

## ClickHouse Schema and Metadata Aggregation

Aggregated session metadata is stored in ClickHouse via the schema defined in [`posthog/session_recordings/sql/session_replay_event_sql.py`](https://github.com/PostHog/posthog/blob/main/posthog/session_recordings/sql/session_replay_event_sql.py).

The system utilizes two primary tables:

- **kafka_session_replay_events** — Raw events written directly via the Kafka engine
- **session_replay_events** — Materialized view aggregating data into `writable_session_replay_events`

The materialized view (`session_replay_events_mv`) uses `AggregatingMergeTree` functions like `argMinState`, `groupArray`, and `sum` to collapse raw events into single rows per `(team_id, session_id)`. According to the source code, each aggregated row contains:

- Temporal bounds: `min_first_timestamp`, `max_last_timestamp`
- Navigation data: `first_url`, `all_urls`
- Interaction metrics: `click_count`, `keypress_count`
- Volume statistics: `size`, `event_count`, `message_count`
- AI-generated classifications: `ai_tags_fixed`, `ai_tags_freeform`

This aggregation allows the UI to list and filter recordings efficiently while fetching full session data from S3 only when playback is requested.

## Running the Pipeline

To wire the session replay ingestion into a Kafka consumer:

```typescript
import { createSessionReplayPipeline, runSessionReplayPipeline } from './session-replay-pipeline'
import { createKafkaConsumer } from '../kafka/consumer'
import { makeSessionBatchManager } from '../session-recording/sessions/factory'

// Build dependencies
const sessionBatchManager = makeSessionBatchManager(/* config from env/flags */)
const pipeline = createSessionReplayPipeline({
    outputs,                                   // IngestionOutputs for DLQ/overflow
    eventIngestionRestrictionManager,         // Rate-limit and team restrictions
    overflowEnabled: true,
    promiseScheduler: new PromiseScheduler(),
    teamService: new TeamService(),
    topHog: new TopHogRegistry(),
    sessionBatchManager,
    isDebugLoggingEnabled: (partition) => partition === 0,
})

// Consume and process
const consumer = createKafkaConsumer('session-replay')
consumer.run({
    eachBatchAutoResolve: false,
    eachBatch: async ({ batch, resolveOffset, heartbeat }) => {
        const msgs = batch.messages.map(m => ({ ...m }))
        await runSessionReplayPipeline(pipeline, msgs)
        // Commit offsets only after successful processing
        await consumer.commitOffsets([
            { topic: batch.topic, partition: batch.partition, offset: batch.lastOffset }
        ])
    },
})

```

This implementation feeds Kafka messages into the pipeline and commits offsets only after `SessionBatchManager` successfully flushes batches to S3 and ClickHouse.

## Summary

- **PostHog's session replay recording and storage** uses a Node.js ingestion pipeline to consume Kafka messages from the `session-replay` topic.
- The `createRecordSessionEventStep` function in [`record-session-event-step.ts`](https://github.com/PostHog/posthog/blob/main/record-session-event-step.ts) delegates raw event storage to `SessionBatchManager`.
- `SessionBatchManager` accumulates events in memory and flushes compressed JSONL files to S3 when size, age, or rate limits trigger.
- Configuration for S3 storage—including bucket, prefix, and credentials—is managed in [`session_replay_v2.py`](https://github.com/PostHog/posthog/blob/main/session_replay_v2.py).
- ClickHouse stores aggregated session metadata via materialized views in [`session_replay_event_sql.py`](https://github.com/PostHog/posthog/blob/main/session_replay_event_sql.py), enabling fast listing while S3 holds the full recording data.

## Frequently Asked Questions

### How does PostHog decide when to flush a session recording batch to S3?

PostHog evaluates three criteria in `SessionBatchManager.shouldFlush()`: **maxBatchSizeBytes** (raw byte limit), **maxBatchAgeMs** (time since last flush), and **maxEventsPerSessionPerBatch** (per-session rate limiting). When any threshold is exceeded, `flush()` compresses the current batch and uploads it to the configured S3 bucket.

### What format does PostHog use to store session replay data on S3?

Session recordings are stored as **compressed JSONL files** containing rows formatted as `[windowId, event]` tuples. These batches are written by `SessionBatchRecorder` and managed through `SessionBatchFileStorage` according to the settings in [`session_replay_v2.py`](https://github.com/PostHog/posthog/blob/main/session_replay_v2.py).

### How does ClickHouse enable fast session replay queries without scanning S3?

The `session_replay_events` materialized view in [`session_replay_event_sql.py`](https://github.com/PostHog/posthog/blob/main/session_replay_event_sql.py) pre-aggregates raw Kafka events into single rows per session using `AggregatingMergeTree`. This row contains summary statistics like timestamps, URL lists, and interaction counts, allowing the UI to filter and list recordings without fetching full data from S3 until playback is initiated.

### Where is the session replay ingestion pipeline defined in the codebase?

The pipeline definition resides in [`nodejs/src/ingestion/session_replay/session-replay-pipeline.ts`](https://github.com/PostHog/posthog/blob/main/nodejs/src/ingestion/session_replay/session-replay-pipeline.ts), which constructs the processing chain using steps defined in separate files including [`record-session-event-step.ts`](https://github.com/PostHog/posthog/blob/main/record-session-event-step.ts) and utilities from [`session-batch-manager.ts`](https://github.com/PostHog/posthog/blob/main/session-batch-manager.ts).