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:
- Parse headers and apply restrictions —
createParseHeadersStep()validates message headers, thencreateApplyEventRestrictionsStep()drops or overflows messages violating ingestion limits. - Team filter —
createTeamFilterStep(teamService)ensures events belong to known teams and enriches messages with team context. - Parse message —
createParseMessageStep()transforms raw Kafka payloads into typedParsedMessageDataobjects. - Library version monitor —
createLibVersionMonitorStep()emits warnings for outdated recording libraries. - Record to batch —
createRecordSessionEventStep()writes parsed events into the current session batch viaSessionBatchManager.
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-replaytopic. - The
createRecordSessionEventStepfunction inrecord-session-event-step.tsdelegates raw event storage toSessionBatchManager. SessionBatchManageraccumulates 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:
curl -s "https://instagit.com/install.md" Maintain an open-source project? Get it listed too →