From a3f4750eae3a10a686b8b19240026b9ba38f335d Mon Sep 17 00:00:00 2001 From: Mike Cao Date: Mon, 23 Mar 2026 17:27:06 -0700 Subject: [PATCH] Send session replay data to Kafka when available instead of direct ClickHouse insert --- src/queries/sql/replays/saveRecording.ts | 32 ++++++++++++++---------- 1 file changed, 19 insertions(+), 13 deletions(-) diff --git a/src/queries/sql/replays/saveRecording.ts b/src/queries/sql/replays/saveRecording.ts index 73d2f9f58..01f58ca05 100644 --- a/src/queries/sql/replays/saveRecording.ts +++ b/src/queries/sql/replays/saveRecording.ts @@ -2,6 +2,7 @@ import { gzipSync } from 'node:zlib'; import clickhouse from '@/lib/clickhouse'; import { uuid } from '@/lib/crypto'; import { CLICKHOUSE, PRISMA, runQuery } from '@/lib/db'; +import kafka from '@/lib/kafka'; import prisma from '@/lib/prisma'; export interface SaveRecordingArgs { @@ -60,18 +61,23 @@ async function clickhouseQuery({ endedAt, }: SaveRecordingArgs) { const { insert, getUTCString } = clickhouse; + const { sendMessage } = kafka; - return insert('session_replay', [ - { - replay_id: uuid(), - website_id: websiteId, - session_id: sessionId, - visit_id: visitId, - chunk_index: chunkIndex, - events: JSON.stringify(events), - event_count: eventCount, - started_at: getUTCString(startedAt), - ended_at: getUTCString(endedAt), - }, - ]); + const message = { + replay_id: uuid(), + website_id: websiteId, + session_id: sessionId, + visit_id: visitId, + chunk_index: chunkIndex, + events: JSON.stringify(events), + event_count: eventCount, + started_at: getUTCString(startedAt), + ended_at: getUTCString(endedAt), + }; + + if (kafka.enabled) { + return sendMessage('session_replay', message); + } + + return insert('session_replay', [message]); }