This is an automated email from the ASF dual-hosted git repository.

davidzollo pushed a commit to branch dev
in repository https://gitbox.apache.org/repos/asf/seatunnel.git


The following commit(s) were added to refs/heads/dev by this push:
     new a1324acdae [Fix][Connector-V2] Avoid Xugu pooled connection isValid 
checks (#11190)
a1324acdae is described below

commit a1324acdae34bdb6bf8e524afb40a5c19691001d
Author: Daniel <[email protected]>
AuthorDate: Mon Jun 29 22:45:04 2026 +0800

    [Fix][Connector-V2] Avoid Xugu pooled connection isValid checks (#11190)
---
 .../connection/JdbcConnectionValidationUtils.java  | 78 ++++++++++++++++++++++
 .../SimpleJdbcConnectionPoolProviderProxy.java     |  5 +-
 .../connection/SimpleJdbcConnectionProvider.java   |  3 +-
 .../seatunnel/jdbc/sink/JdbcSinkWriter.java        | 13 ++++
 .../JdbcConnectionValidationUtilsTest.java         | 78 ++++++++++++++++++++++
 .../seatunnel/jdbc/sink/JdbcSinkWriterTest.java    | 64 ++++++++++++++++++
 6 files changed, 236 insertions(+), 5 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/JdbcConnectionValidationUtils.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/JdbcConnectionValidationUtils.java
new file mode 100644
index 0000000000..15736a8a25
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/JdbcConnectionValidationUtils.java
@@ -0,0 +1,78 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection;
+
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcConnectionConfig;
+
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.util.Optional;
+
+/** Utility methods for JDBC driver-specific connection validation hooks. */
+public final class JdbcConnectionValidationUtils {
+
+    /** Xugu JDBC driver class name used by connector-jdbc ITs. */
+    public static final String XUGU_DRIVER = "com.xugu.cloudjdbc.Driver";
+
+    /** Validation query used when the Xugu driver cannot answer 
Connection.isValid(timeout). */
+    public static final String XUGU_VALIDATION_QUERY = "SELECT 1 FROM DUAL";
+
+    private JdbcConnectionValidationUtils() {}
+
+    /**
+     * Xugu's driver throws during Connection.isValid(timeout), so pooled 
connections need a SQL
+     * probe instead of the JDBC driver validation hook.
+     */
+    public static boolean isConnectionValid(Connection connection, 
JdbcConnectionConfig jdbcConfig)
+            throws SQLException {
+        if (connection == null) {
+            return false;
+        }
+
+        Optional<String> validationQuery = 
getConnectionValidationQuery(jdbcConfig);
+        if (!validationQuery.isPresent()) {
+            return 
connection.isValid(jdbcConfig.getConnectionCheckTimeoutSeconds());
+        }
+
+        try (PreparedStatement preparedStatement =
+                        connection.prepareStatement(validationQuery.get());
+                ResultSet resultSet = preparedStatement.executeQuery()) {
+            return resultSet.next();
+        }
+    }
+
+    /**
+     * Returns an optional validation query for drivers that need SQL-based 
liveness checks instead
+     * of {@link Connection#isValid(int)}.
+     */
+    public static Optional<String> 
getConnectionValidationQuery(JdbcConnectionConfig jdbcConfig) {
+        if (jdbcConfig == null) {
+            return Optional.empty();
+        }
+
+        String driverName = jdbcConfig.getDriverName();
+        String url = jdbcConfig.getUrl();
+        if (XUGU_DRIVER.equals(driverName) || (url != null && 
url.startsWith("jdbc:xugu:"))) {
+            return Optional.of(XUGU_VALIDATION_QUERY);
+        }
+
+        return Optional.empty();
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionPoolProviderProxy.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionPoolProviderProxy.java
index 6bca2db726..6a3f024e56 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionPoolProviderProxy.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionPoolProviderProxy.java
@@ -47,9 +47,8 @@ public class SimpleJdbcConnectionPoolProviderProxy implements 
JdbcConnectionProv
     @Override
     public boolean isConnectionValid() throws SQLException {
         return poolManager.containsConnection(queueIndex)
-                && poolManager
-                        .getConnection(queueIndex)
-                        
.isValid(jdbcConfig.getConnectionCheckTimeoutSeconds());
+                && JdbcConnectionValidationUtils.isConnectionValid(
+                        poolManager.getConnection(queueIndex), jdbcConfig);
     }
 
     @Override
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionProvider.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionProvider.java
index f9f36325af..267b26c713 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionProvider.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/SimpleJdbcConnectionProvider.java
@@ -59,8 +59,7 @@ public class SimpleJdbcConnectionProvider implements 
JdbcConnectionProvider, Ser
 
     @Override
     public boolean isConnectionValid() throws SQLException {
-        return connection != null
-                && 
connection.isValid(jdbcConfig.getConnectionCheckTimeoutSeconds());
+        return JdbcConnectionValidationUtils.isConnectionValid(connection, 
jdbcConfig);
     }
 
     private static Driver loadDriver(String driverName) throws 
ClassNotFoundException {
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkWriter.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkWriter.java
index 3518015b1a..083a2ee85a 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkWriter.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkWriter.java
@@ -24,10 +24,12 @@ import org.apache.seatunnel.api.table.catalog.TablePath;
 import org.apache.seatunnel.api.table.catalog.TableSchema;
 import org.apache.seatunnel.api.table.type.SeaTunnelRow;
 import org.apache.seatunnel.common.exception.CommonErrorCodeDeprecated;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcConnectionConfig;
 import org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcSinkConfig;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorErrorCode;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.JdbcOutputFormatBuilder;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection.JdbcConnectionValidationUtils;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection.SimpleJdbcConnectionPoolProviderProxy;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.DatabaseIdentifier;
 import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
@@ -95,10 +97,21 @@ public class JdbcSinkWriter extends 
AbstractJdbcSinkWriter<ConnectionPoolManager
             
ds.setPassword(jdbcSinkConfig.getJdbcConnectionConfig().getPassword().get());
         }
         
ds.setAutoCommit(jdbcSinkConfig.getJdbcConnectionConfig().isAutoCommit());
+        applyConnectionValidation(ds, 
jdbcSinkConfig.getJdbcConnectionConfig());
         
jdbcSinkConfig.getJdbcConnectionConfig().getProperties().forEach(ds::addDataSourceProperty);
         return new JdbcMultiTableResourceManager(new 
ConnectionPoolManager(ds));
     }
 
+    /**
+     * Configures pool-level validation for JDBC drivers that cannot pass 
Hikari's default
+     * Connection.isValid(timeout) probe.
+     */
+    static void applyConnectionValidation(
+            HikariDataSource dataSource, JdbcConnectionConfig 
jdbcConnectionConfig) {
+        
JdbcConnectionValidationUtils.getConnectionValidationQuery(jdbcConnectionConfig)
+                .ifPresent(dataSource::setConnectionTestQuery);
+    }
+
     @Override
     public void setMultiTableResourceManager(
             MultiTableResourceManager<ConnectionPoolManager> 
multiTableResourceManager,
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/JdbcConnectionValidationUtilsTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/JdbcConnectionValidationUtilsTest.java
new file mode 100644
index 0000000000..7946142b99
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/connection/JdbcConnectionValidationUtilsTest.java
@@ -0,0 +1,78 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection;
+
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcConnectionConfig;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+
+import static org.mockito.ArgumentMatchers.anyInt;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+/** Tests driver-specific JDBC connection validation fallbacks. */
+class JdbcConnectionValidationUtilsTest {
+
+    /** Verifies that Xugu uses an explicit SQL probe instead of 
Connection.isValid(timeout). */
+    @Test
+    void testXuguValidationUsesSqlProbe() throws SQLException {
+        JdbcConnectionConfig jdbcConnectionConfig =
+                JdbcConnectionConfig.builder()
+                        .driverName(JdbcConnectionValidationUtils.XUGU_DRIVER)
+                        .url("jdbc:xugu://localhost:5138/SYSTEM")
+                        .build();
+        Connection connection = mock(Connection.class);
+        PreparedStatement preparedStatement = mock(PreparedStatement.class);
+        ResultSet resultSet = mock(ResultSet.class);
+        
when(connection.prepareStatement(JdbcConnectionValidationUtils.XUGU_VALIDATION_QUERY))
+                .thenReturn(preparedStatement);
+        when(preparedStatement.executeQuery()).thenReturn(resultSet);
+        when(resultSet.next()).thenReturn(true);
+
+        Assertions.assertTrue(
+                JdbcConnectionValidationUtils.isConnectionValid(connection, 
jdbcConnectionConfig));
+        verify(connection, never()).isValid(anyInt());
+    }
+
+    /** Verifies that default drivers still use the standard JDBC validation 
hook. */
+    @Test
+    void testDefaultValidationFallsBackToJdbcIsValid() throws SQLException {
+        JdbcConnectionConfig jdbcConnectionConfig =
+                JdbcConnectionConfig.builder()
+                        .driverName("org.postgresql.Driver")
+                        .url("jdbc:postgresql://localhost:5432/test")
+                        .connectionCheckTimeoutSeconds(12)
+                        .build();
+        Connection connection = mock(Connection.class);
+        when(connection.isValid(12)).thenReturn(true);
+
+        Assertions.assertTrue(
+                JdbcConnectionValidationUtils.isConnectionValid(connection, 
jdbcConnectionConfig));
+        verify(connection).isValid(12);
+        verify(connection, never())
+                
.prepareStatement(JdbcConnectionValidationUtils.XUGU_VALIDATION_QUERY);
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkWriterTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkWriterTest.java
new file mode 100644
index 0000000000..442b9f9493
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/sink/JdbcSinkWriterTest.java
@@ -0,0 +1,64 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one or more
+ * contributor license agreements.  See the NOTICE file distributed with
+ * this work for additional information regarding copyright ownership.
+ * The ASF licenses this file to You under the Apache License, Version 2.0
+ * (the "License"); you may not use this file except in compliance with
+ * the License.  You may obtain a copy of the License at
+ *
+ *    http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing, software
+ * distributed under the License is distributed on an "AS IS" BASIS,
+ * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+ * See the License for the specific language governing permissions and
+ * limitations under the License.
+ */
+
+package org.apache.seatunnel.connectors.seatunnel.jdbc.sink;
+
+import org.apache.seatunnel.shade.com.zaxxer.hikari.HikariDataSource;
+
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.config.JdbcConnectionConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.connection.JdbcConnectionValidationUtils;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+/** Tests JDBC sink connection pool validation query customization. */
+class JdbcSinkWriterTest {
+
+    /** Verifies that Xugu pools use a validation query compatible with the 
driver. */
+    @Test
+    void testApplyConnectionValidationSetsXuguValidationQuery() {
+        HikariDataSource dataSource = new HikariDataSource();
+        JdbcConnectionConfig jdbcConnectionConfig =
+                JdbcConnectionConfig.builder()
+                        .driverName(JdbcConnectionValidationUtils.XUGU_DRIVER)
+                        .url("jdbc:xugu://localhost:5138/SYSTEM")
+                        .build();
+
+        JdbcSinkWriter.applyConnectionValidation(dataSource, 
jdbcConnectionConfig);
+
+        Assertions.assertEquals(
+                JdbcConnectionValidationUtils.XUGU_VALIDATION_QUERY,
+                dataSource.getConnectionTestQuery());
+        dataSource.close();
+    }
+
+    /** Verifies that other drivers keep Hikari's default validation behavior. 
*/
+    @Test
+    void testApplyConnectionValidationKeepsDefaultDriverValidation() {
+        HikariDataSource dataSource = new HikariDataSource();
+        JdbcConnectionConfig jdbcConnectionConfig =
+                JdbcConnectionConfig.builder()
+                        .driverName("org.postgresql.Driver")
+                        .url("jdbc:postgresql://localhost:5432/test")
+                        .build();
+
+        JdbcSinkWriter.applyConnectionValidation(dataSource, 
jdbcConnectionConfig);
+
+        Assertions.assertNull(dataSource.getConnectionTestQuery());
+        dataSource.close();
+    }
+}

Reply via email to