diff --git a/java-bigquery-jdbc/README.MD b/java-bigquery-jdbc/README.MD index cd5cde8e795f..34bbb7b38c54 100644 --- a/java-bigquery-jdbc/README.MD +++ b/java-bigquery-jdbc/README.MD @@ -92,10 +92,6 @@ a connection property (`EnableDiagnosticTelemetry=0`), environment variable, or system property, the service immediately shuts down and remains permanently off for the remainder of the JVM lifecycle. It will only restart when the JVM restarts. -* Global Worker Settings: Because the telemetry batcher is a shared global -resource, properties that configure its behavior (`TelemetryUploadInterval` -and `TelemetryBatchSize`) are exclusively evaluated upon establishing the first -connection. Subsequent overrides for these specific properties are ignored. #### Configuration Keys @@ -103,16 +99,8 @@ To configure these properties, you can use the following keys: * Opt-Out Control: * Connection Property: `EnableDiagnosticTelemetry=0;` - * Environment Variable: `GOOGLE_BIGQUERY_JDBC_TELEMETRY_ENABLED=false` - * System Property: `-Dgoogle.bigquery.jdbc.telemetry.enabled=false` -* Upload Interval: - * Connection Property: `TelemetryUploadInterval=300000;` - * Environment Variable: `GOOGLE_BIGQUERY_JDBC_TELEMETRY_INTERVAL_MS=300000` - * System Property: `-Dgoogle.bigquery.jdbc.telemetry.interval_ms=300000` -* Batch Size: - * Connection Property: `TelemetryBatchSize=5000;` - * Environment Variable: `GOOGLE_BIGQUERY_JDBC_TELEMETRY_BATCH_SIZE=5000` - * System Property: `-Dgoogle.bigquery.jdbc.telemetry.batch_size=5000` + * Environment Variable: `BIGQUERY_JDBC_TELEMETRY_ENABLED=false` + * System Property: `-DBIGQUERY_JDBC_TELEMETRY_ENABLED=false` ## Developer Guide diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryBaseResultSet.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryBaseResultSet.java index b0b981181a3c..f6a380884ad2 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryBaseResultSet.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryBaseResultSet.java @@ -28,6 +28,8 @@ import com.google.cloud.bigquery.StandardSQLTypeName; import com.google.cloud.bigquery.exception.BigQueryConversionException; import com.google.cloud.bigquery.exception.BigQueryJdbcException; +import com.google.cloud.bigquery.jdbc.telemetry.v1.DriverFeature; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import io.opentelemetry.api.trace.Span; import io.opentelemetry.api.trace.SpanContext; import io.opentelemetry.context.Context; @@ -258,6 +260,11 @@ public ResultSetMetaData getMetaData() throws SQLException { metaData = BigQueryResultSetMetadata.of(this.schema.getFields(), this.statement); } } + + TelemetryManager.recordFeatureUsage( + DriverFeature.DRIVER_FEATURE_METADATA_RETRIEVAL, + "DRIVER_FEATURE_RESULTSET_METADATA_RETRIEVAL"); + return BigQueryJdbcContextProxy.wrap(metaData, ResultSetMetaData.class, connectionId); } 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 423c4f6dd65c..af93ac9c5bdb 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 @@ -40,6 +40,8 @@ import com.google.cloud.bigquery.exception.BigQueryJdbcException; import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; import com.google.cloud.bigquery.exception.BigQueryJdbcSqlFeatureNotSupportedException; +import com.google.cloud.bigquery.jdbc.telemetry.v1.DriverFeature; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import com.google.cloud.bigquery.storage.v1.BigQueryReadClient; import com.google.cloud.bigquery.storage.v1.BigQueryReadSettings; import com.google.cloud.bigquery.storage.v1.BigQueryWriteClient; @@ -483,6 +485,9 @@ public Statement createStatement() throws SQLException { BigQueryStatement currentStatement = new BigQueryStatement(this); LOG.fine("Statement %s created.", currentStatement); addOpenStatements(currentStatement); + + TelemetryManager.recordFeatureUsage( + DriverFeature.DRIVER_FEATURE_CUSTOM, "DRIVER_FEATURE_REGULAR_STATEMENT"); return currentStatement; } @@ -541,6 +546,9 @@ public PreparedStatement prepareStatement(String sql) throws SQLException { PreparedStatement currentStatement = new BigQueryPreparedStatement(this, sql); LOG.fine("Prepared Statement %s created.", currentStatement); addOpenStatements(currentStatement); + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_PREPARED_STATEMENT); + return currentStatement; } @@ -681,6 +689,8 @@ private void beginTransaction() { updateSessionInfo(transactionBeginJob.getStatistics().getSessionInfo().getSessionId()); } this.transactionStarted = true; + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_TRANSACTIONS); } catch (InterruptedException ex) { throw new BigQueryJdbcRuntimeException("Failed to begin transaction", ex); } @@ -915,6 +925,12 @@ public void setAutoCommit(boolean autoCommit) throws SQLException { if (!this.autoCommit) { beginTransaction(); } + + if (autoCommit) { + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_AUTOCOMMIT_ENABLED); + } else { + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_AUTOCOMMIT_DISABLED); + } } @Override @@ -972,6 +988,9 @@ public DatabaseMetaData getMetaData() throws SQLException { if (databaseMetaData == null) { databaseMetaData = new BigQueryDatabaseMetaData(this); } + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_METADATA_RETRIEVAL); + return databaseMetaData; } @@ -1474,6 +1493,9 @@ public CallableStatement prepareCall(String sql) throws SQLException { CallableStatement currentStatement = new BigQueryCallableStatement(this, sql); LOG.fine("Callable Statement %s created.", currentStatement); addOpenStatements(currentStatement); + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_CALLABLE_STATEMENT); + return currentStatement; } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDriver.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDriver.java index 8c748b1f52bc..140636632316 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDriver.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryDriver.java @@ -18,6 +18,9 @@ import com.google.cloud.bigquery.exception.BigQueryJdbcException; import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; +import com.google.cloud.bigquery.jdbc.telemetry.v1.AuthenticationType; +import com.google.cloud.bigquery.jdbc.telemetry.v1.Status; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import com.google.cloud.bigquery.jdbc.utils.BigQueryJdbcVersionUtility; import io.grpc.LoadBalancerRegistry; import io.grpc.internal.PickFirstLoadBalancerProvider; @@ -124,6 +127,7 @@ public static BigQueryDriver getRegisteredDriver() throws IllegalStateException @Override public Connection connect(String url, Properties info) throws SQLException { LOG.finest("++enter++"); + AuthenticationType authType = AuthenticationType.AUTHENTICATION_TYPE_UNSPECIFIED; try { if (acceptsURL(url)) { Properties connectInfo = info == null ? new Properties() : (Properties) info.clone(); @@ -132,6 +136,17 @@ public Connection connect(String url, Properties info) throws SQLException { String connectionUri = BigQueryJdbcUrlUtility.appendPropertiesToURL( url.substring(5), this.toString(), connectInfo); + + String telemetryOptOut = + BigQueryJdbcUrlUtility.parseUriPropertyWithoutValidation( + connectionUri, BigQueryJdbcUrlUtility.ENABLE_DIAGNOSTIC_TELEMETRY_PROPERTY_NAME); + + if (telemetryOptOut != null) { + connectInfo.setProperty( + BigQueryJdbcUrlUtility.ENABLE_DIAGNOSTIC_TELEMETRY_PROPERTY_NAME, telemetryOptOut); + } + TelemetryManager.getInstance(connectInfo); + Level logLevel; String logPath; try { @@ -200,14 +215,25 @@ public Connection connect(String url, Properties info) throws SQLException { logLevel, logPath, this.toString()); - return BigQueryJdbcContextProxy.wrap(connection, Connection.class); + + Connection wrapped = BigQueryJdbcContextProxy.wrap(connection, Connection.class); + + authType = TelemetryManager.toAuthenticationType(ds.getOAuthType()); + TelemetryManager.recordConnectionAttempt(Status.STATUS_SUCCESS, 0, authType); + return wrapped; } else { return null; } } catch (IOException e) { - LOG.warning("Getting a warning: " + e.getMessage()); + TelemetryManager.recordConnectionAttempt( + Status.STATUS_ERROR, TelemetryManager.extractErrorCode(e), authType); + LOG.warning("Getting a warning: %s", e.getMessage()); + return null; + } catch (Throwable t) { + TelemetryManager.recordConnectionAttempt( + Status.STATUS_ERROR, TelemetryManager.extractErrorCode(t), authType); + throw t; } - return null; } /** diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryParameterHandler.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryParameterHandler.java index 70719fd91a0a..fcf0500def23 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryParameterHandler.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryParameterHandler.java @@ -21,6 +21,8 @@ import com.google.cloud.bigquery.StandardSQLTypeName; import com.google.cloud.bigquery.exception.BigQueryJdbcException; import com.google.cloud.bigquery.exception.BigQueryJdbcSqlFeatureNotSupportedException; +import com.google.cloud.bigquery.jdbc.telemetry.v1.DriverFeature; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import java.math.BigInteger; import java.sql.SQLException; import java.sql.Time; @@ -157,6 +159,8 @@ void setParameter(int parameterIndex, Object value, Class type) parameter.setParamType(BigQueryStatementParameterType.UNSPECIFIED); parameter.setScale(-1); + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_PARAMETER_BINDING); + LOG.finest("Parameter set { %s }", parameter.toString()); } @@ -239,6 +243,9 @@ void setParameter( if (parameter.getIndex() == -1) { parametersList.add(parameter); } + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_PARAMETER_BINDING); + LOG.finest("Parameter set { %s }", parameter.toString()); } @@ -272,6 +279,8 @@ void setParameter( parameter.setParamType(paramType); parameter.setScale(scale); + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_PARAMETER_BINDING); + LOG.finest("Parameter set { %s }", parameter.toString()); } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPooledConnection.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPooledConnection.java index 99dea20e21ce..ad62ee76bcdb 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPooledConnection.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPooledConnection.java @@ -16,8 +16,11 @@ package com.google.cloud.bigquery.jdbc; +import com.google.cloud.bigquery.jdbc.telemetry.v1.DriverFeature; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import com.google.common.annotations.VisibleForTesting; import java.sql.Connection; +import java.sql.DatabaseMetaData; import java.sql.SQLException; import java.util.UUID; import java.util.concurrent.Executor; @@ -231,7 +234,8 @@ public void rollback() throws SQLException { } @Override - public java.sql.DatabaseMetaData getMetaData() throws SQLException { + public DatabaseMetaData getMetaData() throws SQLException { + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_METADATA_RETRIEVAL); return bqConnectionDelegate.getMetaData(); } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPreparedStatement.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPreparedStatement.java index b7dd3465b224..86ab310bc8c6 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPreparedStatement.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryPreparedStatement.java @@ -28,6 +28,9 @@ import com.google.cloud.bigquery.exception.BigQueryJdbcException; import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; import com.google.cloud.bigquery.exception.BigQueryJdbcSqlFeatureNotSupportedException; +import com.google.cloud.bigquery.jdbc.telemetry.v1.DriverFeature; +import com.google.cloud.bigquery.jdbc.telemetry.v1.StatementExecution; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import com.google.cloud.bigquery.storage.v1.BatchCommitWriteStreamsRequest; import com.google.cloud.bigquery.storage.v1.BatchCommitWriteStreamsResponse; import com.google.cloud.bigquery.storage.v1.BigQueryWriteClient; @@ -323,7 +326,16 @@ public int[] executeBatch() throws SQLException { if (this.batchParameters.isEmpty()) { return result; } + if (useWriteAPI()) { + long startTime = System.currentTimeMillis(); + StatementExecution.Builder writeApiExecutionBuilder = + StatementExecution.newBuilder() + .setStatementType( + com.google.cloud.bigquery.jdbc.telemetry.v1.StatementType.STATEMENT_TYPE_INSERT) + .setQueryApiType( + com.google.cloud.bigquery.jdbc.telemetry.v1.QueryApiType + .QUERY_API_TYPE_WRITE_API); try (BigQueryWriteClient writeClient = this.connection.getBigQueryWriteClient()) { LOG.info("Using Write API for bulk INSERT operation."); ArrayList currentParameterList = this.batchParameters.peek(); @@ -336,10 +348,24 @@ public int[] executeBatch() throws SQLException { long rowCount = bulkInsertWithWriteAPI(writeClient); int[] insertArray = new int[Math.toIntExact(rowCount)]; Arrays.fill(insertArray, 1); + + writeApiExecutionBuilder.setStatus( + com.google.cloud.bigquery.jdbc.telemetry.v1.Status.STATUS_SUCCESS); + return insertArray; } catch (DescriptorValidationException | IOException | InterruptedException e) { + writeApiExecutionBuilder + .setStatus(com.google.cloud.bigquery.jdbc.telemetry.v1.Status.STATUS_ERROR) + .setErrorCode(TelemetryManager.extractErrorCode(e)); + if (e instanceof InterruptedException) { + Thread.currentThread().interrupt(); + } throw new BigQueryJdbcRuntimeException("Failed to execute batch with Write API", e); + } finally { + long durationMs = System.currentTimeMillis() - startTime; + TelemetryManager.recordStatementExecution(writeApiExecutionBuilder, durationMs); + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_BATCH_OPERATIONS); } } else { @@ -370,6 +396,8 @@ public int[] executeBatch() throws SQLException { throw new BigQueryJdbcRuntimeException("Interrupted during individual INSERT batch", ex); } catch (SQLException e) { throw new BigQueryJdbcException("SQL error during individual INSERT batch", e); + } finally { + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_BATCH_OPERATIONS); } } } @@ -542,6 +570,11 @@ public ResultSetMetaData getMetaData() throws SQLException { if (this.insertSchema != null) { return BigQueryResultSetMetadata.of(this.insertSchema.getFields(), this); } + + TelemetryManager.recordFeatureUsage( + DriverFeature.DRIVER_FEATURE_METADATA_RETRIEVAL, + "DRIVER_FEATURE_RESULTSET_METADATA_RETRIEVAL"); + return null; } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java index 2ee127a26af1..603e444c1c3f 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/BigQueryStatement.java @@ -48,6 +48,10 @@ import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; import com.google.cloud.bigquery.exception.BigQueryJdbcSqlFeatureNotSupportedException; import com.google.cloud.bigquery.exception.BigQueryJdbcSqlSyntaxErrorException; +import com.google.cloud.bigquery.jdbc.telemetry.v1.DriverFeature; +import com.google.cloud.bigquery.jdbc.telemetry.v1.QueryApiType; +import com.google.cloud.bigquery.jdbc.telemetry.v1.StatementExecution; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import com.google.cloud.bigquery.storage.v1.ArrowRecordBatch; import com.google.cloud.bigquery.storage.v1.ArrowSchema; import com.google.cloud.bigquery.storage.v1.ArrowSerializationOptions; @@ -106,6 +110,9 @@ public class BigQueryStatement extends BigQueryNoOpsStatement { private static final String JDBC_JOB_PREFIX = "google-jdbc-"; private static final int MAX_RETRY_COUNT = 5; private static final long RETRY_DELAY_MS = 2000L; + // Reported when a failure carries no BigQuery or SQL error code, matching the fallback used by + // TelemetryManager#extractErrorCode. + private static final int UNKNOWN_ERROR_CODE = 1000; protected ResultSet currentResultSet; protected long currentUpdateCount = -1; protected List jobIds = new ArrayList<>(); @@ -149,6 +156,8 @@ public class BigQueryStatement extends BigQueryNoOpsStatement { private static final ThreadFactory JDBC_THREAD_FACTORY = new BigQueryThreadFactory("BigQuery-Thread-"); + protected StatementExecution.Builder currentExecutionBuilder = StatementExecution.newBuilder(); + static { BigQueryDaemonPollingTask.startGcDaemonTask( referenceQueueArrowRs, @@ -171,6 +180,8 @@ private void resetStatementFields() { this.parentJobId = null; this.currentJobIdIndex = -1; this.currentUpdateCount = -1; + + this.currentExecutionBuilder = StatementExecution.newBuilder(); } private BigQuerySettings generateBigQuerySettings() { @@ -487,6 +498,8 @@ public void cancel() throws SQLException { // If a ResultSet exists, then it will be closed as well, closing the // ownedThreads closeStatementResources(); + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_STATEMENT_CANCEL); } @Override @@ -574,6 +587,7 @@ ExecuteResult executeJob(QueryJobConfiguration jobConfiguration) if (result instanceof TableResult) { TableResult tableResult = (TableResult) result; saveSessionIdIfPresent(tableResult); + this.currentExecutionBuilder.setQueryApiType(QueryApiType.QUERY_API_TYPE_JOBLESS_QUERY); return new ExecuteResult(tableResult, null); } @@ -604,6 +618,7 @@ ExecuteResult executeJob(QueryJobConfiguration jobConfiguration) job = refreshedJob; } } + this.currentExecutionBuilder.setQueryApiType(QueryApiType.QUERY_API_TYPE_STANDARD_REST_API); return new ExecuteResult(tableResult, job); } @@ -655,19 +670,51 @@ void runQuery(String query, QueryJobConfiguration jobConfiguration) jobConfiguration.toBuilder().setJobTimeoutMs(Long.valueOf(queryTimeout) * 1000).build(); } + long startTime = System.currentTimeMillis(); + // Overwritten by the catch blocks below. The initial value covers failures that reach none of + // them, such as an uncaught RuntimeException, so a failure is never reported without a code. + int errorCode = UNKNOWN_ERROR_CODE; + try { resetStatementFields(); ExecuteResult executeResult = executeJob(jobConfiguration); StatementType statementType = getStatementType(executeResult); + + this.currentExecutionBuilder.setStatementType( + TelemetryManager.toStatementType(statementType)); + SqlType queryType = getQueryType(jobConfiguration, statementType); handleQueryResult(query, executeResult.tableResult, queryType, executeResult.job); + + this.currentExecutionBuilder.setStatus( + com.google.cloud.bigquery.jdbc.telemetry.v1.Status.STATUS_SUCCESS); + } catch (InterruptedException ex) { + errorCode = TelemetryManager.extractErrorCode(ex); + Thread.currentThread().interrupt(); throw new BigQueryJdbcRuntimeException("Interrupted during runQuery", ex); } catch (BigQueryException ex) { + errorCode = TelemetryManager.extractErrorCode(ex); if (ex.getMessage().contains("Syntax error")) { throw new BigQueryJdbcSqlSyntaxErrorException("BigQueryException during runQuery", ex); } throw new BigQueryJdbcException("BigQueryException during runQuery", ex); + } finally { + long durationMs = System.currentTimeMillis() - startTime; + + // Any path that did not reach STATUS_SUCCESS is a failure, whether it was caught above or + // propagated as an uncaught RuntimeException, so failures are never reported as UNSPECIFIED. + if (this.currentExecutionBuilder.getStatus() + == com.google.cloud.bigquery.jdbc.telemetry.v1.Status.STATUS_UNSPECIFIED) { + this.currentExecutionBuilder + .setStatus( + isCanceled + ? com.google.cloud.bigquery.jdbc.telemetry.v1.Status.STATUS_CANCELLED + : com.google.cloud.bigquery.jdbc.telemetry.v1.Status.STATUS_ERROR) + .setErrorCode(errorCode); + } + + TelemetryManager.recordStatementExecution(this.currentExecutionBuilder, durationMs); } } @@ -831,7 +878,6 @@ private QueryStatistics getQueryStatisticsFromJob(TableResult results, Job job) } private void updateAffectedRowCount(Long count) throws SQLException { - // TODO(neenu): check if this need to be closed vs removed) if (this.currentResultSet != null) { try { this.currentResultSet.close(); @@ -1077,6 +1123,7 @@ void processQueryResponse(String query, TableResult results, Job job) throws SQL try { LOG.info("Using ReadAPI to read the data."); resultSet = processArrowResultSet(results, job); + this.currentExecutionBuilder.setQueryApiType(QueryApiType.QUERY_API_TYPE_READ_API); } catch (SQLException e) { if (!isPermissionDeniedException(e)) { throw e; @@ -1088,6 +1135,13 @@ void processQueryResponse(String query, TableResult results, Job job) throws SQL if (resultSet == null) { LOG.info("Using Standard API to read the data."); resultSet = processJsonResultSet(results, job); + + // Jobless vs Standard REST + if (jobId == null) { + this.currentExecutionBuilder.setQueryApiType(QueryApiType.QUERY_API_TYPE_JOBLESS_QUERY); + } else { + this.currentExecutionBuilder.setQueryApiType(QueryApiType.QUERY_API_TYPE_STANDARD_REST_API); + } } this.currentResultSet = resultSet; this.currentUpdateCount = -1; @@ -1746,6 +1800,9 @@ private int[] executeBatchImpl(String combinedQueries) throws SQLException { } clearBatch(); + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_BATCH_OPERATIONS); + return result; } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/PooledConnectionDataSource.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/PooledConnectionDataSource.java index 7de4516427c3..0a986bd860b5 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/PooledConnectionDataSource.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/PooledConnectionDataSource.java @@ -17,6 +17,8 @@ package com.google.cloud.bigquery.jdbc; import com.google.cloud.bigquery.exception.BigQueryJdbcRuntimeException; +import com.google.cloud.bigquery.jdbc.telemetry.v1.DriverFeature; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import com.google.common.annotations.VisibleForTesting; import java.sql.Connection; import java.sql.SQLException; @@ -54,6 +56,8 @@ public PooledConnection getPooledConnection() throws SQLException { } BigQueryPooledConnection bqPooledConnection = new BigQueryPooledConnection(physicalConnection); bqPooledConnection.addConnectionEventListener(connectionPoolManager); + + TelemetryManager.recordFeatureUsage(DriverFeature.DRIVER_FEATURE_CONNECTION_POOLING); return bqPooledConnection; } diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryBatcher.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryBatcher.java index 7a9e804aa2be..39e28963cd02 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryBatcher.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryBatcher.java @@ -48,7 +48,7 @@ final class TelemetryBatcher implements AutoCloseable { private final boolean ownsExecutor; private final ReentrantLock flushLock = new ReentrantLock(); - // Live telemetry accumulators. Lock-free to eliminate object allocation and GC overhead. + // Live telemetry accumulator. Lock-free to eliminate object allocation and GC overhead. private ConcurrentHashMap metricsMap = new ConcurrentHashMap<>(); diff --git a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManager.java b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManager.java index 2837257b3534..8e742d934ec1 100644 --- a/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManager.java +++ b/java-bigquery-jdbc/src/main/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManager.java @@ -16,9 +16,11 @@ package com.google.cloud.bigquery.jdbc.telemetry.v1; +import com.google.cloud.bigquery.BigQueryException; import com.google.cloud.bigquery.JobStatistics.QueryStatistics; import com.google.cloud.bigquery.jdbc.BigQueryJdbcCustomLogger; import com.google.protobuf.Descriptors.EnumValueDescriptor; +import java.sql.SQLException; import java.util.Properties; import java.util.logging.Level; import java.util.logging.Logger; @@ -56,45 +58,48 @@ public static TelemetryManager getInstance(Properties properties) { return null; } - if (properties != null) { - TelemetryConfiguration configCheck = - TelemetryConfiguration.builder().resolveProperties(properties).build(); - if (!configCheck.isEnabled()) { - synchronized (TelemetryManager.class) { - globallyDisabled = true; - closeInstance(); - } - return null; + if (properties != null + && !TelemetryConfiguration.builder().resolveProperties(properties).build().isEnabled()) { + synchronized (TelemetryManager.class) { + globallyDisabled = true; + closeInstance(); } + return null; } TelemetryManager localRef = instance; - if (localRef == null) { - synchronized (TelemetryManager.class) { - if (globallyDisabled) { - return null; - } - localRef = instance; - if (localRef == null) { - TelemetryConfiguration config = - TelemetryConfiguration.builder().resolveProperties(properties).build(); - ClearcutTransport transport = new ClearcutTransport(config); - TelemetryBatcher batcher = new TelemetryBatcher(config, transport); - localRef = new TelemetryManager(batcher); - instance = localRef; - } + if (localRef != null) { + return localRef; + } + + synchronized (TelemetryManager.class) { + if (globallyDisabled) { + return null; } + localRef = instance; + if (localRef != null) { + return localRef; + } + TelemetryConfiguration config = + TelemetryConfiguration.builder().resolveProperties(properties).build(); + ClearcutTransport transport = new ClearcutTransport(config); + TelemetryBatcher batcher = new TelemetryBatcher(config, transport); + localRef = new TelemetryManager(batcher); + // Registered before the instance is published so that a non-null instance always implies a + // registered shutdown hook. + registerShutdownHook(); + instance = localRef; + return localRef; } - return localRef; } /** Package-private lifecycle initialisation method for explicit configuration or unit testing. */ static synchronized void init(TelemetryConfiguration config, ClearcutTransport transport) { closeInstance(); - if (config != null && config.isEnabled() && transport != null) { - TelemetryBatcher batcher = new TelemetryBatcher(config, transport); - instance = new TelemetryManager(batcher); + if (config == null || !config.isEnabled() || transport == null) { + return; } + instance = new TelemetryManager(new TelemetryBatcher(config, transport)); } /** @@ -105,6 +110,15 @@ TelemetryBatcher getBatcher() { return batcher; } + /** + * Returns the {@link TelemetryBatcher} of the active instance, or {@code null} if telemetry is + * closed or uninitialized. + */ + private static TelemetryBatcher activeBatcher() { + TelemetryManager localRef = instance; + return localRef == null ? null : localRef.getBatcher(); + } + /** * Executes a telemetry logging operation safely inside an exception-isolated block. Guaranteed to * catch all {@link Throwable} exceptions to protect JDBC driver operations. @@ -129,20 +143,19 @@ public static boolean isInitialized() { public static synchronized void closeInstance() { TelemetryManager localRef = instance; instance = null; - if (localRef != null) { - try { - localRef.close(); - } catch (Throwable t) { - logger.log(Level.FINE, "Error closing TelemetryManager instance", t); - } + if (localRef == null) { + return; + } + try { + localRef.close(); + } catch (Throwable t) { + logger.log(Level.FINE, "Error closing TelemetryManager instance", t); } } @Override public void close() { - if (batcher != null) { - batcher.close(); - } + batcher.close(); } // Package-private test helper to reset the global kill switch between test runs @@ -150,7 +163,7 @@ static synchronized void resetGlobalDisableForTest() { globallyDisabled = false; } - static StatementType toStatementType(QueryStatistics.StatementType bqStatementType) { + public static StatementType toStatementType(QueryStatistics.StatementType bqStatementType) { if (bqStatementType == null) { return StatementType.STATEMENT_TYPE_UNSPECIFIED; } @@ -161,40 +174,41 @@ static StatementType toStatementType(QueryStatistics.StatementType bqStatementTy return desc != null ? StatementType.valueOf(desc) : StatementType.STATEMENT_TYPE_OTHER; } - static AuthenticationType toAuthenticationType(int oauthType) { + public static AuthenticationType toAuthenticationType(int oauthType) { switch (oauthType) { case 0: return AuthenticationType.AUTHENTICATION_TYPE_SERVICE_ACCOUNT; case 1: return AuthenticationType.AUTHENTICATION_TYPE_USER_AUTHENTICATION; case 2: - return AuthenticationType.AUTHENTICATION_TYPE_APPLICATION_DEFAULT_CREDENTIALS; + return AuthenticationType.AUTHENTICATION_TYPE_TOKEN; case 3: - return AuthenticationType.AUTHENTICATION_TYPE_EXTERNAL; + return AuthenticationType.AUTHENTICATION_TYPE_APPLICATION_DEFAULT_CREDENTIALS; case 4: - return AuthenticationType.AUTHENTICATION_TYPE_TOKEN; + return AuthenticationType.AUTHENTICATION_TYPE_EXTERNAL; default: return AuthenticationType.AUTHENTICATION_TYPE_CUSTOM; } } - static void recordConnectionAttempt(Status status, int errorCode, AuthenticationType authType) { + public static void recordConnectionAttempt( + Status status, int errorCode, AuthenticationType authType) { runSafely( () -> { - TelemetryManager mgr = instance; - if (mgr != null && mgr.getBatcher() != null) { - mgr.getBatcher() - .offer( - ConnectionAttempt.newBuilder() - .setStatus(status) - .setErrorCode(errorCode) - .setAuthType(authType) - .build()); + TelemetryBatcher activeBatcher = activeBatcher(); + if (activeBatcher == null) { + return; } + activeBatcher.offer( + ConnectionAttempt.newBuilder() + .setStatus(status) + .setErrorCode(errorCode) + .setAuthType(authType) + .build()); }); } - static void recordStatementExecution( + public static void recordStatementExecution( StatementType statementType, QueryApiType apiType, Status status, @@ -202,33 +216,123 @@ static void recordStatementExecution( long durationMs) { runSafely( () -> { - TelemetryManager mgr = instance; - if (mgr != null && mgr.getBatcher() != null) { - mgr.getBatcher() - .offer( - StatementExecution.newBuilder() - .setStatementType(statementType) - .setQueryApiType(apiType) - .setStatus(status) - .setErrorCode(errorCode) - .build(), - durationMs); + TelemetryBatcher activeBatcher = activeBatcher(); + if (activeBatcher == null) { + return; } + activeBatcher.offer( + StatementExecution.newBuilder() + .setStatementType(statementType) + .setQueryApiType(apiType) + .setStatus(status) + .setErrorCode(errorCode) + .build(), + durationMs); }); } - static void recordFeatureUsage(DriverFeature feature, String customFeatureName) { + public static void recordStatementExecution( + StatementExecution.Builder statementExecutionBuilder, long durationMs) { + if (statementExecutionBuilder == null) { + return; + } + runSafely( + () -> { + TelemetryBatcher activeBatcher = activeBatcher(); + if (activeBatcher == null) { + return; + } + activeBatcher.offer(statementExecutionBuilder.build(), durationMs); + }); + } + + public static void recordFeatureUsage(DriverFeature feature, String customFeatureName) { + runSafely( + () -> { + TelemetryBatcher activeBatcher = activeBatcher(); + if (activeBatcher == null) { + return; + } + activeBatcher.offer( + FeatureUsage.newBuilder() + .setDriverFeature(feature) + .setCustomFeatureName(customFeatureName == null ? "" : customFeatureName) + .build()); + }); + } + + public static void recordFeatureUsage(DriverFeature feature) { + recordFeatureUsage(feature, null); + } + + public static void recordError(int errorCode, int errorXdbcCode, String methodName) { runSafely( () -> { - TelemetryManager mgr = instance; - if (mgr != null && mgr.getBatcher() != null) { - mgr.getBatcher() - .offer( - FeatureUsage.newBuilder() - .setDriverFeature(feature) - .setCustomFeatureName(customFeatureName == null ? "" : customFeatureName) - .build()); + TelemetryBatcher activeBatcher = activeBatcher(); + if (activeBatcher == null) { + return; } + activeBatcher.offer( + ErrorMetric.newBuilder() + .setErrorCode(errorCode) + .setErrorXdbcCode(errorXdbcCode) + .setMethodName(methodName == null ? "" : methodName) + .build()); }); } + + /** + * Extracts the numeric error code from the throwable chain. Traverses causes to unpack + * BigQueryException (HTTP status codes) or SQLException error codes. Returns 1000 as the fallback + * driver error code. + */ + public static int extractErrorCode(Throwable t) { + int depth = 0; + while (t != null && depth++ < 20) { + if (t instanceof BigQueryException) { + int code = ((BigQueryException) t).getCode(); + if (code != 0) { + return code; + } + } + if (t instanceof SQLException) { + int code = ((SQLException) t).getErrorCode(); + if (code != 0) { + return code; + } + } + t = t.getCause(); + } + return 1000; + } + + /** + * Registers the JVM shutdown hook that flushes pending telemetry. + * + *

Only reached from the instance-creation critical section of {@link + * #getInstance(Properties)}, which already guarantees single execution, so no further guarding is + * needed here. + */ + private static void registerShutdownHook() { + try { + Runtime.getRuntime() + .addShutdownHook( + new Thread( + TelemetryManager::closeInstanceOnShutdown, + "bigquery-jdbc-telemetry-shutdown-hook")); + } catch (IllegalStateException e) { + // Thrown if the JVM is already in the process of shutting down + } catch (SecurityException e) { + logger.warning("SecurityManager prevented registering telemetry shutdown hook"); + } + } + + /** Shutdown-hook body: flushes and closes the shared instance without propagating failures. */ + private static void closeInstanceOnShutdown() { + try { + closeInstance(); + } catch (Throwable t) { + logger.warning("Error closing TelemetryManager during JVM shutdown"); + } + } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryDriverTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryDriverTest.java index 8acbc5abb8dc..9f0b07ace74d 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryDriverTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/BigQueryDriverTest.java @@ -18,6 +18,7 @@ import static com.google.common.truth.Truth.assertThat; import static org.mockito.Mockito.mock; +import com.google.cloud.bigquery.jdbc.telemetry.v1.TelemetryManager; import com.google.cloud.bigquery.jdbc.utils.BigQueryJdbcVersionUtility; import io.opentelemetry.api.OpenTelemetry; import java.sql.Connection; @@ -187,4 +188,48 @@ public void testInvalidLogLevelExceptionIsLogged() { && r.getMessage().contains("Failed to parse connection URL properties")); assertThat(foundSevere).isTrue(); } + + @Test + public void testConnect_recordsSuccessfulConnectionTelemetry() throws SQLException { + TelemetryManager.closeInstance(); + Connection connection = + bigQueryDriver.connect( + "jdbc:bigquery://https://www.googleapis.com/bigquery/v2:443;" + + "OAuthType=2;ProjectId=MyBigQueryProject;" + + "OAuthAccessToken=redactedToken;OAuthClientId=redactedToken;" + + "OAuthClientSecret=redactedToken;", + new Properties()); + assertThat(connection).isNotNull(); + assertThat(connection.isClosed()).isFalse(); + // Verify TelemetryManager is initialized and recorded the connection + assertThat(TelemetryManager.isInitialized()).isTrue(); + } + + @Test + public void testConnect_recordsFailedConnectionTelemetry() { + TelemetryManager.closeInstance(); + // Malformed URL causing DataSource parsing failure + Assertions.assertThrows( + SQLException.class, + () -> + bigQueryDriver.connect( + "jdbc:bigquery://https://www.googleapis.com/bigquery/v2:443;OAuthType=invalid;", + new Properties())); + assertThat(TelemetryManager.isInitialized()).isTrue(); + } + + @Test + public void testConnect_optOut_noTelemetryRecorded() throws SQLException { + TelemetryManager.closeInstance(); + Connection connection = + bigQueryDriver.connect( + "jdbc:bigquery://https://www.googleapis.com/bigquery/v2:443;" + + "OAuthType=2;ProjectId=MyBigQueryProject;" + + "OAuthAccessToken=redactedToken;OAuthClientId=redactedToken;" + + "OAuthClientSecret=redactedToken;EnableDiagnosticTelemetry=0;", + new Properties()); + assertThat(connection).isNotNull(); + // Since opt-out was requested, TelemetryManager should NOT be initialized + assertThat(TelemetryManager.isInitialized()).isFalse(); + } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManagerTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManagerTest.java index 007236f25962..cf467afaf319 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManagerTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/telemetry/v1/TelemetryManagerTest.java @@ -162,12 +162,12 @@ public void testToAuthenticationType() { AuthenticationType.AUTHENTICATION_TYPE_USER_AUTHENTICATION, TelemetryManager.toAuthenticationType(1)); assertEquals( - AuthenticationType.AUTHENTICATION_TYPE_APPLICATION_DEFAULT_CREDENTIALS, - TelemetryManager.toAuthenticationType(2)); + AuthenticationType.AUTHENTICATION_TYPE_TOKEN, TelemetryManager.toAuthenticationType(2)); assertEquals( - AuthenticationType.AUTHENTICATION_TYPE_EXTERNAL, TelemetryManager.toAuthenticationType(3)); + AuthenticationType.AUTHENTICATION_TYPE_APPLICATION_DEFAULT_CREDENTIALS, + TelemetryManager.toAuthenticationType(3)); assertEquals( - AuthenticationType.AUTHENTICATION_TYPE_TOKEN, TelemetryManager.toAuthenticationType(4)); + AuthenticationType.AUTHENTICATION_TYPE_EXTERNAL, TelemetryManager.toAuthenticationType(4)); assertEquals( AuthenticationType.AUTHENTICATION_TYPE_CUSTOM, TelemetryManager.toAuthenticationType(5));