From b70f83e49278e25e1439c44688aba98c1f91d068 Mon Sep 17 00:00:00 2001 From: Igor Bernstein Date: Mon, 23 Feb 2026 15:16:59 -0500 Subject: [PATCH] chore: decompose stubsettings into perOpSettings & BtClientContext This avoids any overlap between the 2 classes being passed to a stub and avoids setting inconsistencies Change-Id: I22beae49e1e4118930c7ed9636576a42d397f7ce --- .../bigtable/data/v2/BigtableDataClient.java | 12 -- .../data/v2/BigtableDataClientFactory.java | 46 +++----- .../data/v2/stub/ClientOperationSettings.java | 4 +- .../data/v2/stub/EnhancedBigtableStub.java | 108 +++++++++--------- .../v2/stub/EnhancedBigtableStubSettings.java | 5 + .../metrics/BigtableTracerCallableTest.java | 4 +- .../metrics/BuiltinMetricsTracerTest.java | 4 +- .../v2/stub/metrics/MetricsTracerTest.java | 4 +- 8 files changed, 82 insertions(+), 105 deletions(-) diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClient.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClient.java index 5af8e9dc9618..b659a021754c 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClient.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClient.java @@ -54,7 +54,6 @@ import com.google.cloud.bigtable.data.v2.models.sql.PreparedStatement; import com.google.cloud.bigtable.data.v2.models.sql.ResultSet; import com.google.cloud.bigtable.data.v2.models.sql.SqlType; -import com.google.cloud.bigtable.data.v2.stub.BigtableClientContext; import com.google.cloud.bigtable.data.v2.stub.EnhancedBigtableStub; import com.google.cloud.bigtable.data.v2.stub.sql.SqlServerStream; import com.google.common.util.concurrent.MoreExecutors; @@ -180,17 +179,6 @@ public static BigtableDataClient create(BigtableDataSettings settings) throws IO return new BigtableDataClient(stub); } - /** - * Constructs an instance of BigtableDataClient with the provided client context. This is used by - * {@link BigtableDataClientFactory} and the client context will not be closed unless {@link - * BigtableDataClientFactory#close()} is called. - */ - static BigtableDataClient createWithClientContext( - BigtableDataSettings settings, BigtableClientContext context) throws IOException { - EnhancedBigtableStub stub = new EnhancedBigtableStub(settings.getStubSettings(), context); - return new BigtableDataClient(stub); - } - @InternalApi("Visible for testing") BigtableDataClient(EnhancedBigtableStub stub) { this.stub = stub; diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClientFactory.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClientFactory.java index 544d75d6a78e..d73fbe2a12db 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClientFactory.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/BigtableDataClientFactory.java @@ -18,6 +18,8 @@ import com.google.api.core.BetaApi; import com.google.bigtable.v2.InstanceName; import com.google.cloud.bigtable.data.v2.stub.BigtableClientContext; +import com.google.cloud.bigtable.data.v2.stub.ClientOperationSettings; +import com.google.cloud.bigtable.data.v2.stub.EnhancedBigtableStub; import java.io.IOException; import javax.annotation.Nonnull; @@ -61,9 +63,8 @@ */ @BetaApi("This feature is currently experimental and can change in the future") public final class BigtableDataClientFactory implements AutoCloseable { - - private final BigtableDataSettings defaultSettings; private final BigtableClientContext sharedClientContext; + private final ClientOperationSettings perOpSettings; /** * Create a instance of this factory. @@ -75,13 +76,14 @@ public static BigtableDataClientFactory create(BigtableDataSettings defaultSetti throws IOException { BigtableClientContext sharedClientContext = BigtableClientContext.create(defaultSettings.getStubSettings()); - return new BigtableDataClientFactory(sharedClientContext, defaultSettings); + ClientOperationSettings perOpSettings = defaultSettings.getStubSettings().getPerOpSettings(); + return new BigtableDataClientFactory(sharedClientContext, perOpSettings); } private BigtableDataClientFactory( - BigtableClientContext sharedClientContext, BigtableDataSettings defaultSettings) { + BigtableClientContext sharedClientContext, ClientOperationSettings perOpSettings) { this.sharedClientContext = sharedClientContext; - this.defaultSettings = defaultSettings; + this.perOpSettings = perOpSettings; } /** @@ -109,7 +111,7 @@ public BigtableDataClient createDefault() { sharedClientContext.createChild( sharedClientContext.getInstanceName(), sharedClientContext.getAppProfileId()); - return BigtableDataClient.createWithClientContext(defaultSettings, ctx); + return new BigtableDataClient(new EnhancedBigtableStub(perOpSettings, ctx)); } catch (IOException e) { // Should never happen because the connection has been established already throw new RuntimeException( @@ -127,14 +129,10 @@ public BigtableDataClient createDefault() { * release all resources, first close all of the created clients and then this factory instance. */ public BigtableDataClient createForAppProfile(@Nonnull String appProfileId) throws IOException { - BigtableDataSettings settings = - defaultSettings.toBuilder().setAppProfileId(appProfileId).build(); BigtableClientContext ctx = - sharedClientContext.createChild( - InstanceName.of(settings.getProjectId(), settings.getInstanceId()), - settings.getAppProfileId()); + sharedClientContext.createChild(sharedClientContext.getInstanceName(), appProfileId); - return BigtableDataClient.createWithClientContext(settings, ctx); + return new BigtableDataClient(new EnhancedBigtableStub(perOpSettings, ctx)); } /** @@ -148,18 +146,10 @@ public BigtableDataClient createForAppProfile(@Nonnull String appProfileId) thro */ public BigtableDataClient createForInstance(@Nonnull String projectId, @Nonnull String instanceId) throws IOException { - BigtableDataSettings settings = - defaultSettings.toBuilder() - .setProjectId(projectId) - .setInstanceId(instanceId) - .setDefaultAppProfileId() - .build(); BigtableClientContext ctx = - sharedClientContext.createChild( - InstanceName.of(settings.getProjectId(), settings.getInstanceId()), - settings.getAppProfileId()); + sharedClientContext.createChild(InstanceName.of(projectId, instanceId), ""); - return BigtableDataClient.createWithClientContext(settings, ctx); + return new BigtableDataClient(new EnhancedBigtableStub(perOpSettings, ctx)); } /** @@ -174,17 +164,9 @@ public BigtableDataClient createForInstance(@Nonnull String projectId, @Nonnull public BigtableDataClient createForInstance( @Nonnull String projectId, @Nonnull String instanceId, @Nonnull String appProfileId) throws IOException { - BigtableDataSettings settings = - defaultSettings.toBuilder() - .setProjectId(projectId) - .setInstanceId(instanceId) - .setAppProfileId(appProfileId) - .build(); BigtableClientContext ctx = - sharedClientContext.createChild( - InstanceName.of(settings.getProjectId(), settings.getInstanceId()), - settings.getAppProfileId()); + sharedClientContext.createChild(InstanceName.of(projectId, instanceId), appProfileId); - return BigtableDataClient.createWithClientContext(settings, ctx); + return new BigtableDataClient(new EnhancedBigtableStub(perOpSettings, ctx)); } } diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/ClientOperationSettings.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/ClientOperationSettings.java index 8252b4b22a8b..540eb08cc8bd 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/ClientOperationSettings.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/ClientOperationSettings.java @@ -15,6 +15,7 @@ */ package com.google.cloud.bigtable.data.v2.stub; +import com.google.api.core.InternalApi; import com.google.api.gax.batching.BatchingSettings; import com.google.api.gax.batching.FlowControlSettings; import com.google.api.gax.batching.FlowController; @@ -45,7 +46,8 @@ import java.util.Set; import org.threeten.bp.Duration; -class ClientOperationSettings { +@InternalApi +public class ClientOperationSettings { private static final Set IDEMPOTENT_RETRY_CODES = ImmutableSet.of(StatusCode.Code.DEADLINE_EXCEEDED, StatusCode.Code.UNAVAILABLE); diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java index bf94964434af..6f0ffdc60f8f 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStub.java @@ -141,7 +141,7 @@ public class EnhancedBigtableStub implements AutoCloseable { private static final String CLIENT_NAME = "Bigtable"; private static final long FLOW_CONTROL_ADJUSTING_INTERVAL_MS = TimeUnit.SECONDS.toMillis(20); - private final EnhancedBigtableStubSettings settings; + private final ClientOperationSettings perOpSettings; private final BigtableClientContext bigtableClientContext; private final RequestContext requestContext; @@ -175,12 +175,12 @@ public class EnhancedBigtableStub implements AutoCloseable { public static EnhancedBigtableStub create(EnhancedBigtableStubSettings settings) throws IOException { BigtableClientContext bigtableClientContext = BigtableClientContext.create(settings); - return new EnhancedBigtableStub(settings, bigtableClientContext); + return new EnhancedBigtableStub(settings.getPerOpSettings(), bigtableClientContext); } public EnhancedBigtableStub( - EnhancedBigtableStubSettings settings, BigtableClientContext clientContext) { - this.settings = settings; + ClientOperationSettings perOpSettings, BigtableClientContext clientContext) { + this.perOpSettings = perOpSettings; this.bigtableClientContext = clientContext; this.requestContext = RequestContext.create( @@ -188,7 +188,7 @@ public EnhancedBigtableStub( clientContext.getInstanceName().getInstance(), clientContext.getAppProfileId()); this.bulkMutationFlowController = - new FlowController(settings.bulkMutateRowsSettings().getDynamicFlowControlSettings()); + new FlowController(perOpSettings.bulkMutateRowsSettings.getDynamicFlowControlSettings()); this.bulkMutationDynamicFlowControlStats = new DynamicFlowControlStats(); readRowsCallable = createReadRowsCallable(new DefaultRowAdapter()); @@ -230,7 +230,7 @@ public EnhancedBigtableStub( @BetaApi("This surface is stable yet it might be removed in the future.") public ServerStreamingCallable createReadRowsRawCallable( RowAdapter rowAdapter) { - return createReadRowsBaseCallable(settings.readRowsSettings(), rowAdapter) + return createReadRowsBaseCallable(perOpSettings.readRowsSettings, rowAdapter) .withDefaultCallContext(bigtableClientContext.getClientContext().getDefaultCallContext()); } @@ -251,7 +251,7 @@ public ServerStreamingCallable createReadRowsRawCa public ServerStreamingCallable createReadRowsCallable( RowAdapter rowAdapter) { ServerStreamingCallable readRowsCallable = - createReadRowsBaseCallable(settings.readRowsSettings(), rowAdapter); + createReadRowsBaseCallable(perOpSettings.readRowsSettings, rowAdapter); ServerStreamingCallable readRowsUserCallable = new ReadRowsUserCallable<>(readRowsCallable, requestContext); @@ -267,7 +267,7 @@ public ServerStreamingCallable createReadRowsCallable( bigtableClientContext .getClientContext() .getDefaultCallContext() - .withRetrySettings(settings.readRowsSettings().getRetrySettings())); + .withRetrySettings(perOpSettings.readRowsSettings.getRetrySettings())); } /** @@ -290,8 +290,8 @@ public UnaryCallable createReadRowCallable(RowAdapter ServerStreamingCallable readRowsCallable = createReadRowsBaseCallable( ServerStreamingCallSettings.newBuilder() - .setRetryableCodes(settings.readRowSettings().getRetryableCodes()) - .setRetrySettings(settings.readRowSettings().getRetrySettings()) + .setRetryableCodes(perOpSettings.readRowSettings.getRetryableCodes()) + .setRetrySettings(perOpSettings.readRowSettings.getRetrySettings()) .setIdleTimeoutDuration(Duration.ZERO) .setWaitTimeoutDuration(Duration.ZERO) .build(), @@ -307,7 +307,7 @@ public UnaryCallable createReadRowCallable(RowAdapter readRowCallable, clientContext .getDefaultCallContext() - .withRetrySettings(settings.readRowSettings().getRetrySettings()), + .withRetrySettings(perOpSettings.readRowSettings.getRetrySettings()), clientContext.getTracerFactory(), getSpanName("ReadRow"), /* allowNoResponses= */ true); @@ -410,7 +410,7 @@ public ServerStreamingCallable createSkipLargeRowsCall RowAdapter rowAdapter) { ServerStreamingCallSettings readRowsSettings = - (ServerStreamingCallSettings) settings.readRowsSettings(); + (ServerStreamingCallSettings) perOpSettings.readRowsSettings; ServerStreamingCallable base = GrpcRawCallableFactory.createServerStreamingCallable( @@ -501,7 +501,7 @@ public ServerStreamingCallable createSkipLargeRowsCall private UnaryCallable> createBulkReadRowsCallable( RowAdapter rowAdapter) { ServerStreamingCallable readRowsCallable = - createReadRowsBaseCallable(settings.readRowsSettings(), rowAdapter); + createReadRowsBaseCallable(perOpSettings.readRowsSettings, rowAdapter); ServerStreamingCallable readRowsUserCallable = new ReadRowsUserCallable<>(readRowsCallable, requestContext); @@ -521,7 +521,7 @@ private UnaryCallable> createBulkReadRowsCallable( bigtableClientContext .getClientContext() .getDefaultCallContext() - .withRetrySettings(settings.readRowsSettings().getRetrySettings())); + .withRetrySettings(perOpSettings.readRowsSettings.getRetrySettings())); } /** @@ -571,7 +571,7 @@ public ApiFuture> futureCall(String s, ApiCallContext apiCallCon composeRequestParams( r.getAppProfileId(), r.getTableName(), r.getAuthorizedViewName())) .build(), - settings.sampleRowKeysSettings().getRetryableCodes()); + perOpSettings.sampleRowKeysSettings.getRetryableCodes()); UnaryCallable> spoolable = base.all(); @@ -583,7 +583,7 @@ public ApiFuture> futureCall(String s, ApiCallContext apiCallCon withAttemptTracer = new BigtableTracerUnaryCallable<>(withStatsHeaders); UnaryCallable> - retryable = withRetries(withAttemptTracer, settings.sampleRowKeysSettings()); + retryable = withRetries(withAttemptTracer, perOpSettings.sampleRowKeysSettings); return createUserFacingUnaryCallable( methodName, @@ -592,7 +592,7 @@ public ApiFuture> futureCall(String s, ApiCallContext apiCallCon bigtableClientContext .getClientContext() .getDefaultCallContext() - .withRetrySettings(settings.sampleRowKeysSettings().getRetrySettings()))); + .withRetrySettings(perOpSettings.sampleRowKeysSettings.getRetrySettings()))); } /** @@ -609,7 +609,7 @@ private UnaryCallable createMutateRowCallable() { req -> composeRequestParams( req.getAppProfileId(), req.getTableName(), req.getAuthorizedViewName()), - settings.mutateRowSettings(), + perOpSettings.mutateRowSettings, req -> req.toProto(requestContext), resp -> null); } @@ -643,12 +643,12 @@ private UnaryCallable createMutateRowsBas composeRequestParams( r.getAppProfileId(), r.getTableName(), r.getAuthorizedViewName())) .build(), - settings.bulkMutateRowsSettings().getRetryableCodes()); + perOpSettings.bulkMutateRowsSettings.getRetryableCodes()); ServerStreamingCallable callable = new StatsHeadersServerStreamingCallable<>(base); - if (settings.bulkMutateRowsSettings().isServerInitiatedFlowControlEnabled()) { + if (perOpSettings.bulkMutateRowsSettings.isServerInitiatedFlowControlEnabled()) { callable = new RateLimitingServerStreamingCallable(callable); } @@ -672,7 +672,7 @@ private UnaryCallable createMutateRowsBas new RetryAlgorithm<>( mutateRowsPartialErrorRetryAlgorithm, new ExponentialRetryAlgorithm( - settings.bulkMutateRowsSettings().getRetrySettings(), clientContext.getClock())); + perOpSettings.bulkMutateRowsSettings.getRetrySettings(), clientContext.getClock())); RetryingExecutorWithContext retryingExecutor = new ScheduledRetryingExecutor<>(retryAlgorithm, clientContext.getExecutor()); @@ -681,20 +681,20 @@ private UnaryCallable createMutateRowsBas clientContext.getDefaultCallContext(), withAttemptTracer, retryingExecutor, - settings.bulkMutateRowsSettings().getRetryableCodes(), + perOpSettings.bulkMutateRowsSettings.getRetryableCodes(), retryAlgorithm); UnaryCallable withCookie = new CookiesUnaryCallable<>(baseCallable); UnaryCallable flowControlCallable = null; - if (settings.bulkMutateRowsSettings().isLatencyBasedThrottlingEnabled()) { + if (perOpSettings.bulkMutateRowsSettings.isLatencyBasedThrottlingEnabled()) { flowControlCallable = new DynamicFlowControlCallable( withCookie, bulkMutationFlowController, bulkMutationDynamicFlowControlStats, - settings.bulkMutateRowsSettings().getTargetRpcLatencyMs(), + perOpSettings.bulkMutateRowsSettings.getTargetRpcLatencyMs(), FLOW_CONTROL_ADJUSTING_INTERVAL_MS); } UnaryCallable userFacing = @@ -713,7 +713,7 @@ private UnaryCallable createMutateRowsBas return traced.withDefaultCallContext( clientContext .getDefaultCallContext() - .withRetrySettings(settings.bulkMutateRowsSettings().getRetrySettings())); + .withRetrySettings(perOpSettings.bulkMutateRowsSettings.getRetrySettings())); } /** @@ -738,10 +738,10 @@ private UnaryCallable createMutateRowsBas public Batcher newMutateRowsBatcher( @Nonnull String tableId, @Nullable GrpcCallContext ctx) { return new BatcherImpl<>( - settings.bulkMutateRowsSettings().getBatchingDescriptor(), + perOpSettings.bulkMutateRowsSettings.getBatchingDescriptor(), bulkMutateRowsCallable, BulkMutation.create(tableId), - settings.bulkMutateRowsSettings().getBatchingSettings(), + perOpSettings.bulkMutateRowsSettings.getBatchingSettings(), bigtableClientContext.getClientContext().getExecutor(), bulkMutationFlowController, MoreObjects.firstNonNull( @@ -770,10 +770,10 @@ public Batcher newMutateRowsBatcher( public Batcher newMutateRowsBatcher( TargetId targetId, @Nullable GrpcCallContext ctx) { return new BatcherImpl<>( - settings.bulkMutateRowsSettings().getBatchingDescriptor(), + perOpSettings.bulkMutateRowsSettings.getBatchingDescriptor(), bulkMutateRowsCallable, BulkMutation.create(targetId), - settings.bulkMutateRowsSettings().getBatchingSettings(), + perOpSettings.bulkMutateRowsSettings.getBatchingSettings(), bigtableClientContext.getClientContext().getExecutor(), bulkMutationFlowController, MoreObjects.firstNonNull( @@ -799,10 +799,10 @@ public Batcher newBulkReadRowsBatcher( @Nonnull Query query, @Nullable GrpcCallContext ctx) { Preconditions.checkNotNull(query, "query cannot be null"); return new BatcherImpl<>( - settings.bulkReadRowsSettings().getBatchingDescriptor(), + perOpSettings.bulkReadRowsSettings.getBatchingDescriptor(), bulkReadRowsCallable, query, - settings.bulkReadRowsSettings().getBatchingSettings(), + perOpSettings.bulkReadRowsSettings.getBatchingSettings(), bigtableClientContext.getClientContext().getExecutor(), null, MoreObjects.firstNonNull( @@ -824,7 +824,7 @@ private UnaryCallable createCheckAndMutateRowCa req -> composeRequestParams( req.getAppProfileId(), req.getTableName(), req.getAuthorizedViewName()), - settings.checkAndMutateRowSettings(), + perOpSettings.checkAndMutateRowSettings, req -> req.toProto(requestContext), CheckAndMutateRowResponse::getPredicateMatched); } @@ -847,7 +847,7 @@ private UnaryCallable createReadModifyWriteRowCallable( req -> composeRequestParams( req.getAppProfileId(), req.getTableName(), req.getAuthorizedViewName()), - settings.readModifyWriteRowSettings(), + perOpSettings.readModifyWriteRowSettings, req -> req.toProto(requestContext), resp -> rowAdapter.createRowFromProto(resp.getRow())); } @@ -881,7 +881,7 @@ private UnaryCallable createReadModifyWriteRowCallable( .setParamsExtractor( r -> composeRequestParams(r.getAppProfileId(), r.getTableName(), "")) .build(), - settings.generateInitialChangeStreamPartitionsSettings().getRetryableCodes()); + perOpSettings.generateInitialChangeStreamPartitionsSettings.getRetryableCodes()); ServerStreamingCallable userCallable = new GenerateInitialChangeStreamPartitionsUserCallable(base, requestContext); @@ -900,13 +900,13 @@ private UnaryCallable createReadModifyWriteRowCallable( ServerStreamingCallSettings innerSettings = ServerStreamingCallSettings.newBuilder() .setRetryableCodes( - settings.generateInitialChangeStreamPartitionsSettings().getRetryableCodes()) + perOpSettings.generateInitialChangeStreamPartitionsSettings.getRetryableCodes()) .setRetrySettings( - settings.generateInitialChangeStreamPartitionsSettings().getRetrySettings()) + perOpSettings.generateInitialChangeStreamPartitionsSettings.getRetrySettings()) .setIdleTimeout( - settings.generateInitialChangeStreamPartitionsSettings().getIdleTimeout()) + perOpSettings.generateInitialChangeStreamPartitionsSettings.getIdleTimeout()) .setWaitTimeout( - settings.generateInitialChangeStreamPartitionsSettings().getWaitTimeout()) + perOpSettings.generateInitialChangeStreamPartitionsSettings.getWaitTimeout()) .build(); ServerStreamingCallable watched = @@ -926,7 +926,7 @@ private UnaryCallable createReadModifyWriteRowCallable( clientContext .getDefaultCallContext() .withRetrySettings( - settings.generateInitialChangeStreamPartitionsSettings().getRetrySettings())); + perOpSettings.generateInitialChangeStreamPartitionsSettings.getRetrySettings())); } /** @@ -955,7 +955,7 @@ private UnaryCallable createReadModifyWriteRowCallable( .setParamsExtractor( r -> composeRequestParams(r.getAppProfileId(), r.getTableName(), "")) .build(), - settings.readChangeStreamSettings().getRetryableCodes()); + perOpSettings.readChangeStreamSettings.getRetryableCodes()); ServerStreamingCallable withStatsHeaders = new StatsHeadersServerStreamingCallable<>(base); @@ -975,10 +975,10 @@ private UnaryCallable createReadModifyWriteRowCallable( ServerStreamingCallSettings.newBuilder() .setResumptionStrategy( new ReadChangeStreamResumptionStrategy<>(changeStreamRecordAdapter)) - .setRetryableCodes(settings.readChangeStreamSettings().getRetryableCodes()) - .setRetrySettings(settings.readChangeStreamSettings().getRetrySettings()) - .setIdleTimeout(settings.readChangeStreamSettings().getIdleTimeout()) - .setWaitTimeout(settings.readChangeStreamSettings().getWaitTimeout()) + .setRetryableCodes(perOpSettings.readChangeStreamSettings.getRetryableCodes()) + .setRetrySettings(perOpSettings.readChangeStreamSettings.getRetrySettings()) + .setIdleTimeout(perOpSettings.readChangeStreamSettings.getIdleTimeout()) + .setWaitTimeout(perOpSettings.readChangeStreamSettings.getWaitTimeout()) .build(); ServerStreamingCallable watched = @@ -1002,7 +1002,7 @@ private UnaryCallable createReadModifyWriteRowCallable( return traced.withDefaultCallContext( clientContext .getDefaultCallContext() - .withRetrySettings(settings.readChangeStreamSettings().getRetrySettings())); + .withRetrySettings(perOpSettings.readChangeStreamSettings.getRetrySettings())); } /** @@ -1039,7 +1039,7 @@ public Map extract(ExecuteQueryRequest executeQueryRequest) { } }) .build(), - settings.executeQuerySettings().getRetryableCodes()); + perOpSettings.executeQuerySettings.getRetryableCodes()); ServerStreamingCallable withStatsHeaders = new StatsHeadersServerStreamingCallable<>(base); @@ -1060,10 +1060,10 @@ public Map extract(ExecuteQueryRequest executeQueryRequest) { ServerStreamingCallSettings retrySettings = ServerStreamingCallSettings.newBuilder() .setResumptionStrategy(new ExecuteQueryResumptionStrategy()) - .setRetryableCodes(settings.executeQuerySettings().getRetryableCodes()) - .setRetrySettings(settings.executeQuerySettings().getRetrySettings()) - .setIdleTimeout(settings.executeQuerySettings().getIdleTimeout()) - .setWaitTimeout(settings.executeQuerySettings().getWaitTimeout()) + .setRetryableCodes(perOpSettings.executeQuerySettings.getRetryableCodes()) + .setRetrySettings(perOpSettings.executeQuerySettings.getRetrySettings()) + .setIdleTimeout(perOpSettings.executeQuerySettings.getIdleTimeout()) + .setWaitTimeout(perOpSettings.executeQuerySettings.getWaitTimeout()) .build(); // Retries need to happen before row merging, because the resumeToken is part @@ -1078,8 +1078,8 @@ public Map extract(ExecuteQueryRequest executeQueryRequest) { ServerStreamingCallSettings watchdogSettings = ServerStreamingCallSettings.newBuilder() - .setIdleTimeout(settings.executeQuerySettings().getIdleTimeout()) - .setWaitTimeout(settings.executeQuerySettings().getWaitTimeout()) + .setIdleTimeout(perOpSettings.executeQuerySettings.getIdleTimeout()) + .setWaitTimeout(perOpSettings.executeQuerySettings.getWaitTimeout()) .build(); // Watchdog needs to stay above the metadata error handling so that watchdog errors @@ -1099,14 +1099,14 @@ public Map extract(ExecuteQueryRequest executeQueryRequest) { traced.withDefaultCallContext( clientContext .getDefaultCallContext() - .withRetrySettings(settings.executeQuerySettings().getRetrySettings()))); + .withRetrySettings(perOpSettings.executeQuerySettings.getRetrySettings()))); } private UnaryCallable createPrepareQueryCallable() { return createUnaryCallable( BigtableGrpc.getPrepareQueryMethod(), req -> composeInstanceLevelRequestParams(req.getInstanceName(), req.getAppProfileId()), - settings.prepareQuerySettings(), + perOpSettings.prepareQuerySettings, req -> req.toProto(requestContext), PrepareResponse::fromProto); } diff --git a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStubSettings.java b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStubSettings.java index 0ce0c7b299ea..1a416d51e4e2 100644 --- a/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStubSettings.java +++ b/google-cloud-bigtable/src/main/java/com/google/cloud/bigtable/data/v2/stub/EnhancedBigtableStubSettings.java @@ -237,6 +237,11 @@ public boolean areInternalMetricsEnabled() { return areInternalMetricsEnabled; } + @InternalApi + public ClientOperationSettings getPerOpSettings() { + return perOpSettings; + } + /** Returns a builder for the default ChannelProvider for this service. */ public static InstantiatingGrpcChannelProvider.Builder defaultGrpcTransportProviderBuilder() { InstantiatingGrpcChannelProvider.Builder grpcTransportProviderBuilder = diff --git a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableTracerCallableTest.java b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableTracerCallableTest.java index 0f84417a7007..f9b0e56ac527 100644 --- a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableTracerCallableTest.java +++ b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BigtableTracerCallableTest.java @@ -132,7 +132,7 @@ public void sendHeaders(Metadata headers) { attempts = settings.getStubSettings().readRowsSettings().getRetrySettings().getMaxAttempts(); stub = new EnhancedBigtableStub( - settings.getStubSettings(), + settings.getStubSettings().getPerOpSettings(), BigtableClientContext.create( settings.getStubSettings(), Tags.getTagger(), localStats.getStatsRecorder())); @@ -151,7 +151,7 @@ public void sendHeaders(Metadata headers) { noHeaderStub = new EnhancedBigtableStub( - noHeaderSettings.getStubSettings(), + noHeaderSettings.getStubSettings().getPerOpSettings(), BigtableClientContext.create( noHeaderSettings.getStubSettings(), Tags.getTagger(), diff --git a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BuiltinMetricsTracerTest.java b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BuiltinMetricsTracerTest.java index 3d0e6425d937..1ffccab7dd6a 100644 --- a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BuiltinMetricsTracerTest.java +++ b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/BuiltinMetricsTracerTest.java @@ -70,7 +70,6 @@ import com.google.cloud.bigtable.data.v2.models.RowMutation; import com.google.cloud.bigtable.data.v2.models.RowMutationEntry; import com.google.cloud.bigtable.data.v2.models.TableId; -import com.google.cloud.bigtable.data.v2.stub.BigtableClientContext; import com.google.cloud.bigtable.data.v2.stub.EnhancedBigtableStub; import com.google.cloud.bigtable.data.v2.stub.EnhancedBigtableStubSettings; import com.google.common.base.Stopwatch; @@ -287,8 +286,7 @@ public void sendHeaders(Metadata headers) { return builder.proxyDetector(delayProxyDetector).intercept(outstandingRpcCounter); }); stubSettingsBuilder.setTransportChannelProvider(channelProvider.build()); - EnhancedBigtableStubSettings stubSettings = stubSettingsBuilder.build(); - stub = new EnhancedBigtableStub(stubSettings, BigtableClientContext.create(stubSettings)); + stub = EnhancedBigtableStub.create(stubSettingsBuilder.build()); } @After diff --git a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/MetricsTracerTest.java b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/MetricsTracerTest.java index ab2fe6e205ef..da864bf49593 100644 --- a/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/MetricsTracerTest.java +++ b/google-cloud-bigtable/src/test/java/com/google/cloud/bigtable/data/v2/stub/metrics/MetricsTracerTest.java @@ -127,7 +127,9 @@ public void setUp() throws Exception { BigtableClientContext.create( settings.getStubSettings(), Tags.getTagger(), localStats.getStatsRecorder()); - stub = new EnhancedBigtableStub(settings.getStubSettings(), bigtableClientContext); + stub = + new EnhancedBigtableStub( + settings.getStubSettings().getPerOpSettings(), bigtableClientContext); } @After