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();
+ }
+}