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

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. 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 serves as the bridge between the pipeline and persistent storage.

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. 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():

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:

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.

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:

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 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.
  • ClickHouse stores aggregated session metadata via materialized views in 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.

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

The session_replay_events materialized view in 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, which constructs the processing chain using steps defined in separate files including record-session-event-step.ts and utilities from session-batch-manager.ts.

Have a question about this repo?

These articles cover the highlights, but your codebase questions are specific. Give your agent direct access to the source. Share this with your agent to get started:

Share the following with your agent to get started:
curl -s "https://instagit.com/install.md"

Works with
Claude Codex Cursor VS Code OpenClaw Any MCP Client

Maintain an open-source project? Get it listed too →