From f42828d46e2777b6aa4f56d46197781b2aa7e236 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Fri, 4 Sep 2026 14:56:08 -0400 Subject: [PATCH 1/4] fix(bigquery-jdbc): abort session when connection is closed --- .../bigquery/jdbc/BigQueryConnection.java | 36 +++++++++++++++++ .../bigquery/jdbc/BigQueryConnectionTest.java | 40 ++++++++++++++++++- .../bigquery/jdbc/it/ITBigQueryJDBCTest.java | 30 ++++++++++++++ 3 files changed, 105 insertions(+), 1 deletion(-) 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 2f5863054903..7adc6d1afd46 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 @@ -1073,6 +1073,10 @@ private void closeImpl() throws SQLException { } } + if (this.sessionInfoConnectionProperty != null) { + abortSession(); + } + boolean interrupted = Thread.currentThread().isInterrupted(); try { @@ -1467,6 +1471,38 @@ private void commitTransaction() { } } + private void abortSession() { + try { + LOG.fine( + "Aborting session on connection close: " + this.sessionInfoConnectionProperty.getValue()); + QueryJobConfiguration abortSessionJobConfig = + QueryJobConfiguration.newBuilder("CALL BQ.ABORT_SESSION();") + .setConnectionProperties(this.queryProperties) + .build(); + Job abortJob = this.bigQuery.create(JobInfo.of(abortSessionJobConfig)); + abortJob.waitFor(); + } catch (InterruptedException ex) { + Thread.currentThread().interrupt(); + throw new BigQueryJdbcRuntimeException("Interrupted during close", ex); + } catch (BigQueryException ex) { + LOG.warning( + "Failed to abort session during connection close (session may have already ended): " + + ex.getMessage()); + } finally { + this.sessionInfoConnectionProperty = null; + if (this.queryProperties != null) { + List updated = new ArrayList<>(); + for (ConnectionProperty cp : this.queryProperties) { + if (!"session_id".equalsIgnoreCase(cp.getKey())) { + updated.add(cp); + } + } + this.queryProperties = Collections.unmodifiableList(updated); + } + this.transactionStarted = false; + } + } + @Override public CallableStatement prepareCall(String sql) throws SQLException { checkClosed(); 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 7b01f9ac760e..f206dc8a7651 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 @@ -16,7 +16,6 @@ package com.google.cloud.bigquery.jdbc; -import static org.junit.jupiter.api.Assertions.*; import static org.junit.jupiter.api.Assertions.assertEquals; import static org.junit.jupiter.api.Assertions.assertFalse; import static org.junit.jupiter.api.Assertions.assertNotNull; @@ -44,7 +43,10 @@ import com.google.auth.oauth2.GoogleCredentials; import com.google.cloud.bigquery.BigQuery; import com.google.cloud.bigquery.BigQueryException; +import com.google.cloud.bigquery.Job; +import com.google.cloud.bigquery.JobInfo; import com.google.cloud.bigquery.Project; +import com.google.cloud.bigquery.QueryJobConfiguration; import com.google.cloud.bigquery.QueryJobConfiguration.JobCreationMode; import com.google.cloud.bigquery.exception.BigQueryJdbcException; import com.google.cloud.bigquery.storage.v1.BigQueryReadClient; @@ -70,6 +72,7 @@ import org.junit.jupiter.api.extension.RegisterExtension; import org.junit.jupiter.params.ParameterizedTest; import org.junit.jupiter.params.provider.CsvSource; +import org.mockito.ArgumentCaptor; import org.mockito.MockedStatic; public class BigQueryConnectionTest extends BigQueryJdbcLoggingBaseTest { @@ -819,4 +822,39 @@ public void testUserSuppliedSessionId() throws Exception { "user_supplied_session_999", connection.getSessionInfoConnectionProperty().getValue()); } } + + @Test + public void testCloseWithActiveSessionAbortsSession() throws Exception { + try (BigQueryConnection connection = new BigQueryConnection(BASE_URL)) { + BigQuery mockBigQuery = mock(BigQuery.class); + Job mockJob = mock(Job.class); + when(mockBigQuery.create(any(JobInfo.class))).thenReturn(mockJob); + when(mockJob.waitFor()).thenReturn(mockJob); + connection.bigQuery = mockBigQuery; + + connection.updateSessionInfo("test_session_id_to_abort"); + connection.close(); + + ArgumentCaptor jobCaptor = ArgumentCaptor.forClass(JobInfo.class); + verify(mockBigQuery).create(jobCaptor.capture()); + QueryJobConfiguration config = + (QueryJobConfiguration) jobCaptor.getValue().getConfiguration(); + assertEquals("CALL BQ.ABORT_SESSION();", config.getQuery()); + assertNull(connection.getSessionInfoConnectionProperty()); + assertTrue(connection.isClosed()); + } + } + + @Test + public void testCloseWithoutSessionDoesNotAbortSession() throws Exception { + try (BigQueryConnection connection = new BigQueryConnection(BASE_URL)) { + BigQuery mockBigQuery = mock(BigQuery.class); + connection.bigQuery = mockBigQuery; + + connection.close(); + + verify(mockBigQuery, never()).create(any(JobInfo.class)); + assertTrue(connection.isClosed()); + } + } } diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java index 81d172a70018..8e1ea9686d5b 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java @@ -2880,4 +2880,34 @@ public void testPerConnectionLoggingE2E() throws SQLException, IOException { } } } + + @Test + public void testSessionAbortedOnConnectionClose() throws SQLException { + String sessionId; + try (Connection connection = DriverManager.getConnection(session_enabled_connection_uri)) { + try (Statement statement = connection.createStatement()) { + statement.execute("CREATE TEMP TABLE session_temp_table (id INT64);"); + } + BigQueryConnection bqConn = connection.unwrap(BigQueryConnection.class); + assertNotNull(bqConn.getSessionInfoConnectionProperty()); + sessionId = bqConn.getSessionInfoConnectionProperty().getValue(); + assertNotNull(sessionId); + } + + // After connection is closed, the session is aborted on the BigQuery server. + // Attaching to the same session_id in a new connection should fail when running a query. + String urlWithAbortedSession = + connection_uri + "EnableSession=1;QueryProperties=session_id=" + sessionId + ";"; + try (Connection newConnection = DriverManager.getConnection(urlWithAbortedSession)) { + try (Statement statement = newConnection.createStatement()) { + SQLException ex = + assertThrows( + SQLException.class, () -> statement.execute("SELECT * FROM session_temp_table;")); + assertTrue( + ex.getMessage().toLowerCase().contains("session ended") + || ex.getMessage().toLowerCase().contains("not found"), + "Expected session ended error but got: " + ex.getMessage()); + } + } + } } From 488ea7020707304c8b4aae7ae2f7bf40b458ac3b Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Fri, 4 Sep 2026 15:08:04 -0400 Subject: [PATCH 2/4] nit --- .../com/google/cloud/bigquery/jdbc/BigQueryConnection.java | 4 ++-- 1 file changed, 2 insertions(+), 2 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 7adc6d1afd46..c39665fcd8f3 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 @@ -1483,10 +1483,10 @@ private void abortSession() { abortJob.waitFor(); } catch (InterruptedException ex) { Thread.currentThread().interrupt(); - throw new BigQueryJdbcRuntimeException("Interrupted during close", ex); + throw new BigQueryJdbcRuntimeException("Interrupted during session abort", ex); } catch (BigQueryException ex) { LOG.warning( - "Failed to abort session during connection close (session may have already ended): " + "Failed to abort session during session abort (session may have already ended): " + ex.getMessage()); } finally { this.sessionInfoConnectionProperty = null; From 4840c9ed9c0bf11af469cc3ee0d18e50fccebb50 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Fri, 4 Sep 2026 15:53:26 -0400 Subject: [PATCH 3/4] fix test --- .../google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java | 8 +++++--- 1 file changed, 5 insertions(+), 3 deletions(-) diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java index 8e1ea9686d5b..282a776121b3 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java @@ -2904,9 +2904,11 @@ public void testSessionAbortedOnConnectionClose() throws SQLException { assertThrows( SQLException.class, () -> statement.execute("SELECT * FROM session_temp_table;")); assertTrue( - ex.getMessage().toLowerCase().contains("session ended") - || ex.getMessage().toLowerCase().contains("not found"), - "Expected session ended error but got: " + ex.getMessage()); + ex.getMessage().contains(sessionId), + "Expected exception message to not contain session ID: " + + sessionId + + ", but got: " + + ex.getMessage()); } } } From f12d6f4eee1e1312ebab60ab981b88b3fa0363e6 Mon Sep 17 00:00:00 2001 From: Neenu1995 Date: Fri, 4 Sep 2026 20:54:59 -0400 Subject: [PATCH 4/4] nit --- .../google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java | 7 +------ 1 file changed, 1 insertion(+), 6 deletions(-) diff --git a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java index 282a776121b3..625836aa3302 100644 --- a/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java +++ b/java-bigquery-jdbc/src/test/java/com/google/cloud/bigquery/jdbc/it/ITBigQueryJDBCTest.java @@ -2903,12 +2903,7 @@ public void testSessionAbortedOnConnectionClose() throws SQLException { SQLException ex = assertThrows( SQLException.class, () -> statement.execute("SELECT * FROM session_temp_table;")); - assertTrue( - ex.getMessage().contains(sessionId), - "Expected exception message to not contain session ID: " - + sessionId - + ", but got: " - + ex.getMessage()); + assertTrue(ex.getMessage().toLowerCase().contains("not found".toLowerCase())); } } }