From 41f8daacf37d86889e757ef0cdb15b181460971d Mon Sep 17 00:00:00 2001 From: Keshav Dandeva Date: Fri, 7 Aug 2026 15:41:20 +0000 Subject: [PATCH 1/6] fix(bigquery-jdbc): pass connection proxy settings to OpenTelemetry trace and log exporters --- .../bigquery/jdbc/BigQueryConnection.java | 14 +++- .../jdbc/BigQueryJdbcOpenTelemetry.java | 64 +++++++++++++++++-- .../bigquery/jdbc/BigQueryConnectionTest.java | 3 +- .../jdbc/BigQueryJdbcOpenTelemetryTest.java | 28 ++++---- 4 files changed, 86 insertions(+), 23 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java index a4aee6142f40..e212757bc9b7 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java @@ -200,6 +200,7 @@ public class BigQueryConnection extends BigQueryNoOpsConnection { List queryProperties; Map authProperties; Map overrideProperties; + Map proxyProperties; Credentials credentials; boolean useStatelessQueryMode; int numBufferedRows; @@ -299,7 +300,7 @@ public class BigQueryConnection extends BigQueryNoOpsConnection { String.valueOf(ds.getRequestGoogleDriveScope()), BigQueryJdbcUrlUtility.REQUEST_GOOGLE_DRIVE_SCOPE_PROPERTY_NAME); - Map proxyProperties = + this.proxyProperties = BigQueryJdbcProxyUtility.parseProxyProperties(ds, this.connectionClassName); this.sslTrustStorePath = ds.getSSLTrustStorePath(); @@ -1204,14 +1205,21 @@ private OpenTelemetry getOpenTelemetryInstance() { this.customOpenTelemetry, this.gcpTelemetryCredentials, effectiveProjectId, - this.credentials); + this.credentials, + this.proxyProperties); boolean hasExternalOtel = this.customOpenTelemetry != null || this.useGlobalOpenTelemetry; Logging localLoggingClient = null; if (this.enableGcpLogExporter && !hasExternalOtel) { localLoggingClient = BigQueryJdbcOpenTelemetry.createLoggingClient( - true, null, this.gcpTelemetryCredentials, effectiveProjectId, this.credentials); + true, + null, + this.gcpTelemetryCredentials, + effectiveProjectId, + this.credentials, + this.httpTransportOptions, + this.headerProvider); } if (this.enableGcpLogExporter || hasExternalOtel) { diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java index 91404005c330..5f2fa9b06b45 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java @@ -16,9 +16,11 @@ package com.google.cloud.bigquery.jdbc; +import com.google.api.gax.rpc.HeaderProvider; import com.google.auth.Credentials; import com.google.auth.oauth2.GoogleCredentials; import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; +import com.google.cloud.http.HttpTransportOptions; import com.google.cloud.logging.Logging; import com.google.cloud.logging.LoggingOptions; import com.google.common.hash.Hashing; @@ -32,9 +34,16 @@ import io.opentelemetry.context.Context; import io.opentelemetry.context.Scope; import io.opentelemetry.exporter.otlp.http.trace.OtlpHttpSpanExporter; +import io.opentelemetry.exporter.otlp.http.trace.OtlpHttpSpanExporterBuilder; import io.opentelemetry.exporter.otlp.trace.OtlpGrpcSpanExporter; import io.opentelemetry.sdk.OpenTelemetrySdk; import io.opentelemetry.sdk.autoconfigure.AutoConfiguredOpenTelemetrySdk; +import io.opentelemetry.sdk.common.export.ProxyOptions; +import java.io.IOException; +import java.net.InetSocketAddress; +import java.net.Proxy; +import java.net.ProxySelector; +import java.net.SocketAddress; import java.net.URI; import java.nio.charset.StandardCharsets; import java.sql.SQLException; @@ -181,7 +190,9 @@ public static Logging createLoggingClient( OpenTelemetry customOpenTelemetry, String effectiveCredentials, String effectiveProjectId, - Credentials fallbackCredentials) { + Credentials fallbackCredentials, + HttpTransportOptions httpTransportOptions, + HeaderProvider headerProvider) { if (!enableGcpLogExporter || customOpenTelemetry != null) { return null; @@ -200,6 +211,12 @@ public static Logging createLoggingClient( if (credentials != null) { loggingOptionsBuilder.setCredentials(credentials); } + if (httpTransportOptions != null) { + loggingOptionsBuilder.setTransportOptions(httpTransportOptions); + } + if (headerProvider != null) { + loggingOptionsBuilder.setHeaderProvider(headerProvider); + } return loggingOptionsBuilder.build().getService(); } catch (Exception e) { throw new BigQueryJdbcRuntimeException("Failed to initialize Logging client", e); @@ -317,7 +334,8 @@ public static OpenTelemetry getOpenTelemetry( OpenTelemetry customOpenTelemetry, String gcpTelemetryCredentials, String gcpTelemetryProjectId, - Credentials fallbackCredentials) { + Credentials fallbackCredentials, + Map proxyProperties) { if (customOpenTelemetry != null) { return customOpenTelemetry; @@ -415,11 +433,17 @@ public static OpenTelemetry getOpenTelemetry( final Credentials finalCredentials = credentials; if (spanExporter instanceof OtlpHttpSpanExporter) { - return ((OtlpHttpSpanExporter) spanExporter) - .toBuilder() - .setHeaders( - () -> getAuthHeaders(finalCredentials, gcpTelemetryProjectId)) - .build(); + OtlpHttpSpanExporterBuilder builder = + ((OtlpHttpSpanExporter) spanExporter).toBuilder(); + builder.setHeaders( + () -> getAuthHeaders(finalCredentials, gcpTelemetryProjectId)); + + ProxyOptions proxyOptions = createProxyOptions(proxyProperties); + if (proxyOptions != null) { + builder.setProxy(proxyOptions); + } + + return builder.build(); } if (spanExporter instanceof OtlpGrpcSpanExporter) { return ((OtlpGrpcSpanExporter) spanExporter) @@ -506,4 +530,30 @@ public static T withTracing( span.end(); } } + + private static ProxyOptions createProxyOptions(Map proxyProperties) { + if (proxyProperties == null + || !proxyProperties.containsKey(BigQueryJdbcUrlUtility.PROXY_HOST_PROPERTY_NAME) + || !proxyProperties.containsKey(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME)) { + return null; + } + + final String host = proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_HOST_PROPERTY_NAME); + final int port = + Integer.parseInt(proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME)); + + ProxySelector proxySelector = + new ProxySelector() { + @Override + public List select(URI uri) { + return Collections.singletonList( + new Proxy(Proxy.Type.HTTP, new InetSocketAddress(host, port))); + } + + @Override + public void connectFailed(URI uri, SocketAddress sa, IOException ioe) {} + }; + + return ProxyOptions.create(proxySelector); + } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java index 17a714b22bc7..b562c5799c16 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java @@ -620,7 +620,7 @@ public void testOpenTelemetryPrecedenceHierarchy( .when( () -> BigQueryJdbcOpenTelemetry.createLoggingClient( - anyBoolean(), any(), any(), any(), any())) + anyBoolean(), any(), any(), any(), any(), any(), any())) .thenReturn(mockLogging); // Stub getOpenTelemetry to return the expected mock based on inputs @@ -634,6 +634,7 @@ public void testOpenTelemetryPrecedenceHierarchy( hasCustom ? eq(mockCustomOtel) : isNull(), any(), any(), + any(), any())) .thenAnswer( invocation -> { diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java index 93d2550a3dbb..e1577e6010fd 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetryTest.java @@ -52,7 +52,7 @@ public void testGetOpenTelemetry_withCustomSdk_returnsCustom() { OpenTelemetry result = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, false, false, mockCustomOtel, null, null, null); + false, false, false, mockCustomOtel, null, null, null, null); assertThat(result).isSameInstanceAs(mockCustomOtel); } @@ -64,7 +64,7 @@ public void testGetOpenTelemetry_withCustomSdkAndFlags_returnsCustom() { // Custom SDK always takes precedence over individual flags OpenTelemetry result = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, true, mockCustomOtel, null, null, null); + false, true, true, mockCustomOtel, null, null, null, null); assertThat(result).isSameInstanceAs(mockCustomOtel); } @@ -72,7 +72,8 @@ public void testGetOpenTelemetry_withCustomSdkAndFlags_returnsCustom() { @Test public void testGetOpenTelemetry_noFlags_returnsNoop() { OpenTelemetry result = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(false, false, false, null, null, null, null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + false, false, false, null, null, null, null, null); assertThat(result).isSameInstanceAs(OpenTelemetry.noop()); } @@ -88,10 +89,10 @@ public void testGetTracer_respectsScopeName() { public void testGetOpenTelemetry_cachesSdkInstances() { OpenTelemetry result1 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); assertThat(result1).isSameInstanceAs(result2); } @@ -100,10 +101,10 @@ public void testGetOpenTelemetry_cachesSdkInstances() { public void testGetOpenTelemetry_createsNewInstanceForDifferentKey() { OpenTelemetry result1 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project2", null); + false, true, false, null, null, "project2", null, null); assertThat(result1).isNotSameInstanceAs(result2); } @@ -111,10 +112,11 @@ public void testGetOpenTelemetry_createsNewInstanceForDifferentKey() { @Test public void testGetOpenTelemetry_createsNewInstanceForDifferentTraceFlag() { OpenTelemetry result1 = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(false, true, true, null, null, "project1", null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + false, true, true, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, false, true, null, null, "project1", null); + false, false, true, null, null, "project1", null, null); assertThat(result1).isNotSameInstanceAs(result2); } @@ -122,10 +124,11 @@ public void testGetOpenTelemetry_createsNewInstanceForDifferentTraceFlag() { @Test public void testGetOpenTelemetry_ignoresEnableLogFlagInCacheKey() { OpenTelemetry result1 = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(false, true, true, null, null, "project1", null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + false, true, true, null, null, "project1", null, null); OpenTelemetry result2 = BigQueryJdbcOpenTelemetry.getOpenTelemetry( - false, true, false, null, null, "project1", null); + false, true, false, null, null, "project1", null, null); assertThat(result1).isSameInstanceAs(result2); } @@ -133,7 +136,8 @@ public void testGetOpenTelemetry_ignoresEnableLogFlagInCacheKey() { @Test public void testGetOpenTelemetry_withUseGlobalOTel_returnsGlobal() { OpenTelemetry result = - BigQueryJdbcOpenTelemetry.getOpenTelemetry(true, false, false, null, null, null, null); + BigQueryJdbcOpenTelemetry.getOpenTelemetry( + true, false, false, null, null, null, null, null); assertThat(result).isSameInstanceAs(GlobalOpenTelemetry.get()); } From 5dbb7e9e451fd7c7275fac4ca791999e7f7a1604 Mon Sep 17 00:00:00 2001 From: Keshav Dandeva Date: Fri, 7 Aug 2026 17:40:47 +0000 Subject: [PATCH 2/6] add debug statemetns --- .../bigquery/jdbc/BigQueryConnection.java | 1 - .../jdbc/BigQueryJdbcOpenTelemetry.java | 5 -- .../bigquery/jdbc/BigQueryConnectionTest.java | 2 +- .../bigquery/jdbc/it/ITOpenTelemetryTest.java | 65 ++++++++++++++++++- 4 files changed, 63 insertions(+), 10 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java index e212757bc9b7..1e7c0e3fb45f 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryConnection.java @@ -1218,7 +1218,6 @@ private OpenTelemetry getOpenTelemetryInstance() { this.gcpTelemetryCredentials, effectiveProjectId, this.credentials, - this.httpTransportOptions, this.headerProvider); } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java index 5f2fa9b06b45..8dd946c9b251 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java @@ -20,7 +20,6 @@ import com.google.auth.Credentials; import com.google.auth.oauth2.GoogleCredentials; import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; -import com.google.cloud.http.HttpTransportOptions; import com.google.cloud.logging.Logging; import com.google.cloud.logging.LoggingOptions; import com.google.common.hash.Hashing; @@ -191,7 +190,6 @@ public static Logging createLoggingClient( String effectiveCredentials, String effectiveProjectId, Credentials fallbackCredentials, - HttpTransportOptions httpTransportOptions, HeaderProvider headerProvider) { if (!enableGcpLogExporter || customOpenTelemetry != null) { @@ -211,9 +209,6 @@ public static Logging createLoggingClient( if (credentials != null) { loggingOptionsBuilder.setCredentials(credentials); } - if (httpTransportOptions != null) { - loggingOptionsBuilder.setTransportOptions(httpTransportOptions); - } if (headerProvider != null) { loggingOptionsBuilder.setHeaderProvider(headerProvider); } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java index b562c5799c16..5798438bcfde 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryConnectionTest.java @@ -620,7 +620,7 @@ public void testOpenTelemetryPrecedenceHierarchy( .when( () -> BigQueryJdbcOpenTelemetry.createLoggingClient( - anyBoolean(), any(), any(), any(), any(), any(), any())) + anyBoolean(), any(), any(), any(), any(), any())) .thenReturn(mockLogging); // Stub getOpenTelemetry to return the expected mock based on inputs diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java index 444b92fcefc0..10ea52852b78 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java @@ -250,6 +250,8 @@ private void verifyTraceDelivery(DataSource ds) throws Exception { } private String verifyAndFetchLogs(String connectionUuid) throws Exception { + System.out.println( + "[DEBUG-OTEL-TEST] verifyAndFetchLogs called for connectionUuid: " + connectionUuid); GoogleCredentials credentials = getCredentials(); try (Logging logging = @@ -259,9 +261,15 @@ private String verifyAndFetchLogs(String connectionUuid) throws Exception { .build() .getService()) { String filter = "labels.\"jdbc.connection_id\"=\"" + connectionUuid + "\""; + System.out.println("[DEBUG-OTEL-TEST] Querying Cloud Logging with filter: " + filter); List entries = fetchLogsWithRetry(logging, filter); - assertFalse(entries.isEmpty(), "Telemetry logs should be exported to GCP"); + System.out.println( + "[DEBUG-OTEL-TEST] Retrieved " + entries.size() + " log entries from Cloud Logging."); + + assertFalse( + entries.isEmpty(), + "Telemetry logs should be exported to GCP for connectionUuid: " + connectionUuid); String traceId = null; String hexSpanId = null; @@ -273,6 +281,8 @@ private String verifyAndFetchLogs(String connectionUuid) throws Exception { } } + System.out.println( + "[DEBUG-OTEL-TEST] Extracted traceId: " + traceId + ", hexSpanId: " + hexSpanId); assertNotNull(traceId, "Log entry must contain TraceId"); assertNotNull(hexSpanId, "Log entry must contain SpanId"); @@ -282,6 +292,14 @@ private String verifyAndFetchLogs(String connectionUuid) throws Exception { } return traceId; + } catch (Exception e) { + System.err.println( + "[DEBUG-OTEL-TEST] Error in verifyAndFetchLogs for connectionUuid " + + connectionUuid + + ": " + + e.getMessage()); + e.printStackTrace(); + throw e; } } @@ -290,6 +308,8 @@ private Trace verifyAndFetchTrace(String traceId) throws Exception { if (traceId.contains("/traces/")) { hexTraceId = traceId.substring(traceId.lastIndexOf("/traces/") + 8); } + System.out.println( + "[DEBUG-OTEL-TEST] verifyAndFetchTrace called for hexTraceId: " + hexTraceId); GoogleCredentials credentials = getCredentials(); @@ -300,8 +320,34 @@ private Trace verifyAndFetchTrace(String traceId) throws Exception { try (TraceServiceClient traceClient = TraceServiceClient.create(settings)) { Trace trace = fetchTraceWithRetry(traceClient, PROJECT_ID, hexTraceId); + if (trace != null) { + System.out.println( + "[DEBUG-OTEL-TEST] Successfully retrieved trace. Spans count: " + + trace.getSpansCount()); + for (TraceSpan span : trace.getSpansList()) { + System.out.println( + "[DEBUG-OTEL-TEST] -> Span Name: " + + span.getName() + + ", SpanId: " + + span.getSpanId() + + ", ParentSpanId: " + + span.getParentSpanId()); + } + } else { + System.err.println( + "[DEBUG-OTEL-TEST] Trace NOT found in Cloud Trace API after retries for hexTraceId: " + + hexTraceId); + } assertNotNull(trace, "Trace must be found in Cloud Trace API: " + hexTraceId); return trace; + } catch (Exception e) { + System.err.println( + "[DEBUG-OTEL-TEST] Error in verifyAndFetchTrace for hexTraceId " + + hexTraceId + + ": " + + e.getMessage()); + e.printStackTrace(); + throw e; } } @@ -310,27 +356,40 @@ private T pollWithRetry(java.util.concurrent.Callable task) throws Interr int maxAttempts = 10; long delayMs = 10000; - // 10 second wait for GCP to ingest data + System.out.println("[DEBUG-OTEL-TEST] Waiting initial 10s for GCP telemetry ingestion..."); Thread.sleep(10000); while (attempts < maxAttempts) { attempts++; + System.out.println("[DEBUG-OTEL-TEST] Poll attempt " + attempts + "/" + maxAttempts + "..."); try { T result = task.call(); if (result != null) { + System.out.println("[DEBUG-OTEL-TEST] Poll attempt " + attempts + " SUCCEEDED!"); return result; } } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Test execution interrupted", e); } catch (Exception e) { - // Ignore exceptions during remote lookup and retry + System.out.println( + "[DEBUG-OTEL-TEST] Poll attempt " + + attempts + + " encountered exception: " + + e.getMessage()); e.printStackTrace(); } if (attempts < maxAttempts) { + System.out.println( + "[DEBUG-OTEL-TEST] Poll attempt " + + attempts + + " returned no data, sleeping " + + (delayMs / 1000) + + "s..."); Thread.sleep(delayMs); } } + System.err.println("[DEBUG-OTEL-TEST] Poll timed out after " + maxAttempts + " attempts."); return null; } From 9712fa9c1f5bce52b6492df394d69ae6c3836595 Mon Sep 17 00:00:00 2001 From: Keshav Dandeva Date: Fri, 7 Aug 2026 23:46:25 +0000 Subject: [PATCH 3/6] revert debug statements --- .../jdbc/BigQueryJdbcOpenTelemetry.java | 27 ++++++-- .../bigquery/jdbc/it/ITOpenTelemetryTest.java | 65 +------------------ 2 files changed, 25 insertions(+), 67 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java index 8dd946c9b251..41acbfef5a27 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java @@ -93,11 +93,25 @@ private static final class SdkCacheKey { private final String projectId; private final String credentialsHashOrPath; private final boolean enableTrace; - - SdkCacheKey(String projectId, String credentialsHashOrPath, boolean enableTrace) { + private final String proxyHost; + private final String proxyPort; + + SdkCacheKey( + String projectId, + String credentialsHashOrPath, + boolean enableTrace, + Map proxyProperties) { this.projectId = projectId; this.credentialsHashOrPath = credentialsHashOrPath; this.enableTrace = enableTrace; + this.proxyHost = + proxyProperties != null + ? proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_HOST_PROPERTY_NAME) + : null; + this.proxyPort = + proxyProperties != null + ? proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME) + : null; } @Override @@ -107,12 +121,14 @@ public boolean equals(Object o) { SdkCacheKey that = (SdkCacheKey) o; return enableTrace == that.enableTrace && Objects.equals(projectId, that.projectId) - && Objects.equals(credentialsHashOrPath, that.credentialsHashOrPath); + && Objects.equals(credentialsHashOrPath, that.credentialsHashOrPath) + && Objects.equals(proxyHost, that.proxyHost) + && Objects.equals(proxyPort, that.proxyPort); } @Override public int hashCode() { - return Objects.hash(projectId, credentialsHashOrPath, enableTrace); + return Objects.hash(projectId, credentialsHashOrPath, enableTrace, proxyHost, proxyPort); } } @@ -348,7 +364,8 @@ public static OpenTelemetry getOpenTelemetry( new SdkCacheKey( gcpTelemetryProjectId, getCredentialsIdentifier(gcpTelemetryCredentials), - enableGcpTraceExporter); + enableGcpTraceExporter, + proxyProperties); CachedSdk fastCheck = sdkCache.get(key); if (fastCheck != null) { CachedSdk result = diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java index 10ea52852b78..444b92fcefc0 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java @@ -250,8 +250,6 @@ private void verifyTraceDelivery(DataSource ds) throws Exception { } private String verifyAndFetchLogs(String connectionUuid) throws Exception { - System.out.println( - "[DEBUG-OTEL-TEST] verifyAndFetchLogs called for connectionUuid: " + connectionUuid); GoogleCredentials credentials = getCredentials(); try (Logging logging = @@ -261,15 +259,9 @@ private String verifyAndFetchLogs(String connectionUuid) throws Exception { .build() .getService()) { String filter = "labels.\"jdbc.connection_id\"=\"" + connectionUuid + "\""; - System.out.println("[DEBUG-OTEL-TEST] Querying Cloud Logging with filter: " + filter); List entries = fetchLogsWithRetry(logging, filter); - System.out.println( - "[DEBUG-OTEL-TEST] Retrieved " + entries.size() + " log entries from Cloud Logging."); - - assertFalse( - entries.isEmpty(), - "Telemetry logs should be exported to GCP for connectionUuid: " + connectionUuid); + assertFalse(entries.isEmpty(), "Telemetry logs should be exported to GCP"); String traceId = null; String hexSpanId = null; @@ -281,8 +273,6 @@ private String verifyAndFetchLogs(String connectionUuid) throws Exception { } } - System.out.println( - "[DEBUG-OTEL-TEST] Extracted traceId: " + traceId + ", hexSpanId: " + hexSpanId); assertNotNull(traceId, "Log entry must contain TraceId"); assertNotNull(hexSpanId, "Log entry must contain SpanId"); @@ -292,14 +282,6 @@ private String verifyAndFetchLogs(String connectionUuid) throws Exception { } return traceId; - } catch (Exception e) { - System.err.println( - "[DEBUG-OTEL-TEST] Error in verifyAndFetchLogs for connectionUuid " - + connectionUuid - + ": " - + e.getMessage()); - e.printStackTrace(); - throw e; } } @@ -308,8 +290,6 @@ private Trace verifyAndFetchTrace(String traceId) throws Exception { if (traceId.contains("/traces/")) { hexTraceId = traceId.substring(traceId.lastIndexOf("/traces/") + 8); } - System.out.println( - "[DEBUG-OTEL-TEST] verifyAndFetchTrace called for hexTraceId: " + hexTraceId); GoogleCredentials credentials = getCredentials(); @@ -320,34 +300,8 @@ private Trace verifyAndFetchTrace(String traceId) throws Exception { try (TraceServiceClient traceClient = TraceServiceClient.create(settings)) { Trace trace = fetchTraceWithRetry(traceClient, PROJECT_ID, hexTraceId); - if (trace != null) { - System.out.println( - "[DEBUG-OTEL-TEST] Successfully retrieved trace. Spans count: " - + trace.getSpansCount()); - for (TraceSpan span : trace.getSpansList()) { - System.out.println( - "[DEBUG-OTEL-TEST] -> Span Name: " - + span.getName() - + ", SpanId: " - + span.getSpanId() - + ", ParentSpanId: " - + span.getParentSpanId()); - } - } else { - System.err.println( - "[DEBUG-OTEL-TEST] Trace NOT found in Cloud Trace API after retries for hexTraceId: " - + hexTraceId); - } assertNotNull(trace, "Trace must be found in Cloud Trace API: " + hexTraceId); return trace; - } catch (Exception e) { - System.err.println( - "[DEBUG-OTEL-TEST] Error in verifyAndFetchTrace for hexTraceId " - + hexTraceId - + ": " - + e.getMessage()); - e.printStackTrace(); - throw e; } } @@ -356,40 +310,27 @@ private T pollWithRetry(java.util.concurrent.Callable task) throws Interr int maxAttempts = 10; long delayMs = 10000; - System.out.println("[DEBUG-OTEL-TEST] Waiting initial 10s for GCP telemetry ingestion..."); + // 10 second wait for GCP to ingest data Thread.sleep(10000); while (attempts < maxAttempts) { attempts++; - System.out.println("[DEBUG-OTEL-TEST] Poll attempt " + attempts + "/" + maxAttempts + "..."); try { T result = task.call(); if (result != null) { - System.out.println("[DEBUG-OTEL-TEST] Poll attempt " + attempts + " SUCCEEDED!"); return result; } } catch (InterruptedException e) { Thread.currentThread().interrupt(); throw new RuntimeException("Test execution interrupted", e); } catch (Exception e) { - System.out.println( - "[DEBUG-OTEL-TEST] Poll attempt " - + attempts - + " encountered exception: " - + e.getMessage()); + // Ignore exceptions during remote lookup and retry e.printStackTrace(); } if (attempts < maxAttempts) { - System.out.println( - "[DEBUG-OTEL-TEST] Poll attempt " - + attempts - + " returned no data, sleeping " - + (delayMs / 1000) - + "s..."); Thread.sleep(delayMs); } } - System.err.println("[DEBUG-OTEL-TEST] Poll timed out after " + maxAttempts + " attempts."); return null; } From ece6c0b98c7ba06215fda7631ec8c04ba5cd8f41 Mon Sep 17 00:00:00 2001 From: Keshav Dandeva Date: Sat, 8 Aug 2026 00:15:54 +0000 Subject: [PATCH 4/6] add proxy test for trace exporter --- .../bigquery/jdbc/it/ITOpenTelemetryTest.java | 86 +++++++++++++++++-- 1 file changed, 81 insertions(+), 5 deletions(-) diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java index 444b92fcefc0..fcd870d38920 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java @@ -23,7 +23,9 @@ import static org.junit.jupiter.api.Assertions.assertTrue; import com.google.api.gax.core.FixedCredentialsProvider; +import com.google.api.gax.grpc.InstantiatingGrpcChannelProvider; import com.google.api.gax.paging.Page; +import com.google.api.gax.rpc.TransportChannelProvider; import com.google.auth.oauth2.GoogleCredentials; import com.google.cloud.ServiceOptions; import com.google.cloud.bigquery.jdbc.BigQueryConnection; @@ -36,7 +38,18 @@ import com.google.devtools.cloudtrace.v1.Trace; import com.google.devtools.cloudtrace.v1.TraceSpan; import com.google.gson.JsonObject; +import io.grpc.HttpConnectProxiedSocketAddress; +import io.grpc.ProxiedSocketAddress; +import io.grpc.ProxyDetector; +import io.opentelemetry.api.trace.Span; +import io.opentelemetry.api.trace.SpanContext; +import io.opentelemetry.api.trace.TraceFlags; +import io.opentelemetry.api.trace.TraceState; +import io.opentelemetry.context.Scope; +import io.opentelemetry.sdk.trace.IdGenerator; import java.io.File; +import java.net.InetSocketAddress; +import java.net.SocketAddress; import java.nio.charset.StandardCharsets; import java.nio.file.Files; import java.sql.Connection; @@ -45,6 +58,7 @@ import java.sql.Statement; import java.util.ArrayList; import java.util.List; +import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; public class ITOpenTelemetryTest extends ITBase { @@ -53,6 +67,7 @@ public class ITOpenTelemetryTest extends ITBase { private static final String CONNECTION_URL = connectionUrl; @Test + @Tag("known_issue") public void testExecute_withOpenTelemetryGcpExporter() throws Exception { // Step 1: Connect with GCP Exporters enabled via DataSource @@ -124,8 +139,8 @@ public void testExecute_withOpenTelemetryGcpExporter() throws Exception { } @Test + @Tag("known_issue") public void testExecute_withErrorCorrelation() throws Exception { - // Step 1: Connect with GCP Exporters enabled via DataSource DataSource ds = DataSource.fromUrl(CONNECTION_URL); ds.setEnableGcpTraceExporter(true); @@ -168,6 +183,7 @@ public void testExecute_withErrorCorrelation() throws Exception { } @Test + @Tag("known_issue") public void testExecute_withCustomCredentialsJson() throws Exception { JsonObject authJson = getAuthJson(); DataSource ds = DataSource.fromUrl(CONNECTION_URL); @@ -179,6 +195,7 @@ public void testExecute_withCustomCredentialsJson() throws Exception { } @Test + @Tag("known_issue") public void testExecute_withCustomCredentialsFilePath() throws Exception { JsonObject authJson = getAuthJson(); File tempFile = File.createTempFile("auth", ".json"); @@ -194,6 +211,7 @@ public void testExecute_withCustomCredentialsFilePath() throws Exception { } @Test + @Tag("known_issue") public void testExecute_withHttpProtocol() throws Exception { JsonObject authJson = getAuthJson(); System.setProperty("otel.exporter.otlp.protocol", "http/protobuf"); @@ -211,6 +229,7 @@ public void testExecute_withHttpProtocol() throws Exception { } @Test + @Tag("known_issue") public void testExecute_withGrpcProtocol() throws Exception { JsonObject authJson = getAuthJson(); System.setProperty("otel.exporter.otlp.protocol", "grpc"); @@ -227,6 +246,39 @@ public void testExecute_withGrpcProtocol() throws Exception { } } + @Test + public void testExecute_withHttpProtocol_andDirectTraceVerification() throws Exception { + JsonObject authJson = getAuthJson(); + System.setProperty("otel.exporter.otlp.protocol", "http/protobuf"); + + try { + DataSource ds = DataSource.fromUrl(CONNECTION_URL); + ds.setEnableGcpTraceExporter(true); + ds.setGcpTelemetryProjectId(PROJECT_ID); + ds.setGcpTelemetryCredentials(authJson.toString()); + + String traceId = IdGenerator.random().generateTraceId(); + String spanId = IdGenerator.random().generateSpanId(); + SpanContext parentContext = + SpanContext.create(traceId, spanId, TraceFlags.getSampled(), TraceState.getDefault()); + + try (Scope scope = Span.wrap(parentContext).makeCurrent()) { + try (Connection connection = ds.getConnection(); + Statement statement = connection.createStatement()) { + String query = "SELECT 1;"; + try (ResultSet rs = statement.executeQuery(query)) { + assertTrue(rs.next()); + } + } + } + + Trace trace = verifyAndFetchTrace(traceId); + assertNotNull(trace, "Trace must be found in Cloud Trace API: " + traceId); + } finally { + System.clearProperty("otel.exporter.otlp.protocol"); + } + } + private void verifyTraceDelivery(DataSource ds) throws Exception { ds.setEnableGcpLogExporter(true); ds.setLogLevel("5"); @@ -293,18 +345,42 @@ private Trace verifyAndFetchTrace(String traceId) throws Exception { GoogleCredentials credentials = getCredentials(); - TraceServiceSettings settings = + TraceServiceSettings.Builder settingsBuilder = TraceServiceSettings.newBuilder() - .setCredentialsProvider(FixedCredentialsProvider.create(credentials)) - .build(); + .setCredentialsProvider(FixedCredentialsProvider.create(credentials)); + + DataSource ds = DataSource.fromUrl(CONNECTION_URL); + String proxyHost = ds.getProxyHost(); + String proxyPortStr = ds.getProxyPort(); + if (proxyHost != null && proxyPortStr != null) { + settingsBuilder.setTransportChannelProvider( + createProxyTransportChannelProvider(proxyHost, Integer.parseInt(proxyPortStr))); + } - try (TraceServiceClient traceClient = TraceServiceClient.create(settings)) { + try (TraceServiceClient traceClient = TraceServiceClient.create(settingsBuilder.build())) { Trace trace = fetchTraceWithRetry(traceClient, PROJECT_ID, hexTraceId); assertNotNull(trace, "Trace must be found in Cloud Trace API: " + hexTraceId); return trace; } } + private TransportChannelProvider createProxyTransportChannelProvider(String host, int port) { + return InstantiatingGrpcChannelProvider.newBuilder() + .setChannelConfigurator( + builder -> + builder.proxyDetector( + new ProxyDetector() { + @Override + public ProxiedSocketAddress proxyFor(SocketAddress socketAddress) { + return HttpConnectProxiedSocketAddress.newBuilder() + .setProxyAddress(new InetSocketAddress(host, port)) + .setTargetAddress((InetSocketAddress) socketAddress) + .build(); + } + })) + .build(); + } + private T pollWithRetry(java.util.concurrent.Callable task) throws InterruptedException { int attempts = 0; int maxAttempts = 10; From dacab8cdbe2c07a40c8848fff62e61f93af8ca92 Mon Sep 17 00:00:00 2001 From: Keshav Dandeva Date: Sat, 8 Aug 2026 00:39:27 +0000 Subject: [PATCH 5/6] test(bigquery-jdbc): clean up integration tests and finalize trace proxy test --- .../google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java | 7 ------- 1 file changed, 7 deletions(-) diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java index fcd870d38920..814c8329bb5f 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITOpenTelemetryTest.java @@ -58,7 +58,6 @@ import java.sql.Statement; import java.util.ArrayList; import java.util.List; -import org.junit.jupiter.api.Tag; import org.junit.jupiter.api.Test; public class ITOpenTelemetryTest extends ITBase { @@ -67,7 +66,6 @@ public class ITOpenTelemetryTest extends ITBase { private static final String CONNECTION_URL = connectionUrl; @Test - @Tag("known_issue") public void testExecute_withOpenTelemetryGcpExporter() throws Exception { // Step 1: Connect with GCP Exporters enabled via DataSource @@ -139,7 +137,6 @@ public void testExecute_withOpenTelemetryGcpExporter() throws Exception { } @Test - @Tag("known_issue") public void testExecute_withErrorCorrelation() throws Exception { // Step 1: Connect with GCP Exporters enabled via DataSource DataSource ds = DataSource.fromUrl(CONNECTION_URL); @@ -183,7 +180,6 @@ public void testExecute_withErrorCorrelation() throws Exception { } @Test - @Tag("known_issue") public void testExecute_withCustomCredentialsJson() throws Exception { JsonObject authJson = getAuthJson(); DataSource ds = DataSource.fromUrl(CONNECTION_URL); @@ -195,7 +191,6 @@ public void testExecute_withCustomCredentialsJson() throws Exception { } @Test - @Tag("known_issue") public void testExecute_withCustomCredentialsFilePath() throws Exception { JsonObject authJson = getAuthJson(); File tempFile = File.createTempFile("auth", ".json"); @@ -211,7 +206,6 @@ public void testExecute_withCustomCredentialsFilePath() throws Exception { } @Test - @Tag("known_issue") public void testExecute_withHttpProtocol() throws Exception { JsonObject authJson = getAuthJson(); System.setProperty("otel.exporter.otlp.protocol", "http/protobuf"); @@ -229,7 +223,6 @@ public void testExecute_withHttpProtocol() throws Exception { } @Test - @Tag("known_issue") public void testExecute_withGrpcProtocol() throws Exception { JsonObject authJson = getAuthJson(); System.setProperty("otel.exporter.otlp.protocol", "grpc"); From e9e6b2b59e4eb16a0d51feb5af5f7f7bead62601 Mon Sep 17 00:00:00 2001 From: Keshav Dandeva Date: Sat, 8 Aug 2026 00:44:39 +0000 Subject: [PATCH 6/6] fix(bigquery-jdbc): address PR review comments for createProxyOptions --- .../jdbc/BigQueryJdbcOpenTelemetry.java | 19 +++++++++++++------ 1 file changed, 13 insertions(+), 6 deletions(-) diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java index 41acbfef5a27..3ecca298480d 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryJdbcOpenTelemetry.java @@ -544,22 +544,29 @@ public static T withTracing( } private static ProxyOptions createProxyOptions(Map proxyProperties) { - if (proxyProperties == null - || !proxyProperties.containsKey(BigQueryJdbcUrlUtility.PROXY_HOST_PROPERTY_NAME) - || !proxyProperties.containsKey(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME)) { + if (proxyProperties == null) { return null; } final String host = proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_HOST_PROPERTY_NAME); - final int port = - Integer.parseInt(proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME)); + final String portStr = proxyProperties.get(BigQueryJdbcUrlUtility.PROXY_PORT_PROPERTY_NAME); + if (host == null || host.isEmpty() || portStr == null || portStr.isEmpty()) { + return null; + } + + int port; + try { + port = Integer.parseInt(portStr); + } catch (NumberFormatException e) { + throw new BigQueryJdbcRuntimeException("Invalid proxy port number: " + portStr, e); + } ProxySelector proxySelector = new ProxySelector() { @Override public List select(URI uri) { return Collections.singletonList( - new Proxy(Proxy.Type.HTTP, new InetSocketAddress(host, port))); + new Proxy(Proxy.Type.HTTP, InetSocketAddress.createUnresolved(host, port))); } @Override