From 8fd7eef1cfe34be7fe65f5f51a6917d16affaebf Mon Sep 17 00:00:00 2001 From: Igor Bernstein Date: Wed, 25 Feb 2026 20:57:05 -0500 Subject: [PATCH 1/4] chore: update the exporter to use MetricRegistry Change-Id: Ieb93b515d704f6ae67a200f48c32963dc326a3e0 --- .../data/v2/internal/csm/MetricRegistry.java | 2 +- .../data/v2/internal/csm/MetricsImpl.java | 7 +- .../internal/csm/attributes/ClientInfo.java | 2 +- .../v2/internal/csm/attributes/EnvInfo.java | 2 +- .../data/v2/stub/BigtableClientContext.java | 4 + .../BigtableCloudMonitoringExporter.java | 345 ++++++------------ .../data/v2/stub/metrics/Converter.java | 221 +++++++++++ .../BigtableCloudMonitoringExporterTest.java | 59 +-- 8 files changed, 372 insertions(+), 270 deletions(-) create mode 100644 google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricRegistry.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricRegistry.java index f485e79e4fcc..266ac7bc130f 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricRegistry.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricRegistry.java @@ -165,7 +165,7 @@ List getGrpcMetricNames() { return ImmutableList.copyOf(grpcMetricNames); } - MetricWrapper getMetric(String name) { + public MetricWrapper getMetric(String name) { return metrics.get(name); } diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java index 3ae54c531369..df9b4b4eaa9b 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java @@ -23,6 +23,7 @@ import com.google.cloud.bigtable.data.v2.BigtableDataSettings; import com.google.cloud.bigtable.data.v2.internal.csm.MetricRegistry.RecorderRegistry; import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo; +import com.google.cloud.bigtable.data.v2.internal.csm.attributes.EnvInfo; import com.google.cloud.bigtable.data.v2.stub.metrics.BigtableCloudMonitoringExporter; import com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants; import com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsTracerFactory; @@ -72,6 +73,7 @@ public class MetricsImpl implements Metrics, Closeable { private final List> tasks = new ArrayList<>(); public MetricsImpl( + MetricRegistry metricRegistry, ClientInfo clientInfo, ApiTracerFactory userTracerFactory, @Nullable OpenTelemetrySdk internalOtel, @@ -79,7 +81,7 @@ public MetricsImpl( Tagger ocTagger, StatsRecorder ocRecorder, ScheduledExecutorService executor) { - metricRegistry = new MetricRegistry(); + this.metricRegistry = metricRegistry; this.userTracerFactory = Preconditions.checkNotNull(userTracerFactory); this.internalOtel = internalOtel; @@ -168,6 +170,7 @@ public ChannelPoolMetricsTracer getChannelPoolMetricsTracer() { } public static OpenTelemetrySdk createBuiltinOtel( + MetricRegistry metricRegistry, ClientInfo clientInfo, @Nullable Credentials defaultCredentials, @Nullable String metricsEndpoint, @@ -194,7 +197,7 @@ public static OpenTelemetrySdk createBuiltinOtel( MetricExporter publicExporter = BigtableCloudMonitoringExporter.create( - clientInfo, credentials, metricsEndpoint, universeDomain, executor); + metricRegistry, EnvInfo::detect, clientInfo, credentials, metricsEndpoint, universeDomain); PeriodicMetricReaderBuilder readerBuilder = PeriodicMetricReader.builder(publicExporter).setExecutor(executor); meterProvider.registerMetricReader(readerBuilder.build()); diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/ClientInfo.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/ClientInfo.java index 0d4717dfe9a2..64c4b211b297 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/ClientInfo.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/ClientInfo.java @@ -42,7 +42,7 @@ public static Builder builder() { @AutoValue.Builder public abstract static class Builder { - protected abstract Builder setClientName(String name); + public abstract Builder setClientName(String name); public abstract Builder setInstanceName(InstanceName name); diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/EnvInfo.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/EnvInfo.java index cfc9182881b6..b7afb73ee9e3 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/EnvInfo.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/attributes/EnvInfo.java @@ -77,7 +77,7 @@ public static Builder builder() { @AutoValue.Builder public abstract static class Builder { - protected abstract Builder setUid(String uid); + public abstract Builder setUid(String uid); public abstract Builder setPlatform(String platform); diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/BigtableClientContext.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/BigtableClientContext.java index c82b0a8d02ba..c4bef24798f5 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/BigtableClientContext.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/BigtableClientContext.java @@ -28,6 +28,7 @@ import com.google.auth.oauth2.ServiceAccountJwtAccessCredentials; import com.google.bigtable.v2.InstanceName; import com.google.cloud.bigtable.data.v2.internal.JwtCredentialsWithAudience; +import com.google.cloud.bigtable.data.v2.internal.csm.MetricRegistry; import com.google.cloud.bigtable.data.v2.internal.csm.Metrics; import com.google.cloud.bigtable.data.v2.internal.csm.MetricsImpl; import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo; @@ -101,6 +102,7 @@ public static BigtableClientContext create( FixedExecutorProvider.create(backgroundExecutor, shouldAutoClose); builder.setBackgroundExecutorProvider(executorProvider); + MetricRegistry metricRegistry = new MetricRegistry(); // Set up OpenTelemetry @Nullable OpenTelemetry userOtel = null; if (settings.getMetricsProvider() instanceof CustomOpenTelemetryMetricsProvider) { @@ -113,6 +115,7 @@ public static BigtableClientContext create( if (settings.areInternalMetricsEnabled()) { builtinOtel = MetricsImpl.createBuiltinOtel( + metricRegistry, clientInfo, credentials, settings.getMetricsEndpoint(), @@ -125,6 +128,7 @@ public static BigtableClientContext create( Metrics metrics = new MetricsImpl( + metricRegistry, clientInfo, settings.getTracerFactory(), builtinOtel, diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java index 2aba290aff25..eec69d47e58c 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java @@ -15,40 +15,20 @@ */ package com.google.cloud.bigtable.data.v2.stub.metrics; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.APPLICATION_BLOCKING_LATENCIES_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.ATTEMPT_LATENCIES2_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.ATTEMPT_LATENCIES_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.CLIENT_BLOCKING_LATENCIES_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.CONNECTIVITY_ERROR_COUNT_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.FIRST_RESPONSE_LATENCIES_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.METER_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.OPERATION_LATENCIES_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.PER_CONNECTION_ERROR_COUNT_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.REMAINING_DEADLINE_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.RETRY_COUNT_NAME; -import static com.google.cloud.bigtable.data.v2.stub.metrics.BuiltinMetricsConstants.SERVER_LATENCIES_NAME; - -import com.google.api.MonitoredResource; import com.google.api.core.ApiFuture; import com.google.api.core.ApiFutureCallback; import com.google.api.core.ApiFutures; -import com.google.api.core.InternalApi; import com.google.api.gax.core.CredentialsProvider; import com.google.api.gax.core.FixedCredentialsProvider; -import com.google.api.gax.core.FixedExecutorProvider; import com.google.api.gax.core.NoCredentialsProvider; import com.google.api.gax.rpc.PermissionDeniedException; import com.google.auth.Credentials; +import com.google.cloud.bigtable.data.v2.internal.csm.MetricRegistry; import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo; +import com.google.cloud.bigtable.data.v2.internal.csm.attributes.EnvInfo; import com.google.cloud.monitoring.v3.MetricServiceClient; import com.google.cloud.monitoring.v3.MetricServiceSettings; -import com.google.common.annotations.VisibleForTesting; import com.google.common.base.Preconditions; -import com.google.common.base.Supplier; -import com.google.common.base.Suppliers; -import com.google.common.collect.ImmutableList; -import com.google.common.collect.ImmutableMap; -import com.google.common.collect.ImmutableSet; import com.google.common.collect.Iterables; import com.google.common.util.concurrent.MoreExecutors; import com.google.monitoring.v3.CreateTimeSeriesRequest; @@ -65,86 +45,59 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.Collection; -import java.util.HashMap; import java.util.List; import java.util.Map; +import java.util.Map.Entry; import java.util.Optional; -import java.util.concurrent.ScheduledExecutorService; import java.util.concurrent.atomic.AtomicBoolean; -import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; +import java.util.function.Supplier; import java.util.logging.Level; import java.util.logging.Logger; -import java.util.stream.Collectors; import javax.annotation.Nullable; -/** - * Bigtable Cloud Monitoring OpenTelemetry Exporter. - * - *

The exporter will look for all bigtable owned metrics under bigtable.googleapis.com - * instrumentation scope and upload it via the Google Cloud Monitoring API. - */ -@InternalApi -public final class BigtableCloudMonitoringExporter implements MetricExporter { - - private static final Logger logger = +public class BigtableCloudMonitoringExporter implements MetricExporter { + private static final Logger LOGGER = Logger.getLogger(BigtableCloudMonitoringExporter.class.getName()); - // This system property can be used to override the monitoring endpoint - // to a different environment. It's meant for internal testing only and - // will be removed in future versions. Use settings in EnhancedBigtableStubSettings - // to override the endpoint. - @Deprecated @Nullable - private static final String MONITORING_ENDPOINT_OVERRIDE_SYS_PROP = - System.getProperty("bigtable.test-monitoring-endpoint"); - - private static final String APPLICATION_RESOURCE_PROJECT_ID = "project_id"; - // This the quota limit from Cloud Monitoring. More details in // https://cloud.google.com/monitoring/quotas#custom_metrics_quotas. private static final int EXPORT_BATCH_SIZE_LIMIT = 200; + private final Supplier envInfo; + private final ClientInfo clientInfo; + private MetricRegistry metricRegistry; private final MetricServiceClient client; - private final List timeSeriesConverters; - - private final AtomicBoolean isShutdown = new AtomicBoolean(false); - + private final AtomicReference state; private CompletableResultCode lastExportCode; - private final AtomicBoolean exportFailureLogged = new AtomicBoolean(false); + private enum State { + Running, + Closing, + Closed + } + public static BigtableCloudMonitoringExporter create( + MetricRegistry metricRegistry, + Supplier envInfo, ClientInfo clientInfo, @Nullable Credentials credentials, @Nullable String endpoint, - String universeDomain, - @Nullable ScheduledExecutorService executorService) + String universeDomain) throws IOException { + Preconditions.checkNotNull(universeDomain); - MetricServiceSettings.Builder settingsBuilder = MetricServiceSettings.newBuilder(); - CredentialsProvider credentialsProvider = - Optional.ofNullable(credentials) - .map(FixedCredentialsProvider::create) - .orElse(NoCredentialsProvider.create()); - settingsBuilder.setCredentialsProvider(credentialsProvider); - - settingsBuilder.setUniverseDomain(universeDomain); - - // If background executor is not null, use it for the monitoring client. This allows us to - // share the same background executor with the data client. When it's null, the monitoring - // client will create a new executor service from InstantiatingExecutorProvider. It could be - // null if someone uses a CustomOpenTelemetryMetricsProvider#setupSdkMeterProvider without - // the executor. - if (executorService != null) { - settingsBuilder.setBackgroundExecutorProvider(FixedExecutorProvider.create(executorService)); - } - if (MONITORING_ENDPOINT_OVERRIDE_SYS_PROP != null) { - logger.warning( - "Setting the monitoring endpoint through system variable will be removed in future" - + " versions"); - settingsBuilder.setEndpoint(MONITORING_ENDPOINT_OVERRIDE_SYS_PROP); - } + MetricServiceSettings.Builder settingsBuilder = + MetricServiceSettings.newBuilder() + .setUniverseDomain(universeDomain) + .setCredentialsProvider( + Optional.ofNullable(credentials) + .map(FixedCredentialsProvider::create) + .orElse(NoCredentialsProvider.create())); + if (endpoint != null) { settingsBuilder.setEndpoint(endpoint); } @@ -154,105 +107,101 @@ public static BigtableCloudMonitoringExporter create( // it as not retried for now. settingsBuilder.createServiceTimeSeriesSettings().setSimpleTimeoutNoRetriesDuration(timeout); - ImmutableList converters = - ImmutableList.of( - new PublicTimeSeriesConverter(), - new InternalTimeSeriesConverter( - Suppliers.memoize( - () -> BigtableExporterUtils.createInternalMonitoredResource(clientInfo)))); - return new BigtableCloudMonitoringExporter( - MetricServiceClient.create(settingsBuilder.build()), converters); + metricRegistry, envInfo, clientInfo, MetricServiceClient.create(settingsBuilder.build())); } - @VisibleForTesting BigtableCloudMonitoringExporter( - MetricServiceClient client, List converters) { + MetricRegistry metricRegistry, + Supplier envInfo, + ClientInfo clientInfo, + MetricServiceClient client) { + this.metricRegistry = metricRegistry; + this.envInfo = envInfo; + this.clientInfo = clientInfo; this.client = client; - this.timeSeriesConverters = ImmutableList.copyOf(converters); + this.state = new AtomicReference<>(State.Running); + } + + public void close() { + client.close(); } @Override public CompletableResultCode export(Collection metricData) { - Preconditions.checkState(!isShutdown.get(), "Exporter is shutting down"); + Preconditions.checkState(state.get() != State.Closed, "Exporter is closed"); + + if (metricRegistry == null) { + String msg = "Bigtable exporter tried to export before fully configured"; + LOGGER.warning(msg); + return CompletableResultCode.ofExceptionalFailure(new IllegalStateException(msg)); + } lastExportCode = doExport(metricData); return lastExportCode; } - /** Export metrics associated with a BigtableTable resource. */ private CompletableResultCode doExport(Collection metricData) { - Map> bigtableTimeSeries = new HashMap<>(); - - List results = new ArrayList<>(); - - for (TimeSeriesConverter c : timeSeriesConverters) { - try { - for (Map.Entry> e : c.convert(metricData).entrySet()) { - bigtableTimeSeries - .computeIfAbsent(e.getKey(), (k) -> new ArrayList<>()) - .addAll(e.getValue()); - } - results.add(CompletableResultCode.ofSuccess()); - } catch (Throwable t) { - logger.log( - Level.WARNING, - String.format( - "Failed to convert %s metric data to cloud monitoring timeseries.", c.name), - t); - results.add(CompletableResultCode.ofExceptionalFailure(t)); + Map> converted; + + try { + converted = new Converter(metricRegistry, envInfo.get(), clientInfo).convertAll(metricData); + } catch (Throwable t) { + if (exportFailureLogged.compareAndSet(false, true)) { + LOGGER.log(Level.WARNING, "Failed to compose metrics for export", t); } - } - CompletableResultCode overall = CompletableResultCode.ofAll(results); - if (!overall.isSuccess()) { - return overall; + + return CompletableResultCode.ofExceptionalFailure(t); } - // Skips exporting if there's none - if (bigtableTimeSeries.isEmpty()) { - return CompletableResultCode.ofSuccess(); + List> futures = new ArrayList<>(); + + for (Entry> e : converted.entrySet()) { + futures.addAll(exportTimeSeries(e.getKey(), e.getValue())); } CompletableResultCode exportCode = new CompletableResultCode(); - bigtableTimeSeries.forEach( - (projectName, ts) -> { - ApiFuture> future = exportTimeSeries(projectName, ts); - ApiFutures.addCallback( - future, - new ApiFutureCallback>() { - @Override - public void onFailure(Throwable throwable) { - if (exportFailureLogged.compareAndSet(false, true)) { - String msg = "createServiceTimeSeries request failed"; - if (throwable instanceof PermissionDeniedException) { - msg += - String.format( - " Need monitoring metric writer permission on project=%s. Follow" - + " https://cloud.google.com/bigtable/docs/client-side-metrics-setup" - + " to set up permissions.", - projectName.getProject()); - } - logger.log(Level.WARNING, msg, throwable); - } - exportCode.fail(); - } - - @Override - public void onSuccess(List emptyList) { - // When an export succeeded reset the export failure flag to false so if there's a - // transient failure it'll be logged. - exportFailureLogged.set(false); - exportCode.succeed(); - } - }, - MoreExecutors.directExecutor()); - }); + StackTraceElement[] stackTrace = Thread.currentThread().getStackTrace(); + + ApiFutures.addCallback( + ApiFutures.allAsList(futures), + new ApiFutureCallback>() { + @Override + public void onFailure(Throwable throwable) { + if (exportFailureLogged.compareAndSet(false, true)) { + String msg = "createServiceTimeSeries request failed"; + if (throwable instanceof PermissionDeniedException) { + msg += + String.format( + " Need monitoring metric writer permission on project=%s. Follow" + + " https://cloud.google.com/bigtable/docs/client-side-metrics-setup" + + " to set up permissions.", + clientInfo.getInstanceName().getProject()); + } + RuntimeException asyncWrapper = new RuntimeException("export failed", throwable); + asyncWrapper.setStackTrace(stackTrace); + + if (state.get() != State.Closing) { + // ignore the export warning when client is shutting down + LOGGER.log(Level.WARNING, msg, asyncWrapper); + } + } + exportCode.fail(); + } + + @Override + public void onSuccess(List objects) { + exportFailureLogged.set(false); + exportCode.succeed(); + } + }, + MoreExecutors.directExecutor()); return exportCode; } - private ApiFuture> exportTimeSeries( - ProjectName projectName, List timeSeries) { + private List> exportTimeSeries( + ProjectName projectName, Collection timeSeries) { List> batchResults = new ArrayList<>(); for (List batch : Iterables.partition(timeSeries, EXPORT_BATCH_SIZE_LIMIT)) { @@ -265,7 +214,7 @@ private ApiFuture> exportTimeSeries( batchResults.add(f); } - return ApiFutures.allAsList(batchResults); + return batchResults; } @Override @@ -278,8 +227,9 @@ public CompletableResultCode flush() { @Override public CompletableResultCode shutdown() { - if (!isShutdown.compareAndSet(false, true)) { - logger.log(Level.WARNING, "shutdown is called multiple times"); + State prevState = state.getAndSet(State.Closed); + if (prevState == State.Closed) { + LOGGER.log(Level.WARNING, "shutdown is called multiple times"); return CompletableResultCode.ofSuccess(); } CompletableResultCode flushResult = flush(); @@ -290,7 +240,7 @@ public CompletableResultCode shutdown() { try { client.shutdown(); } catch (Throwable e) { - logger.log(Level.WARNING, "failed to shutdown the monitoring client", e); + LOGGER.log(Level.WARNING, "failed to shutdown the monitoring client", e); throwable = e; } if (throwable != null) { @@ -299,97 +249,16 @@ public CompletableResultCode shutdown() { shutdownResult.succeed(); } }); + return CompletableResultCode.ofAll(Arrays.asList(flushResult, shutdownResult)); } - /** - * For Google Cloud Monitoring always return CUMULATIVE to keep track of the cumulative value of a - * metric over time. - */ @Override public AggregationTemporality getAggregationTemporality(InstrumentType instrumentType) { return AggregationTemporality.CUMULATIVE; } - abstract static class TimeSeriesConverter { - private final String name; - - TimeSeriesConverter(String name) { - this.name = name; - } - - abstract Map> convert(Collection metricData); - } - - static class PublicTimeSeriesConverter extends TimeSeriesConverter { - private static final ImmutableList BIGTABLE_TABLE_METRICS = - ImmutableSet.of( - OPERATION_LATENCIES_NAME, - ATTEMPT_LATENCIES_NAME, - ATTEMPT_LATENCIES2_NAME, - SERVER_LATENCIES_NAME, - FIRST_RESPONSE_LATENCIES_NAME, - CLIENT_BLOCKING_LATENCIES_NAME, - APPLICATION_BLOCKING_LATENCIES_NAME, - RETRY_COUNT_NAME, - CONNECTIVITY_ERROR_COUNT_NAME, - REMAINING_DEADLINE_NAME) - .stream() - .map(m -> METER_NAME + m) - .collect(ImmutableList.toImmutableList()); - - private static final AtomicLong nextTaskIdSuffix = new AtomicLong(); - private final String taskId; - - PublicTimeSeriesConverter() { - this( - BigtableExporterUtils.DEFAULT_TASK_VALUE.get() - + "-" - + nextTaskIdSuffix.getAndIncrement()); - } - - PublicTimeSeriesConverter(String taskId) { - super("table metrics"); - this.taskId = taskId; - } - - @Override - public Map> convert(Collection metricData) { - List relevantData = - metricData.stream() - .filter(md -> BIGTABLE_TABLE_METRICS.contains(md.getName())) - .collect(Collectors.toList()); - if (relevantData.isEmpty()) { - return ImmutableMap.of(); - } - return BigtableExporterUtils.convertToBigtableTimeSeries(relevantData, taskId); - } - } - - static class InternalTimeSeriesConverter extends TimeSeriesConverter { - private static final ImmutableList APPLICATION_METRICS = - ImmutableSet.of(PER_CONNECTION_ERROR_COUNT_NAME).stream() - .map(m -> METER_NAME + m) - .collect(ImmutableList.toImmutableList()); - - private final Supplier monitoredResource; - - InternalTimeSeriesConverter(Supplier monitoredResource) { - super("client metrics"); - this.monitoredResource = monitoredResource; - } - - @Override - public Map> convert(Collection metricData) { - MonitoredResource monitoredResource = this.monitoredResource.get(); - if (monitoredResource == null) { - return ImmutableMap.of(); - } - - return ImmutableMap.of( - ProjectName.of(monitoredResource.getLabelsOrThrow(APPLICATION_RESOURCE_PROJECT_ID)), - BigtableExporterUtils.convertToApplicationResourceTimeSeries( - metricData, monitoredResource)); - } + public void prepareForShutdown() { + state.compareAndSet(State.Running, State.Closing); } } diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java new file mode 100644 index 000000000000..03b1f1519e27 --- /dev/null +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java @@ -0,0 +1,221 @@ +/* + * Copyright 2025 Google LLC + * + * Licensed under the Apache License, Version 2.0 (the "License"); + * you may not use this file except in compliance with the License. + * You may obtain a copy of the License at + * + * https://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package com.google.cloud.bigtable.data.v2.stub.metrics; + +import static com.google.api.MetricDescriptor.MetricKind.CUMULATIVE; +import static com.google.api.MetricDescriptor.MetricKind.GAUGE; +import static com.google.api.MetricDescriptor.MetricKind.UNRECOGNIZED; +import static com.google.api.MetricDescriptor.ValueType.DISTRIBUTION; +import static com.google.api.MetricDescriptor.ValueType.DOUBLE; +import static com.google.api.MetricDescriptor.ValueType.INT64; + +import com.google.api.Distribution; +import com.google.api.Distribution.BucketOptions; +import com.google.api.Distribution.BucketOptions.Explicit; +import com.google.api.Metric; +import com.google.api.MetricDescriptor.MetricKind; +import com.google.api.MetricDescriptor.ValueType; +import com.google.cloud.bigtable.data.v2.internal.csm.MetricRegistry; +import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo; +import com.google.cloud.bigtable.data.v2.internal.csm.attributes.EnvInfo; +import com.google.cloud.bigtable.data.v2.internal.csm.metrics.GrpcMetric; +import com.google.cloud.bigtable.data.v2.internal.csm.metrics.MetricWrapper; +import com.google.common.collect.ImmutableListMultimap; +import com.google.common.collect.ImmutableMultimap; +import com.google.common.collect.ImmutableSet; +import com.google.common.collect.Multimap; +import com.google.monitoring.v3.Point; +import com.google.monitoring.v3.ProjectName; +import com.google.monitoring.v3.TimeInterval; +import com.google.monitoring.v3.TimeSeries; +import com.google.monitoring.v3.TypedValue; +import com.google.protobuf.util.Timestamps; +import io.opentelemetry.sdk.metrics.data.AggregationTemporality; +import io.opentelemetry.sdk.metrics.data.DoublePointData; +import io.opentelemetry.sdk.metrics.data.HistogramData; +import io.opentelemetry.sdk.metrics.data.HistogramPointData; +import io.opentelemetry.sdk.metrics.data.LongPointData; +import io.opentelemetry.sdk.metrics.data.MetricData; +import io.opentelemetry.sdk.metrics.data.MetricDataType; +import io.opentelemetry.sdk.metrics.data.PointData; +import io.opentelemetry.sdk.metrics.data.SumData; +import java.util.Collection; +import java.util.Map; +import java.util.Set; +import java.util.logging.Level; +import java.util.logging.Logger; + +/** + * Helper for exporting metrics from Opentelemetry to Cloud Monitoring. + * + *

Takes collection {@link MetricData} and uses the {@link MetricWrapper}s defined in {@link + * MetricRegistry} to compose both the {@link com.google.api.MonitoredResource} and {@link Point}. + */ +class Converter { + private static final Logger LOGGER = Logger.getLogger(Converter.class.getName()); + + private final MetricRegistry metricRegistry; + private final EnvInfo envInfo; + private final ClientInfo clientInfo; + + Converter(MetricRegistry metricRegistry, EnvInfo envInfo, ClientInfo clientInfo) { + this.metricRegistry = metricRegistry; + this.envInfo = envInfo; + this.clientInfo = clientInfo; + } + + Map> convertAll(Collection otelMetrics) { + ImmutableMultimap.Builder builder = ImmutableMultimap.builder(); + + for (MetricData metricData : otelMetrics) { + Multimap perProject = convertMetricData(metricData); + builder.putAll(perProject); + } + return builder.build().asMap(); + } + + private Multimap convertMetricData(MetricData metricData) { + MetricWrapper metricDef = metricRegistry.getMetric(metricData.getName()); + if (metricDef == null) { + LOGGER.log(Level.FINE, "Skipping unexpected metric: {}", metricData.getName()); + return ImmutableListMultimap.of(); + } + + ImmutableMultimap.Builder builder = ImmutableMultimap.builder(); + + for (PointData pd : metricData.getData().getPoints()) { + ProjectName projectName = + metricDef.getSchema().extractProjectName(pd.getAttributes(), envInfo, clientInfo); + + TimeSeries timeSeries = + TimeSeries.newBuilder() + .setMetricKind(convertMetricKind(metricData)) + .setValueType(convertValueType(metricData.getType())) + .setResource( + metricDef + .getSchema() + .extractMonitoredResource(pd.getAttributes(), envInfo, clientInfo)) + .setMetric( + Metric.newBuilder() + .setType(metricDef.getExternalName()) + .putAllLabels( + metricDef.extractMetricLabels(pd.getAttributes(), envInfo, clientInfo))) + .addPoints(convertPointData(metricData.getType(), pd)) + .build(); + + builder.put(projectName, timeSeries); + } + return builder.build(); + } + + private Point convertPointData(MetricDataType type, PointData pointData) { + TimeInterval timeInterval = + TimeInterval.newBuilder() + .setStartTime(Timestamps.fromNanos(pointData.getStartEpochNanos())) + .setEndTime(Timestamps.fromNanos(pointData.getEpochNanos())) + .build(); + + Point.Builder builder = Point.newBuilder().setInterval(timeInterval); + switch (type) { + case HISTOGRAM: + case EXPONENTIAL_HISTOGRAM: + return builder + .setValue( + TypedValue.newBuilder() + .setDistributionValue(convertHistogramData((HistogramPointData) pointData)) + .build()) + .build(); + case DOUBLE_GAUGE: + case DOUBLE_SUM: + return builder + .setValue( + TypedValue.newBuilder() + .setDoubleValue(((DoublePointData) pointData).getValue()) + .build()) + .build(); + case LONG_GAUGE: + case LONG_SUM: + return builder + .setValue(TypedValue.newBuilder().setInt64Value(((LongPointData) pointData).getValue())) + .build(); + default: + LOGGER.log(Level.WARNING, "unsupported metric type %s", type); + return builder.build(); + } + } + + private static Distribution convertHistogramData(HistogramPointData pointData) { + return Distribution.newBuilder() + .setCount(pointData.getCount()) + .setMean(pointData.getCount() == 0L ? 0.0D : pointData.getSum() / pointData.getCount()) + .setBucketOptions( + BucketOptions.newBuilder() + .setExplicitBuckets(Explicit.newBuilder().addAllBounds(pointData.getBoundaries()))) + .addAllBucketCounts(pointData.getCounts()) + .build(); + } + + private static MetricKind convertMetricKind(MetricData metricData) { + switch (metricData.getType()) { + case HISTOGRAM: + case EXPONENTIAL_HISTOGRAM: + return convertHistogramType(metricData.getHistogramData()); + case LONG_GAUGE: + case DOUBLE_GAUGE: + return GAUGE; + case LONG_SUM: + return convertSumDataType(metricData.getLongSumData()); + case DOUBLE_SUM: + return convertSumDataType(metricData.getDoubleSumData()); + default: + return UNRECOGNIZED; + } + } + + private static MetricKind convertHistogramType(HistogramData histogramData) { + if (histogramData.getAggregationTemporality() == AggregationTemporality.CUMULATIVE) { + return CUMULATIVE; + } + return UNRECOGNIZED; + } + + private static MetricKind convertSumDataType(SumData sum) { + if (!sum.isMonotonic()) { + return GAUGE; + } + if (sum.getAggregationTemporality() == AggregationTemporality.CUMULATIVE) { + return CUMULATIVE; + } + return UNRECOGNIZED; + } + + private static ValueType convertValueType(MetricDataType metricDataType) { + switch (metricDataType) { + case LONG_GAUGE: + case LONG_SUM: + return INT64; + case DOUBLE_GAUGE: + case DOUBLE_SUM: + return DOUBLE; + case HISTOGRAM: + case EXPONENTIAL_HISTOGRAM: + return DISTRIBUTION; + default: + return ValueType.UNRECOGNIZED; + } + } +} diff --git a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java index 285206e94937..b68911de4b69 100644 --- a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java +++ b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java @@ -35,6 +35,10 @@ import com.google.api.core.ApiFuture; import com.google.api.core.ApiFutures; import com.google.api.gax.rpc.UnaryCallable; +import com.google.bigtable.v2.InstanceName; +import com.google.cloud.bigtable.data.v2.internal.csm.MetricRegistry; +import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo; +import com.google.cloud.bigtable.data.v2.internal.csm.attributes.EnvInfo; import com.google.cloud.monitoring.v3.MetricServiceClient; import com.google.cloud.monitoring.v3.stub.MetricServiceStub; import com.google.common.base.Suppliers; @@ -90,15 +94,30 @@ public class BigtableCloudMonitoringExporterTest { private Resource resource; private InstrumentationScopeInfo scope; + private EnvInfo envInfo = EnvInfo.builder() + .setProject("client-project") + .setPlatform("gce_instance") + .setRegion("cleint-region") + .setHostName("harold") + .setHostId("1234567890") + .setUid(taskId) + .build(); + private ClientInfo clientInfo = ClientInfo.builder() + .setInstanceName(InstanceName.of(projectId, instanceId)) + .setAppProfileId(appProfileId) + .setClientName(clientName) + .build(); + @Before public void setUp() { fakeMetricServiceClient = new FakeMetricServiceClient(mockMetricServiceStub); exporter = new BigtableCloudMonitoringExporter( - fakeMetricServiceClient, - ImmutableList.of( - new BigtableCloudMonitoringExporter.PublicTimeSeriesConverter(taskId))); + new MetricRegistry(), + Suppliers.ofInstance(envInfo), + clientInfo, + fakeMetricServiceClient); attributes = Attributes.builder() @@ -307,26 +326,12 @@ public void testExportingSumDataInBatches() { @Test public void testTimeSeriesForMetricWithGceOrGkeResource() { - String gceProjectId = "fake-gce-project"; BigtableCloudMonitoringExporter exporter = new BigtableCloudMonitoringExporter( - fakeMetricServiceClient, - ImmutableList.of( - new BigtableCloudMonitoringExporter.InternalTimeSeriesConverter( - Suppliers.ofInstance( - MonitoredResource.newBuilder() - .setType("bigtable_client") - .putLabels("project_id", gceProjectId) - .putLabels("instance", "resource-instance") - .putLabels("app_profile", "resource-app-profile") - .putLabels("client_project", "client-project") - .putLabels("region", "cleint-region") - .putLabels("cloud_platform", "gce_instance") - .putLabels("host_id", "1234567890") - .putLabels("host_name", "harold") - .putLabels("client_name", "java/1234") - .putLabels("uuid", "something") - .build())))); + new MetricRegistry(), + Suppliers.ofInstance(envInfo), + clientInfo, + fakeMetricServiceClient); ArgumentCaptor argumentCaptor = ArgumentCaptor.forClass(CreateTimeSeriesRequest.class); @@ -372,7 +377,7 @@ public void testTimeSeriesForMetricWithGceOrGkeResource() { CreateTimeSeriesRequest request = argumentCaptor.getValue(); - assertThat(request.getName()).isEqualTo("projects/" + gceProjectId); + assertThat(request.getName()).isEqualTo("projects/" + projectId); assertThat(request.getTimeSeriesList()).hasSize(1); com.google.monitoring.v3.TimeSeries timeSeries = request.getTimeSeriesList().get(0); @@ -380,16 +385,16 @@ public void testTimeSeriesForMetricWithGceOrGkeResource() { assertThat(timeSeries.getResource().getLabelsMap()) .isEqualTo( ImmutableMap.builder() - .put("project_id", gceProjectId) - .put("instance", "resource-instance") - .put("app_profile", "resource-app-profile") + .put("project_id", projectId) + .put("instance", instanceId) + .put("app_profile", appProfileId) .put("client_project", "client-project") .put("region", "cleint-region") .put("cloud_platform", "gce_instance") .put("host_id", "1234567890") .put("host_name", "harold") - .put("client_name", "java/1234") - .put("uuid", "something") + .put("client_name", clientName) + .put("uuid", taskId) .build()); assertThat(timeSeries.getMetric().getLabelsMap()) From 7eb0ba88e97c3acf4023aa5c7c67984e79e2e654 Mon Sep 17 00:00:00 2001 From: Igor Bernstein Date: Thu, 26 Feb 2026 09:31:57 -0500 Subject: [PATCH 2/4] re-add missing endpoint prop Change-Id: Ifcb6af1b9cb041ce79bc38a67817ae94904de83b --- .../BigtableCloudMonitoringExporter.java | 23 +++++++++++++------ 1 file changed, 16 insertions(+), 7 deletions(-) diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java index eec69d47e58c..5c301717191c 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java @@ -60,13 +60,21 @@ public class BigtableCloudMonitoringExporter implements MetricExporter { private static final Logger LOGGER = Logger.getLogger(BigtableCloudMonitoringExporter.class.getName()); + // This system property can be used to override the monitoring endpoint + // to a different environment. It's meant for internal testing only and + // will be removed in future versions. Use settings in EnhancedBigtableStubSettings + // to override the endpoint. + @Deprecated @Nullable + private static final String MONITORING_ENDPOINT_OVERRIDE_SYS_PROP = + System.getProperty("bigtable.test-monitoring-endpoint"); + // This the quota limit from Cloud Monitoring. More details in // https://cloud.google.com/monitoring/quotas#custom_metrics_quotas. private static final int EXPORT_BATCH_SIZE_LIMIT = 200; private final Supplier envInfo; private final ClientInfo clientInfo; - private MetricRegistry metricRegistry; + private final MetricRegistry metricRegistry; private final MetricServiceClient client; private final AtomicReference state; @@ -98,6 +106,13 @@ public static BigtableCloudMonitoringExporter create( .map(FixedCredentialsProvider::create) .orElse(NoCredentialsProvider.create())); + if (MONITORING_ENDPOINT_OVERRIDE_SYS_PROP != null) { + LOGGER.warning( + "Setting the monitoring endpoint through system variable will be removed in future" + + " versions"); + settingsBuilder.setEndpoint(MONITORING_ENDPOINT_OVERRIDE_SYS_PROP); + } + if (endpoint != null) { settingsBuilder.setEndpoint(endpoint); } @@ -131,12 +146,6 @@ public void close() { public CompletableResultCode export(Collection metricData) { Preconditions.checkState(state.get() != State.Closed, "Exporter is closed"); - if (metricRegistry == null) { - String msg = "Bigtable exporter tried to export before fully configured"; - LOGGER.warning(msg); - return CompletableResultCode.ofExceptionalFailure(new IllegalStateException(msg)); - } - lastExportCode = doExport(metricData); return lastExportCode; } From ca87acf1bf252faebc8f4ae6c25430cfe23b465d Mon Sep 17 00:00:00 2001 From: cloud-java-bot Date: Thu, 26 Feb 2026 14:37:33 +0000 Subject: [PATCH 3/4] chore: generate libraries at Thu Feb 26 14:34:57 UTC 2026 --- .../data/v2/internal/csm/MetricsImpl.java | 7 ++++- .../data/v2/stub/metrics/Converter.java | 3 -- .../BigtableCloudMonitoringExporterTest.java | 29 ++++++++++--------- 3 files changed, 21 insertions(+), 18 deletions(-) diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java index df9b4b4eaa9b..c7bf85943166 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/internal/csm/MetricsImpl.java @@ -197,7 +197,12 @@ public static OpenTelemetrySdk createBuiltinOtel( MetricExporter publicExporter = BigtableCloudMonitoringExporter.create( - metricRegistry, EnvInfo::detect, clientInfo, credentials, metricsEndpoint, universeDomain); + metricRegistry, + EnvInfo::detect, + clientInfo, + credentials, + metricsEndpoint, + universeDomain); PeriodicMetricReaderBuilder readerBuilder = PeriodicMetricReader.builder(publicExporter).setExecutor(executor); meterProvider.registerMetricReader(readerBuilder.build()); diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java index 03b1f1519e27..4a2ca946f125 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/Converter.java @@ -32,11 +32,9 @@ import com.google.cloud.bigtable.data.v2.internal.csm.MetricRegistry; import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo; import com.google.cloud.bigtable.data.v2.internal.csm.attributes.EnvInfo; -import com.google.cloud.bigtable.data.v2.internal.csm.metrics.GrpcMetric; import com.google.cloud.bigtable.data.v2.internal.csm.metrics.MetricWrapper; import com.google.common.collect.ImmutableListMultimap; import com.google.common.collect.ImmutableMultimap; -import com.google.common.collect.ImmutableSet; import com.google.common.collect.Multimap; import com.google.monitoring.v3.Point; import com.google.monitoring.v3.ProjectName; @@ -55,7 +53,6 @@ import io.opentelemetry.sdk.metrics.data.SumData; import java.util.Collection; import java.util.Map; -import java.util.Set; import java.util.logging.Level; import java.util.logging.Logger; diff --git a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java index b68911de4b69..d9ccad187e46 100644 --- a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java +++ b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporterTest.java @@ -31,7 +31,6 @@ import static org.mockito.Mockito.when; import com.google.api.Distribution; -import com.google.api.MonitoredResource; import com.google.api.core.ApiFuture; import com.google.api.core.ApiFutures; import com.google.api.gax.rpc.UnaryCallable; @@ -94,19 +93,21 @@ public class BigtableCloudMonitoringExporterTest { private Resource resource; private InstrumentationScopeInfo scope; - private EnvInfo envInfo = EnvInfo.builder() - .setProject("client-project") - .setPlatform("gce_instance") - .setRegion("cleint-region") - .setHostName("harold") - .setHostId("1234567890") - .setUid(taskId) - .build(); - private ClientInfo clientInfo = ClientInfo.builder() - .setInstanceName(InstanceName.of(projectId, instanceId)) - .setAppProfileId(appProfileId) - .setClientName(clientName) - .build(); + private EnvInfo envInfo = + EnvInfo.builder() + .setProject("client-project") + .setPlatform("gce_instance") + .setRegion("cleint-region") + .setHostName("harold") + .setHostId("1234567890") + .setUid(taskId) + .build(); + private ClientInfo clientInfo = + ClientInfo.builder() + .setInstanceName(InstanceName.of(projectId, instanceId)) + .setAppProfileId(appProfileId) + .setClientName(clientName) + .build(); @Before public void setUp() { From 11d9e386638dd256b40a6f0c7d8af2ca6808b757 Mon Sep 17 00:00:00 2001 From: Igor Bernstein Date: Thu, 26 Feb 2026 14:50:44 -0500 Subject: [PATCH 4/4] dont log failures when closed Change-Id: Id5f40c3919b1f94fc28ec6d38f929709aefa8687 --- .../data/v2/stub/metrics/BigtableCloudMonitoringExporter.java | 2 +- 1 file changed, 1 insertion(+), 1 deletion(-) diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java index 5c301717191c..3bec1fc1e7bc 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableCloudMonitoringExporter.java @@ -191,7 +191,7 @@ public void onFailure(Throwable throwable) { RuntimeException asyncWrapper = new RuntimeException("export failed", throwable); asyncWrapper.setStackTrace(stackTrace); - if (state.get() != State.Closing) { + if (state.get() != State.Closing || state.get() != State.Closed) { // ignore the export warning when client is shutting down LOGGER.log(Level.WARNING, msg, asyncWrapper); }