1717 */
1818package org .apache .beam .sdk .io .gcp .spanner ;
1919
20- import com .google .api .gax .core .ExecutorProvider ;
21- import com .google .api .gax .grpc .InstantiatingGrpcChannelProvider ;
2220import com .google .api .gax .retrying .RetrySettings ;
23- import com .google .api .gax .rpc .HeaderProvider ;
21+ import com .google .api .gax .rpc .FixedHeaderProvider ;
2422import com .google .api .gax .rpc .ServerStreamingCallSettings ;
2523import com .google .api .gax .rpc .UnaryCallSettings ;
2624import com .google .cloud .NoCredentials ;
3129import com .google .cloud .spanner .DatabaseId ;
3230import com .google .cloud .spanner .Spanner ;
3331import com .google .cloud .spanner .SpannerOptions ;
34- import com .google .cloud .spanner .spi .v1 .SpannerInterceptorProvider ;
3532import com .google .spanner .v1 .CommitRequest ;
3633import com .google .spanner .v1 .CommitResponse ;
3734import com .google .spanner .v1 .ExecuteSqlRequest ;
4138import io .grpc .ClientCall ;
4239import io .grpc .ClientInterceptor ;
4340import io .grpc .MethodDescriptor ;
44- import java .net .MalformedURLException ;
45- import java .net .URL ;
46- import java .util .ArrayList ;
47- import java .util .HashMap ;
48- import java .util .List ;
49- import java .util .Map ;
5041import java .util .concurrent .ConcurrentHashMap ;
51- import java .util .concurrent .ScheduledExecutorService ;
52- import java .util .concurrent .ScheduledThreadPoolExecutor ;
53- import java .util .concurrent .ThreadFactory ;
5442import java .util .concurrent .TimeUnit ;
5543import java .util .concurrent .atomic .AtomicInteger ;
5644import org .apache .beam .sdk .options .ValueProvider ;
5745import org .apache .beam .sdk .util .ReleaseInfo ;
58- import org .apache .beam .vendor .guava .v26_0_jre .com .google .common .util .concurrent .ThreadFactoryBuilder ;
5946import org .joda .time .Duration ;
6047import org .slf4j .Logger ;
6148import org .slf4j .LoggerFactory ;
@@ -88,12 +75,6 @@ public class SpannerAccessor implements AutoCloseable {
8875 private final DatabaseAdminClient databaseAdminClient ;
8976 private final SpannerConfig spannerConfig ;
9077
91- private static final int MAX_MESSAGE_SIZE = 100 * 1024 * 1024 ;
92- private static final int MAX_METADATA_SIZE = 32 * 1024 ; // bytes
93- private static final int NUM_CHANNELS = 4 ;
94- public static final org .threeten .bp .Duration GRPC_KEEP_ALIVE_SECONDS =
95- org .threeten .bp .Duration .ofSeconds (120 );
96-
9778 private SpannerAccessor (
9879 Spanner spanner ,
9980 DatabaseClient databaseClient ,
@@ -161,23 +142,6 @@ private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig) {
161142 .setTotalTimeout (org .threeten .bp .Duration .ofMinutes (120 ))
162143 .build ());
163144
164- ManagedInstantiatingExecutorProvider executorProvider =
165- new ManagedInstantiatingExecutorProvider (
166- new ThreadFactoryBuilder ()
167- .setDaemon (true )
168- .setNameFormat ("Cloud-Spanner-TransportChannel-%d" )
169- .build ());
170-
171- InstantiatingGrpcChannelProvider .Builder instantiatingGrpcChannelProvider =
172- InstantiatingGrpcChannelProvider .newBuilder ()
173- .setMaxInboundMessageSize (MAX_MESSAGE_SIZE )
174- .setMaxInboundMetadataSize (MAX_METADATA_SIZE )
175- .setPoolSize (NUM_CHANNELS )
176- .setExecutorProvider (executorProvider )
177- .setKeepAliveTime (GRPC_KEEP_ALIVE_SECONDS )
178- .setInterceptorProvider (SpannerInterceptorProvider .createDefault ())
179- .setAttemptDirectPath (true );
180-
181145 ValueProvider <String > projectId = spannerConfig .getProjectId ();
182146 if (projectId != null ) {
183147 builder .setProjectId (projectId .get ());
@@ -189,34 +153,14 @@ private static SpannerAccessor createAndConnect(SpannerConfig spannerConfig) {
189153 ValueProvider <String > host = spannerConfig .getHost ();
190154 if (host != null ) {
191155 builder .setHost (host .get ());
192- instantiatingGrpcChannelProvider .setEndpoint (getEndpoint (host .get ()));
193156 }
194157 ValueProvider <String > emulatorHost = spannerConfig .getEmulatorHost ();
195158 if (emulatorHost != null ) {
196159 builder .setEmulatorHost (emulatorHost .get ());
197160 builder .setCredentials (NoCredentials .getInstance ());
198- } else {
199- String userAgentString = USER_AGENT_PREFIX + "/" + ReleaseInfo .getReleaseInfo ().getVersion ();
200- /* Workaround to setup user-agent string.
201- * InstantiatingGrpcChannelProvider will override the settings provided.
202- * The section below and all associated artifacts will be removed once the bug
203- * that prevents setting user-agent is fixed.
204- * https://github.com/googleapis/java-spanner/pull/871
205- *
206- * Code to be replaced:
207- * builder.setHeaderProvider(FixedHeaderProvider.create("user-agent", userAgentString));
208- */
209- instantiatingGrpcChannelProvider .setHeaderProvider (
210- new HeaderProvider () {
211- @ Override
212- public Map <String , String > getHeaders () {
213- final Map <String , String > headers = new HashMap <>();
214- headers .put ("user-agent" , userAgentString );
215- return headers ;
216- }
217- });
218- builder .setChannelProvider (instantiatingGrpcChannelProvider .build ());
219161 }
162+ String userAgentString = USER_AGENT_PREFIX + "/" + ReleaseInfo .getReleaseInfo ().getVersion ();
163+ builder .setHeaderProvider (FixedHeaderProvider .create ("user-agent" , userAgentString ));
220164 SpannerOptions options = builder .build ();
221165
222166 Spanner spanner = options .getService ();
@@ -232,17 +176,6 @@ public Map<String, String> getHeaders() {
232176 spanner , databaseClient , databaseAdminClient , batchClient , spannerConfig );
233177 }
234178
235- private static String getEndpoint (String host ) {
236- URL url ;
237- try {
238- url = new URL (host );
239- } catch (MalformedURLException e ) {
240- throw new IllegalArgumentException ("Invalid host: " + host , e );
241- }
242- return String .format (
243- "%s:%s" , url .getHost (), url .getPort () < 0 ? url .getDefaultPort () : url .getPort ());
244- }
245-
246179 public DatabaseClient getDatabaseClient () {
247180 return databaseClient ;
248181 }
@@ -291,32 +224,4 @@ public <ReqT, RespT> ClientCall<ReqT, RespT> interceptCall(
291224 return next .newCall (method , callOptions );
292225 }
293226 }
294-
295- private static final class ManagedInstantiatingExecutorProvider implements ExecutorProvider {
296- // 4 Gapic clients * 4 channels per client.
297- private static final int DEFAULT_MIN_THREAD_COUNT = 16 ;
298- private final List <ScheduledExecutorService > executors = new ArrayList <>();
299- private final ThreadFactory threadFactory ;
300-
301- private ManagedInstantiatingExecutorProvider (ThreadFactory threadFactory ) {
302- this .threadFactory = threadFactory ;
303- }
304-
305- @ Override
306- public boolean shouldAutoClose () {
307- return false ;
308- }
309-
310- @ Override
311- public ScheduledExecutorService getExecutor () {
312- int numCpus = Runtime .getRuntime ().availableProcessors ();
313- int numThreads = Math .max (DEFAULT_MIN_THREAD_COUNT , numCpus );
314- ScheduledExecutorService executor =
315- new ScheduledThreadPoolExecutor (numThreads , threadFactory );
316- synchronized (this ) {
317- executors .add (executor );
318- }
319- return executor ;
320- }
321- }
322227}
0 commit comments