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 b1e21631ff [Fix][Connector-V2] Close JDBC source read transactions 
after each split (#11770)
b1e21631ff is described below

commit b1e21631ffdfd8cb33e6adca618b2d95224a9d66
Author: Yu, Wang <[email protected]>
AuthorDate: Fri Aug 14 08:09:07 2026 -0700

    [Fix][Connector-V2] Close JDBC source read transactions after each split 
(#11770)
    
    Co-authored-by: Brownie Woom 
<[email protected]>
---
 .../seatunnel/jdbc/internal/JdbcInputFormat.java   | 133 +++++++++-
 .../jdbc/internal/JdbcInputFormatTest.java         | 272 +++++++++++++++++++++
 2 files changed, 398 insertions(+), 7 deletions(-)

diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/JdbcInputFormat.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/JdbcInputFormat.java
index 90dbe021fc..b837827e5d 100644
--- 
a/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/JdbcInputFormat.java
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/main/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/JdbcInputFormat.java
@@ -38,6 +38,7 @@ import org.slf4j.LoggerFactory;
 
 import java.io.IOException;
 import java.io.Serializable;
+import java.sql.Connection;
 import java.sql.PreparedStatement;
 import java.sql.ResultSet;
 import java.sql.SQLException;
@@ -57,6 +58,7 @@ public class JdbcInputFormat implements Serializable {
     private final JdbcRowConverter jdbcRowConverter;
     private final Map<TablePath, CatalogTable> tables;
     private final ChunkSplitter chunkSplitter;
+    private final boolean configuredAutoCommit;
 
     private transient String splitTableId;
     private transient TableSchema splitTableSchema;
@@ -65,23 +67,37 @@ public class JdbcInputFormat implements Serializable {
     private volatile boolean hasNext;
 
     public JdbcInputFormat(JdbcSourceConfig config, Map<TablePath, 
CatalogTable> tables) {
-        this.jdbcDialect =
+        this(
                 JdbcDialectLoader.load(
                         config.getJdbcConnectionConfig().getUrl(),
                         config.getJdbcConnectionConfig().getDialect(),
-                        config.getCompatibleMode());
-        this.chunkSplitter = ChunkSplitter.create(config);
+                        config.getCompatibleMode()),
+                ChunkSplitter.create(config),
+                tables,
+                config.getJdbcConnectionConfig().isAutoCommit());
+    }
+
+    JdbcInputFormat(
+            JdbcDialect jdbcDialect,
+            ChunkSplitter chunkSplitter,
+            Map<TablePath, CatalogTable> tables,
+            boolean configuredAutoCommit) {
+        this.jdbcDialect = jdbcDialect;
+        this.chunkSplitter = chunkSplitter;
         this.jdbcRowConverter = jdbcDialect.getRowConverter();
         this.tables = tables;
+        this.configuredAutoCommit = configuredAutoCommit;
     }
 
     public void openInputFormat() {}
 
     public void closeInputFormat() throws IOException {
-        close();
-
-        if (chunkSplitter != null) {
-            chunkSplitter.close();
+        try {
+            close();
+        } finally {
+            if (chunkSplitter != null) {
+                chunkSplitter.close();
+            }
         }
     }
 
@@ -101,10 +117,31 @@ public class JdbcInputFormat implements Serializable {
             resultSet = statement.executeQuery();
             hasNext = resultSet.next();
         } catch (SQLException se) {
+            cleanupAfterOpenFailure(se);
             throw new JdbcConnectorException(
                     JdbcConnectorErrorCode.CONNECT_DATABASE_FAILED,
                     "open() failed." + se.getMessage(),
                     se);
+        } catch (RuntimeException runtimeException) {
+            cleanupAfterOpenFailure(runtimeException);
+            throw runtimeException;
+        }
+    }
+
+    private void cleanupAfterOpenFailure(Throwable openException) {
+        boolean shouldDiscardConnection = statement == null;
+        try {
+            close();
+        } catch (IOException cleanupException) {
+            openException.addSuppressed(cleanupException);
+            shouldDiscardConnection = true;
+        } finally {
+            if (shouldDiscardConnection) {
+                // Statement creation may establish or mutate the cached 
connection before failing
+                // without returning a statement. Discard it because close() 
cannot identify and
+                // finish that transaction safely.
+                chunkSplitter.close();
+            }
         }
     }
 
@@ -114,11 +151,14 @@ public class JdbcInputFormat implements Serializable {
      * @throws IOException Indicates that a resource could not be closed.
      */
     public void close() throws IOException {
+        Connection connection = getStatementConnection();
         if (resultSet != null) {
             try {
                 resultSet.close();
             } catch (SQLException e) {
                 LOG.info("ResultSet couldn't be closed - " + e.getMessage());
+            } finally {
+                resultSet = null;
             }
         }
         if (statement != null) {
@@ -126,7 +166,86 @@ public class JdbcInputFormat implements Serializable {
                 statement.close();
             } catch (SQLException e) {
                 LOG.info("Statement couldn't be closed - " + e.getMessage());
+            } finally {
+                statement = null;
+            }
+        }
+
+        hasNext = false;
+        splitTableSchema = null;
+        splitTableId = null;
+        finishReadTransaction(connection);
+    }
+
+    private Connection getStatementConnection() {
+        if (statement == null) {
+            return null;
+        }
+        try {
+            Connection connection = statement.getConnection();
+            if (connection == null) {
+                LOG.warn(
+                        "The JDBC source statement returned no connection. "
+                                + "Closing the cached connection to avoid 
reusing an unknown "
+                                + "transaction.");
+                chunkSplitter.close();
             }
+            return connection;
+        } catch (SQLException e) {
+            LOG.warn(
+                    "Failed to get the JDBC source connection from the current 
statement. "
+                            + "Closing the cached connection to avoid reusing 
an unknown "
+                            + "transaction.",
+                    e);
+            chunkSplitter.close();
+            return null;
+        }
+    }
+
+    private void finishReadTransaction(Connection connection) {
+        try {
+            finishReadTransaction(connection, configuredAutoCommit);
+        } catch (SQLException e) {
+            LOG.warn(
+                    "Failed to finish the JDBC source read transaction. "
+                            + "Closing the connection to avoid leaving or 
reusing an idle "
+                            + "transaction.",
+                    e);
+            discardConnection(connection, e);
+        }
+    }
+
+    private void discardConnection(Connection connection, SQLException 
cleanupException) {
+        try {
+            if (connection != null) {
+                connection.close();
+            }
+        } catch (SQLException closeException) {
+            cleanupException.addSuppressed(closeException);
+            LOG.warn(
+                    "Failed to close the JDBC source connection after 
transaction cleanup failed.",
+                    cleanupException);
+        } finally {
+            // Clear the provider's cached reference and retry close for 
drivers whose first close
+            // attempt failed.
+            chunkSplitter.close();
+        }
+    }
+
+    static void finishReadTransaction(Connection connection, boolean 
configuredAutoCommit)
+            throws SQLException {
+        if (connection == null || connection.isClosed()) {
+            return;
+        }
+
+        boolean currentAutoCommit = connection.getAutoCommit();
+        if (!currentAutoCommit) {
+            // JDBC source reads do not have changes to commit. Rollback ends 
the server-side cursor
+            // transaction, releases its snapshot, and also recovers a 
transaction in failed state.
+            connection.rollback();
+        }
+        if (currentAutoCommit != configuredAutoCommit) {
+            connection.setAutoCommit(configuredAutoCommit);
         }
     }
 
diff --git 
a/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/JdbcInputFormatTest.java
 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/JdbcInputFormatTest.java
new file mode 100644
index 0000000000..13f5fb72f5
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-jdbc/src/test/java/org/apache/seatunnel/connectors/seatunnel/jdbc/internal/JdbcInputFormatTest.java
@@ -0,0 +1,272 @@
+/*
+ * 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;
+
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import org.apache.seatunnel.api.table.catalog.TableSchema;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.exception.JdbcConnectorException;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.converter.JdbcRowConverter;
+import 
org.apache.seatunnel.connectors.seatunnel.jdbc.internal.dialect.JdbcDialect;
+import org.apache.seatunnel.connectors.seatunnel.jdbc.source.ChunkSplitter;
+import org.apache.seatunnel.connectors.seatunnel.jdbc.source.JdbcSourceSplit;
+
+import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
+
+import java.io.IOException;
+import java.sql.Connection;
+import java.sql.PreparedStatement;
+import java.sql.ResultSet;
+import java.sql.SQLException;
+import java.util.Collections;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.inOrder;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.never;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
+
+class JdbcInputFormatTest {
+
+    private static final TablePath TABLE_PATH = TablePath.of("test", "public", 
"source_table");
+    private static final TableSchema TABLE_SCHEMA = 
TableSchema.builder().build();
+    private static final JdbcSourceSplit SPLIT =
+            new JdbcSourceSplit(TABLE_PATH, "split-0", null, null, null, null, 
null);
+
+    @Test
+    void shouldRollbackAndRestoreConfiguredAutoCommitAfterClosingSplit() 
throws Exception {
+        TestContext context = createContext(true);
+        context.openEmptySplit();
+        when(context.connection.isClosed()).thenReturn(false);
+        when(context.connection.getAutoCommit()).thenReturn(false);
+
+        context.inputFormat.close();
+
+        InOrder inOrder = inOrder(context.statement, context.resultSet, 
context.connection);
+        inOrder.verify(context.statement).getConnection();
+        inOrder.verify(context.resultSet).close();
+        inOrder.verify(context.statement).close();
+        inOrder.verify(context.connection).rollback();
+        inOrder.verify(context.connection).setAutoCommit(true);
+
+        context.inputFormat.close();
+        verify(context.connection, times(1)).rollback();
+    }
+
+    @Test
+    void shouldRollbackAndKeepConfiguredManualCommit() throws Exception {
+        TestContext context = createContext(false);
+        context.openEmptySplit();
+        when(context.connection.isClosed()).thenReturn(false);
+        when(context.connection.getAutoCommit()).thenReturn(false);
+
+        context.inputFormat.close();
+
+        verify(context.connection).rollback();
+        verify(context.connection, never()).setAutoCommit(true);
+    }
+
+    @Test
+    void shouldNotRollbackAutoCommitConnection() throws Exception {
+        TestContext context = createContext(true);
+        context.openEmptySplit();
+        when(context.connection.isClosed()).thenReturn(false);
+        when(context.connection.getAutoCommit()).thenReturn(true);
+
+        context.inputFormat.close();
+
+        verify(context.connection, never()).rollback();
+        verify(context.connection, never()).setAutoCommit(true);
+    }
+
+    @Test
+    void shouldFinishTransactionWhenResourceCloseFails() throws Exception {
+        TestContext context = createContext(true);
+        context.openEmptySplit();
+        when(context.connection.isClosed()).thenReturn(false);
+        when(context.connection.getAutoCommit()).thenReturn(false);
+        doThrow(new SQLException("result set close 
failed")).when(context.resultSet).close();
+        doThrow(new SQLException("statement close 
failed")).when(context.statement).close();
+
+        context.inputFormat.close();
+
+        verify(context.connection).rollback();
+        verify(context.connection).setAutoCommit(true);
+    }
+
+    @Test
+    void shouldKeepConnectionReusableAcrossSuccessfulSplits() throws Exception 
{
+        TestContext context = createContext(true);
+        PreparedStatement secondStatement = mock(PreparedStatement.class);
+        ResultSet secondResultSet = mock(ResultSet.class);
+
+        when(context.chunkSplitter.generateSplitStatement(SPLIT, TABLE_SCHEMA))
+                .thenReturn(context.statement, secondStatement);
+        when(context.statement.executeQuery()).thenReturn(context.resultSet);
+        when(secondStatement.executeQuery()).thenReturn(secondResultSet);
+        when(context.resultSet.next()).thenReturn(false);
+        when(secondResultSet.next()).thenReturn(false);
+        when(context.statement.getConnection()).thenReturn(context.connection);
+        when(secondStatement.getConnection()).thenReturn(context.connection);
+        when(context.connection.isClosed()).thenReturn(false);
+        when(context.connection.getAutoCommit()).thenReturn(false);
+
+        context.inputFormat.open(SPLIT);
+        context.inputFormat.close();
+        context.inputFormat.open(SPLIT);
+        context.inputFormat.close();
+
+        verify(context.chunkSplitter, times(2)).generateSplitStatement(SPLIT, 
TABLE_SCHEMA);
+        verify(context.connection, times(2)).rollback();
+        verify(context.connection, times(2)).setAutoCommit(true);
+        verify(context.connection, never()).close();
+        verify(context.chunkSplitter, never()).close();
+    }
+
+    @Test
+    void shouldDiscardConnectionWhenTransactionCleanupFails() throws Exception 
{
+        TestContext context = createContext(true);
+        context.openEmptySplit();
+        when(context.connection.isClosed()).thenReturn(false);
+        when(context.connection.getAutoCommit()).thenReturn(false);
+        doThrow(new SQLException("rollback 
failed")).when(context.connection).rollback();
+
+        context.inputFormat.close();
+
+        verify(context.connection).close();
+        verify(context.chunkSplitter).close();
+        verify(context.connection, never()).setAutoCommit(true);
+    }
+
+    @Test
+    void shouldCloseCachedConnectionWhenStatementCannotExposeConnection() 
throws Exception {
+        TestContext context = createContext(true);
+        context.openEmptySplit();
+        when(context.statement.getConnection())
+                .thenThrow(new SQLException("get connection failed"));
+
+        context.inputFormat.close();
+
+        verify(context.chunkSplitter).close();
+        verify(context.resultSet).close();
+        verify(context.statement).close();
+    }
+
+    @Test
+    void shouldDiscardCachedConnectionWhenStatementCreationFails() throws 
Exception {
+        TestContext context = createContext(true);
+        when(context.chunkSplitter.generateSplitStatement(SPLIT, TABLE_SCHEMA))
+                .thenThrow(new SQLException("prepare failed"));
+
+        assertThrows(JdbcConnectorException.class, () -> 
context.inputFormat.open(SPLIT));
+
+        verify(context.chunkSplitter).close();
+        verify(context.statement, never()).executeQuery();
+    }
+
+    @Test
+    void shouldRollbackAndKeepConnectionWhenExecuteQueryFails() throws 
Exception {
+        TestContext context = createContext(true);
+        when(context.chunkSplitter.generateSplitStatement(SPLIT, TABLE_SCHEMA))
+                .thenReturn(context.statement);
+        when(context.statement.getConnection()).thenReturn(context.connection);
+        when(context.statement.executeQuery()).thenThrow(new 
SQLException("execute failed"));
+        when(context.connection.isClosed()).thenReturn(false);
+        when(context.connection.getAutoCommit()).thenReturn(false);
+
+        assertThrows(JdbcConnectorException.class, () -> 
context.inputFormat.open(SPLIT));
+
+        verify(context.statement).close();
+        verify(context.connection).rollback();
+        verify(context.connection).setAutoCommit(true);
+        verify(context.connection, never()).close();
+        verify(context.chunkSplitter, never()).close();
+    }
+
+    @Test
+    void shouldAlwaysCloseCachedConnectionWhenInputFormatCloses() throws 
Exception {
+        TestContext context = createContext(true);
+
+        context.inputFormat.closeInputFormat();
+
+        verify(context.chunkSplitter).close();
+    }
+
+    @Test
+    void shouldIgnoreClosedConnection() throws SQLException {
+        Connection connection = mock(Connection.class);
+        when(connection.isClosed()).thenReturn(true);
+
+        JdbcInputFormat.finishReadTransaction(connection, true);
+
+        verify(connection, never()).getAutoCommit();
+        verify(connection, never()).rollback();
+    }
+
+    private static TestContext createContext(boolean configuredAutoCommit) 
throws SQLException {
+        JdbcDialect dialect = mock(JdbcDialect.class);
+        JdbcRowConverter rowConverter = mock(JdbcRowConverter.class);
+        ChunkSplitter chunkSplitter = mock(ChunkSplitter.class);
+        PreparedStatement statement = mock(PreparedStatement.class);
+        ResultSet resultSet = mock(ResultSet.class);
+        Connection connection = mock(Connection.class);
+        CatalogTable catalogTable = mock(CatalogTable.class);
+        Map<TablePath, CatalogTable> tables = 
Collections.singletonMap(TABLE_PATH, catalogTable);
+
+        when(dialect.getRowConverter()).thenReturn(rowConverter);
+        when(catalogTable.getTableSchema()).thenReturn(TABLE_SCHEMA);
+
+        JdbcInputFormat inputFormat =
+                new JdbcInputFormat(dialect, chunkSplitter, tables, 
configuredAutoCommit);
+        return new TestContext(inputFormat, chunkSplitter, statement, 
resultSet, connection);
+    }
+
+    private static final class TestContext {
+        private final JdbcInputFormat inputFormat;
+        private final ChunkSplitter chunkSplitter;
+        private final PreparedStatement statement;
+        private final ResultSet resultSet;
+        private final Connection connection;
+
+        private TestContext(
+                JdbcInputFormat inputFormat,
+                ChunkSplitter chunkSplitter,
+                PreparedStatement statement,
+                ResultSet resultSet,
+                Connection connection) {
+            this.inputFormat = inputFormat;
+            this.chunkSplitter = chunkSplitter;
+            this.statement = statement;
+            this.resultSet = resultSet;
+            this.connection = connection;
+        }
+
+        private void openEmptySplit() throws IOException, SQLException {
+            when(chunkSplitter.generateSplitStatement(SPLIT, 
TABLE_SCHEMA)).thenReturn(statement);
+            when(statement.executeQuery()).thenReturn(resultSet);
+            when(resultSet.next()).thenReturn(false);
+            when(statement.getConnection()).thenReturn(connection);
+            inputFormat.open(SPLIT);
+        }
+    }
+}

Reply via email to