Skip to content
Merged
Show file tree
Hide file tree
Changes from 1 commit
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
Original file line number Diff line number Diff line change
Expand Up @@ -74,12 +74,8 @@ public final class BigtableDataClientFactory implements AutoCloseable {
*/
public static BigtableDataClientFactory create(BigtableDataSettings defaultSettings)
throws IOException {
BigtableDataSettings.Builder builder = defaultSettings.toBuilder();
builder.stubSettings().setSessionsEnabled(false);
defaultSettings = builder.build();

BigtableClientContext sharedClientContext =
BigtableClientContext.create(defaultSettings.getStubSettings());
BigtableClientContext.createForFactory(defaultSettings.getStubSettings());
ClientOperationSettings perOpSettings = defaultSettings.getStubSettings().getPerOpSettings();
return new BigtableDataClientFactory(sharedClientContext, perOpSettings);
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -73,7 +73,7 @@ public class Client implements AutoCloseable {
private final Resource<ScheduledExecutorService> backgroundExecutor;

private final CallOptions defaultCallOptions;
private final ChannelPool channelPool;
private final Resource<ChannelPool> channelPool;
private final Resource<Metrics> metrics;
private final Resource<ClientConfigurationManager> configManager;

Expand Down Expand Up @@ -147,19 +147,23 @@ public static Client create(ClientSettings settings) throws IOException {
return new Client(
featureFlags,
clientInfo,
settings.getChannelProvider(),
Resource.createOwned(metrics, metrics::close),
Resource.createOwned(configManager, configManager::close),
Resource.createOwned(backgroundExecutor, backgroundExecutor::shutdown));
Resource.createOwned(backgroundExecutor, backgroundExecutor::shutdown),
settings.getChannelProvider());
}

/**
* Standard constructor used by non-factory clients. Builds an owned {@link SwitchingChannelPool}
* from the provided {@link ChannelProvider}.
*/
public Client(
FeatureFlags featureFlags,
ClientInfo clientInfo,
ChannelProvider channelProvider,
Resource<Metrics> metrics,
Resource<ClientConfigurationManager> configManager,
Resource<ScheduledExecutorService> bgExecutor)
Resource<ScheduledExecutorService> bgExecutor,
ChannelProvider channelProvider)
throws IOException {
this.featureFlags = featureFlags;
this.clientInfo = clientInfo;
Expand All @@ -181,13 +185,34 @@ public Client(
// TODO: consider localizing this for large reads
.maxInboundMessageSize(256 * 1024 * 1024));

channelPool =
SwitchingChannelPool switchingPool =
new SwitchingChannelPool(
configuredChannelProvider,
configManager.get(),
metrics.get(),
backgroundExecutor.get());
channelPool.start();
switchingPool.start();
this.channelPool = Resource.createOwned(switchingPool, switchingPool::close);
}

/**
* Factory-child constructor. Uses a pre-built, shared {@link ChannelPool} and {@link
* ClientConfigurationManager}. The pool is already started and must not be closed by this client.

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

nit: maybe reword "The pool is already started and must not be closed by this client." to something like "The pool should be created with Resource.createShared() and closed when the factory closes"

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

The two public Client(...) constructors are distinguished only by the final parameter (ChannelProvider vs Resource). This is fragile — a factory (a static factory method or clearly named builders) would be safer than type-only overload resolution, especially since one throws IOException and the other doesn't.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Sorry, I can wrap it in static factory, but looks like we'll still need two constructors like these. Leaving it as is for now.

*/
public Client(
FeatureFlags featureFlags,
ClientInfo clientInfo,
Resource<Metrics> metrics,
Resource<ClientConfigurationManager> configManager,
Resource<ScheduledExecutorService> bgExecutor,
Resource<ChannelPool> sharedChannelPool) {
this.featureFlags = featureFlags;
this.clientInfo = clientInfo;
this.metrics = metrics;
this.configManager = configManager;
this.backgroundExecutor = bgExecutor;
this.channelPool = sharedChannelPool;
defaultCallOptions = CallOptions.DEFAULT;
}

@Override
Expand All @@ -200,7 +225,7 @@ public void close() {
.setDescription("Client closing")
.build()));
metrics.close();
channelPool.close();
channelPool.close(); // no-op when Resource.createShared (factory child)
Comment thread
nimf marked this conversation as resolved.
Outdated
configManager.close();
backgroundExecutor.close();
}
Expand All @@ -211,7 +236,7 @@ public TableAsync openTableAsync(String tableId, Permission permission) {
featureFlags,
clientInfo,
configManager.get(),
channelPool,
channelPool.get(),
defaultCallOptions,
tableId,
permission,
Expand All @@ -228,7 +253,7 @@ public AuthorizedViewAsync openAuthorizedViewAsync(
featureFlags,
clientInfo,
configManager.get(),
channelPool,
channelPool.get(),
defaultCallOptions,
tableId,
viewId,
Expand All @@ -246,7 +271,7 @@ public MaterializedViewAsync openMaterializedViewAsync(
featureFlags,
clientInfo,
configManager.get(),
channelPool,
channelPool.get(),
defaultCallOptions,
viewId,
permission,
Expand All @@ -256,6 +281,16 @@ public MaterializedViewAsync openMaterializedViewAsync(
return viewAsync;
}

/** Returns the underlying channel pool (e.g. for sharing with factory children). */
public ChannelPool getChannelPool() {
return channelPool.get();
}

/** Returns the feature flags (e.g. for sharing with factory children). */
public FeatureFlags getFeatureFlags() {
return featureFlags;
}

public static class Resource<T> {
private T value;
private Runnable closer;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -22,6 +22,8 @@
import com.google.bigtable.v2.SessionReadRowRequest;
import com.google.bigtable.v2.SessionReadRowResponse;
import com.google.cloud.bigtable.data.v2.internal.channels.ChannelPool;
import com.google.cloud.bigtable.data.v2.internal.channels.ChannelPoolOptions;
import com.google.cloud.bigtable.data.v2.internal.channels.TenantKey;
import com.google.cloud.bigtable.data.v2.internal.csm.Metrics;
import com.google.cloud.bigtable.data.v2.internal.csm.attributes.ClientInfo;
import com.google.cloud.bigtable.data.v2.internal.csm.tracers.VRpcTracer;
Expand Down Expand Up @@ -63,14 +65,20 @@ static <ReqT extends Message> TableBase createAndStart(
Metrics metrics,
ScheduledExecutorService executor) {

// Stamp the tenant key so ChannelPoolDpImpl can make tenant-aware placement decisions.
CallOptions stamped =
callOptions.withOption(
ChannelPoolOptions.TENANT_KEY_OPTION,
new TenantKey(clientInfo.getInstanceName(), clientInfo.getAppProfileId()));

SessionPool<ReqT> sessionPool =
new SessionPoolImpl<>(
metrics,
featureFlags,
clientInfo,
configManager,
channelPool,
callOptions,
stamped,
sessionDescriptor,
sessionPoolName,
executor);
Expand Down
Loading
Loading