Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
26 changes: 24 additions & 2 deletions apps/worker/src/jobs/import.ts
Original file line number Diff line number Diff line change
Expand Up @@ -173,7 +173,18 @@ export async function importJob(job: Job<ImportQueuePayload>) {
) => 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 = [];
Expand Down Expand Up @@ -211,7 +222,18 @@ export async function importJob(job: Job<ImportQueuePayload>) {
) => 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 = [];
Expand Down
5 changes: 5 additions & 0 deletions packages/db/src/services/import.service.ts
Original file line number Diff line number Diff line change
Expand Up @@ -16,6 +16,7 @@ export interface ImportStageResult {
importId: string;
totalEvents: number;
insertedEvents: number;
csvDataSizeBytes?: number;
}

export interface ImportProgress {
Expand Down Expand Up @@ -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
Expand All @@ -87,6 +91,7 @@ export async function insertImportBatch(
importId,
totalEvents: events.length,
insertedEvents: events.length,
csvDataSizeBytes: csvDataSize,
};
}

Expand Down