From e2bc8de898921fbe4b7ff23227db1b0224680243 Mon Sep 17 00:00:00 2001 From: Matt Aitken Date: Sat, 4 Apr 2026 08:43:14 +0100 Subject: [PATCH] We don't need separate stores for metrics --- apps/webapp/app/v3/otlpExporter.server.ts | 70 ++++++++--------------- 1 file changed, 24 insertions(+), 46 deletions(-) diff --git a/apps/webapp/app/v3/otlpExporter.server.ts b/apps/webapp/app/v3/otlpExporter.server.ts index a371212dd..878793111 100644 --- a/apps/webapp/app/v3/otlpExporter.server.ts +++ b/apps/webapp/app/v3/otlpExporter.server.ts @@ -74,19 +74,15 @@ class OTLPExporter { async exportMetrics(request: ExportMetricsServiceRequest): Promise { return await startSpan(this._tracer, "exportMetrics", async (span) => { - const metricsWithStores = this.#filterResourceMetrics(request.resourceMetrics).map( + const rows = this.#filterResourceMetrics(request.resourceMetrics).flatMap( (resourceMetrics) => - convertResourceMetricsToRowsWithStore( - resourceMetrics, - this._spanAttributeValueLengthLimit - ) + convertMetricsToClickhouseRows(resourceMetrics, this._spanAttributeValueLengthLimit) ); - const rowCount = metricsWithStores.reduce((acc, m) => acc + m.rows.length, 0); - span.setAttribute("metric_row_count", rowCount); + span.setAttribute("metric_row_count", rows.length); - if (rowCount > 0) { - await this.#exportMetricRows(metricsWithStores); + if (rows.length > 0) { + await this.#exportMetricRows(rows); } return ExportMetricsServiceResponse.create(); @@ -155,34 +151,31 @@ class OTLPExporter { return eventCount; } - async #exportMetricRows( - metricsWithStores: { rows: MetricsV1Input[]; taskEventStore: string }[] - ): Promise { + async #exportMetricRows(rows: MetricsV1Input[]): Promise { const routeCache = new Map(); const groups = new Map(); - for (const { rows, taskEventStore } of metricsWithStores) { - for (const row of rows) { - const routeKey = `${row.organization_id}\0${taskEventStore}`; - let resolved = routeCache.get(routeKey); - if (!resolved) { - resolved = this._clickhouseFactory.getEventRepositoryForOrganizationSync( - taskEventStore, - row.organization_id - ); - routeCache.set(routeKey, resolved); - } - let group = groups.get(resolved.key); - if (!group) { - group = { repository: resolved.repository, rows: [] }; - groups.set(resolved.key, group); - } - group.rows.push(row); + for (const row of rows) { + const routeKey = row.organization_id; + let resolved = routeCache.get(routeKey); + if (!resolved) { + resolved = this._clickhouseFactory.getEventRepositoryForOrganizationSync( + "clickhouse_v2", + row.organization_id + ); + routeCache.set(routeKey, resolved); } + + let group = groups.get(resolved.key); + if (!group) { + group = { repository: resolved.repository, rows: [] }; + groups.set(resolved.key, group); + } + group.rows.push(row); } - for (const [, { repository, rows }] of groups) { - repository.insertManyMetrics(rows); + for (const [, { repository, rows: groupedRows }] of groups) { + repository.insertManyMetrics(groupedRows); } } @@ -601,21 +594,6 @@ function convertMetricsToClickhouseRows( return rows; } -function convertResourceMetricsToRowsWithStore( - resourceMetrics: ResourceMetrics, - spanAttributeValueLengthLimit: number -): { rows: MetricsV1Input[]; taskEventStore: string } { - const resourceAttributes = resourceMetrics.resource?.attributes ?? []; - const taskEventStore = - extractStringAttribute(resourceAttributes, [SemanticInternalAttributes.TASK_EVENT_STORE]) ?? - env.EVENT_REPOSITORY_DEFAULT_STORE; - - return { - rows: convertMetricsToClickhouseRows(resourceMetrics, spanAttributeValueLengthLimit), - taskEventStore, - }; -} - // Prefixes injected by TaskContextMetricExporter — these are extracted into // the nested `trigger` key and should not appear as top-level user attributes. const INTERNAL_METRIC_ATTRIBUTE_PREFIXES = ["ctx.", "worker."];