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 66214e4600 [Fix][Connector-V2] Do not fail FTP listings on logout
reset (#11695)
66214e4600 is described below
commit 66214e4600cc2baa4b07441eb08c08f3c97a78d1
Author: Jast <[email protected]>
AuthorDate: Thu Aug 20 23:11:48 2026 +0800
[Fix][Connector-V2] Do not fail FTP listings on logout reset (#11695)
Co-authored-by: zhangshenghang <[email protected]>
Co-authored-by: davidzollo <[email protected]>
---
.../file/ftp/system/SeaTunnelFTPFileSystem.java | 32 ++++++++++++++--------
.../ftp/system/SeaTunnelFTPFileSystemTest.java | 18 ++++++++++++
2 files changed, 38 insertions(+), 12 deletions(-)
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
index f5cb9bfaf0..7cf322fd87 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystem.java
@@ -249,19 +249,31 @@ public class SeaTunnelFTPFileSystem extends FileSystem
implements StreamingFileS
* Logout and disconnect the given FTPClient. *
*
* @param client FTPClient
- * @throws IOException IOException
*/
- private void disconnect(FTPClient client) throws IOException {
- if (client != null) {
- if (!client.isConnected()) {
- throw new FTPException("Client not connected");
- }
+ void disconnect(FTPClient client) {
+ if (client == null || !client.isConnected()) {
+ return;
+ }
+
+ try {
boolean logoutSuccess = client.logout();
- client.disconnect();
if (!logoutSuccess) {
LOG.warn(
"Logout failed while disconnecting, error code - " +
client.getReplyCode());
}
+ } catch (IOException e) {
+ // Some FTP servers close the control connection before responding
to QUIT. The
+ // preceding operation has already completed, so do not turn a
successful operation
+ // into a failure while releasing the connection.
+ LOG.warn("Failed to logout from FTP server while disconnecting",
e);
+ } finally {
+ if (client.isConnected()) {
+ try {
+ client.disconnect();
+ } catch (IOException e) {
+ LOG.warn("Failed to disconnect from FTP server", e);
+ }
+ }
}
}
@@ -854,11 +866,7 @@ public class SeaTunnelFTPFileSystem extends FileSystem
implements StreamingFileS
} catch (IOException ioe) {
throw new FTPException("Failed to get home directory", ioe);
} finally {
- try {
- disconnect(client);
- } catch (IOException ioe) {
- throw new FTPException("Failed to disconnect", ioe);
- }
+ disconnect(client);
}
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
index a29351bc8b..6c75ca9bd1 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-ftp/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/ftp/system/SeaTunnelFTPFileSystemTest.java
@@ -19,6 +19,7 @@ package
org.apache.seatunnel.connectors.seatunnel.file.ftp.system;
import
org.apache.seatunnel.connectors.seatunnel.file.hadoop.FileStatusListingSession;
+import org.apache.commons.net.ftp.FTPClient;
import org.apache.commons.net.ftp.FTPFile;
import org.apache.hadoop.conf.Configuration;
import org.apache.hadoop.fs.FSDataInputStream;
@@ -38,17 +39,23 @@ import org.mockftpserver.fake.filesystem.FileSystem;
import org.mockftpserver.fake.filesystem.UnixFakeFileSystem;
import java.io.IOException;
+import java.net.SocketException;
import java.net.URI;
import java.nio.charset.StandardCharsets;
import java.util.ArrayList;
import java.util.Calendar;
import java.util.List;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
import static org.junit.jupiter.api.Assertions.assertEquals;
import static org.junit.jupiter.api.Assertions.assertFalse;
import static org.junit.jupiter.api.Assertions.assertNotNull;
import static org.junit.jupiter.api.Assertions.assertThrows;
import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.Mockito.doThrow;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
/** Unit tests for SeaTunnelFTPFileSystem. */
public class SeaTunnelFTPFileSystemTest {
@@ -169,6 +176,17 @@ public class SeaTunnelFTPFileSystemTest {
ftpFileSystem.delete(testDir, true);
}
+ @Test
+ public void testDisconnectIgnoresConnectionResetDuringLogout() throws
IOException {
+ FTPClient client = mock(FTPClient.class);
+ when(client.isConnected()).thenReturn(true);
+ doThrow(new SocketException("Connection reset")).when(client).logout();
+
+ assertDoesNotThrow(() -> ftpFileSystem.disconnect(client));
+
+ verify(client).disconnect();
+ }
+
@Test
public void testStreamingListingSkipsUnparseableEntries() throws
IOException {
FTPFile file = new FTPFile();