diff --git a/apps/worker/src/jobs/import.ts b/apps/worker/src/jobs/import.ts index b9986bd87..5d4dc2bb2 100644 --- a/apps/worker/src/jobs/import.ts +++ b/apps/worker/src/jobs/import.ts @@ -173,7 +173,18 @@ export async function importJob(job: Job) { ) => providerInstance!.transformEvent(event), ); - await insertImportBatch(transformedEvents, importId); + const result = await insertImportBatch(transformedEvents, importId); + + // Log actual CSV data size sent to ClickHouse + const csvSizeMB = result.csvDataSizeBytes + ? (result.csvDataSizeBytes / (1024 * 1024)).toFixed(2) + : 'N/A'; + + jobLogger.info('Inserted batch into ClickHouse', { + batchSize: result.insertedEvents, + csvDataSizeMB: csvSizeMB, + csvDataSizeBytes: result.csvDataSizeBytes || 0, + }); processedEvents += eventBatch.length; eventBatch = []; @@ -211,7 +222,18 @@ export async function importJob(job: Job) { ) => providerInstance!.transformEvent(event), ); - await insertImportBatch(transformedEvents, importId); + const result = await insertImportBatch(transformedEvents, importId); + + // Log actual CSV data size sent to ClickHouse + const csvSizeMB = result.csvDataSizeBytes + ? (result.csvDataSizeBytes / (1024 * 1024)).toFixed(2) + : 'N/A'; + + jobLogger.info('Inserted final batch into ClickHouse', { + batchSize: result.insertedEvents, + csvDataSizeMB: csvSizeMB, + csvDataSizeBytes: result.csvDataSizeBytes || 0, + }); processedEvents += eventBatch.length; eventBatch = []; diff --git a/packages/db/src/services/import.service.ts b/packages/db/src/services/import.service.ts index 6251ad712..8349920f5 100644 --- a/packages/db/src/services/import.service.ts +++ b/packages/db/src/services/import.service.ts @@ -16,6 +16,7 @@ export interface ImportStageResult { importId: string; totalEvents: number; insertedEvents: number; + csvDataSizeBytes?: number; } export interface ImportProgress { @@ -78,6 +79,9 @@ export async function insertImportBatch( return fields.join(','); }); + // Calculate actual CSV data size being sent to ClickHouse + const csvDataSize = csvRows.reduce((sum, row) => sum + row.length + 1, 0); + await chInsertCSV(TABLE_NAMES.events_imports, csvRows); // Explicitly release memory @@ -87,6 +91,7 @@ export async function insertImportBatch( importId, totalEvents: events.length, insertedEvents: events.length, + csvDataSizeBytes: csvDataSize, }; }