From d98aada9045b20f496b95e8407bfe096a803e86c Mon Sep 17 00:00:00 2001 From: Alexander Platov Date: Sat, 26 Jul 2025 20:21:57 +0700 Subject: [PATCH] Fix PayloadTooLargeError in analytics collector (#9604) Signed-off-by: Alexander Platov --- dev/docker-compose.yaml | 1 + .../src/analyticsCollector.ts | 61 ++++++++++++++----- .../pod-analytics-collector/src/config.ts | 4 +- .../pod-analytics-collector/src/server.ts | 22 ++++++- 4 files changed, 71 insertions(+), 17 deletions(-) diff --git a/dev/docker-compose.yaml b/dev/docker-compose.yaml index b554ff930f..f9e487ac6a 100644 --- a/dev/docker-compose.yaml +++ b/dev/docker-compose.yaml @@ -414,6 +414,7 @@ services: # - STATS_URL=http://huly.local:4900 # - POSTHOG_HOST=${POSTHOG_HOST} # - POSTHOG_API_KEY=${POSTHOG_API_KEY} + # - MAX_PAYLOAD_SIZE=10mb msg2file: image: hardcoreeng/msg2file ports: diff --git a/packages/analytics-providers/src/analyticsCollector.ts b/packages/analytics-providers/src/analyticsCollector.ts index fb3e69b2e0..b12d4346e6 100644 --- a/packages/analytics-providers/src/analyticsCollector.ts +++ b/packages/analytics-providers/src/analyticsCollector.ts @@ -23,6 +23,8 @@ import { type QueuedEvent } from './types' export class AnalyticsCollectorProvider implements AnalyticProvider { private readonly collectIntervalMs = 5000 private readonly maxRetries = 3 + private readonly maxBatchSize = 100 + private readonly maxBatchSizeBytes = 5 * 1024 * 1024 // 5MB private readonly events: QueuedEvent[] = [] private collectTimer: any = null private url: string = '' @@ -70,26 +72,57 @@ export class AnalyticsCollectorProvider implements AnalyticProvider { const token = getMetadata(presentation.metadata.Token) ?? '' if (token === '') return - const eventsToSend = this.events.splice(0, this.events.length) + const batches = this.createBatches(this.events) + this.events.length = 0 - try { - const response = await fetch(`${this.url}/collect`, { - method: 'POST', - headers: { - Authorization: 'Bearer ' + token, - 'Content-Type': 'application/json' - }, - body: JSON.stringify(eventsToSend) - }) + for (const batch of batches) { + try { + const response = await fetch(`${this.url}/collect`, { + method: 'POST', + headers: { + Authorization: 'Bearer ' + token, + 'Content-Type': 'application/json' + }, + body: JSON.stringify(batch) + }) - if (!response.ok) { - this.handleFailedEvents(eventsToSend) + if (!response.ok) { + this.handleFailedEvents(batch) + } + } catch (err) { + this.handleFailedEvents(batch) } - } catch (err) { - this.handleFailedEvents(eventsToSend) } } + private createBatches (events: QueuedEvent[]): QueuedEvent[][] { + const batches: QueuedEvent[][] = [] + let currentBatch: QueuedEvent[] = [] + let currentBatchSize = 0 + + for (const event of events) { + const eventSize = JSON.stringify(event).length + + if ( + currentBatch.length >= this.maxBatchSize || + (currentBatchSize + eventSize > this.maxBatchSizeBytes && currentBatch.length > 0) + ) { + batches.push(currentBatch) + currentBatch = [] + currentBatchSize = 0 + } + + currentBatch.push(event) + currentBatchSize += eventSize + } + + if (currentBatch.length > 0) { + batches.push(currentBatch) + } + + return batches + } + private handleFailedEvents (failedEvents: QueuedEvent[]): void { const eventsToRetry: QueuedEvent[] = [] diff --git a/services/analytics-collector/pod-analytics-collector/src/config.ts b/services/analytics-collector/pod-analytics-collector/src/config.ts index a603f3c5f7..9b44d03e67 100644 --- a/services/analytics-collector/pod-analytics-collector/src/config.ts +++ b/services/analytics-collector/pod-analytics-collector/src/config.ts @@ -21,6 +21,7 @@ export interface Config { PostHogHost: string PostHogAPI: string SentryDSN?: string + MaxPayloadSize?: string } const parseNumber = (str: string | undefined): number | undefined => (str !== undefined ? Number(str) : undefined) @@ -33,7 +34,8 @@ const config: Config = (() => { AccountsUrl: process.env.ACCOUNTS_URL, PostHogHost: process.env.POSTHOG_HOST, PostHogAPI: process.env.POSTHOG_API_KEY, - SentryDSN: process.env.SENTRY_DSN ?? '' + SentryDSN: process.env.SENTRY_DSN ?? '', + MaxPayloadSize: process.env.MAX_PAYLOAD_SIZE ?? '10mb' } const missingEnv = (Object.keys(params) as Array).filter((key) => params[key] === undefined) diff --git a/services/analytics-collector/pod-analytics-collector/src/server.ts b/services/analytics-collector/pod-analytics-collector/src/server.ts index c3450508c2..87751c2ddc 100644 --- a/services/analytics-collector/pod-analytics-collector/src/server.ts +++ b/services/analytics-collector/pod-analytics-collector/src/server.ts @@ -50,7 +50,16 @@ function isContentValid (body: any[]): boolean { if (it == null) return true if (!('event' in it)) return true if (!('properties' in it)) return true - return !('timestamp' in it) + if (!('timestamp' in it)) return true + + const eventSize = JSON.stringify(it).length + if (eventSize > 1024 * 1024) { + // maximum 1MB per event + console.warn(`Event too large: ${eventSize} bytes, event: ${it.event}`) + return true + } + + return false }) } @@ -212,7 +221,7 @@ function preparePostHogEvent (event: AnalyticEvent, req: Request): Record { return preparePostHogEvent(event, req) @@ -236,6 +248,9 @@ export function createServer (): Express { batch: posthogEvents.reverse() } + const posthogPayloadSize = JSON.stringify(payload).length + console.log(`Sending to PostHog: ${posthogEvents.length} events, ${posthogPayloadSize} bytes`) + try { const response = await fetch(`${config.PostHogHost}/batch/`, { method: 'POST', @@ -249,6 +264,8 @@ export function createServer (): Express { if (!response.ok) { const errorText = await response.text() console.error(`PostHog API error: ${response.status} ${response.statusText}`, errorText) + } else { + console.log(`Successfully sent ${posthogEvents.length} events to PostHog`) } } catch (error) { console.error('Failed to send events to PostHog:', error) @@ -258,6 +275,7 @@ export function createServer (): Express { res.json({ received: events.length, processed: posthogEvents.length, + payloadSize, timestamp: new Date().toISOString() }) })