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

github-merge-queue[bot] pushed a commit to branch 
gh-readonly-queue/dev/pr-11554-b2045b14c0ad3eeffd3e14c5d8f42f85f0bec2d0
in repository https://gitbox.apache.org/repos/asf/seatunnel.git

commit 62ae6a3a10d70a5f971bb5d9df0e443482463667
Author: Jast <[email protected]>
AuthorDate: Mon Sep 28 11:32:46 2026 +0000

    [Fix][Connector-V2] Support FTP append data mode (#11554)
    
    Co-authored-by: zhangshenghang <[email protected]>
    Co-authored-by: jast <[email protected]>
---
 docs/en/connectors/sink/FtpFile.md                 |  11 +-
 docs/zh/connectors/sink/FtpFile.md                 |   8 +-
 .../seatunnel/file/config/BaseFileSinkConfig.java  |   6 ++
 .../file/hadoop/HadoopFileSystemProxy.java         | 114 +++++++++++++++++++++
 .../file/sink/BaseMultipleTableFileSink.java       |   4 +-
 .../sink/commit/FileSinkAggregatedCommitter.java   |  52 ++++++++--
 .../commit/FileSinkAggregatedCommitterTest.java    |  68 ++++++++++++
 .../seatunnel/file/writer/FileSinkConfigTest.java  |  27 +++++
 .../file/ftp/system/SeaTunnelFTPFileSystem.java    |  69 ++++++++++++-
 .../ftp/system/SeaTunnelFTPFileSystemTest.java     |  15 +++
 10 files changed, 358 insertions(+), 16 deletions(-)

diff --git a/docs/en/connectors/sink/FtpFile.md 
b/docs/en/connectors/sink/FtpFile.md
index 4588315b95..ed0b765aa4 100644
--- a/docs/en/connectors/sink/FtpFile.md
+++ b/docs/en/connectors/sink/FtpFile.md
@@ -287,7 +287,16 @@ Existing dir processing method.
 Existing data processing method.
 
 - DROP_DATA: preserve dir and delete data files
-- APPEND_DATA: preserve dir, preserve data files
+- APPEND_DATA: preserve dir and data files. For FTP sinks, new rows are 
appended to
+  existing target files only when `data_save_mode = "APPEND_DATA"` is 
explicitly
+  configured in the job config. If this option is omitted and the value only 
comes from
+  the default, FTP sinks keep the legacy commit path and do not use FTP 
byte-level append.
+  Byte-level append additionally requires a stable target filename across 
commits, i.e.
+  `custom_filename = true` with a `file_name_expression` that does not vary per
+  transaction; with the default expression each checkpoint writes a new file, 
so there is
+  nothing to append to. FTP append is at-least-once: if a checkpoint is 
aborted while a
+  commit is only partially applied, the rows of that commit can appear twice 
in the target
+  file, so verify this is acceptable for your target before enabling the mode
 - ERROR_WHEN_DATA_EXISTS: when there is data files, an error is reported
 
 ### schema_evolution_enabled [boolean]
diff --git a/docs/zh/connectors/sink/FtpFile.md 
b/docs/zh/connectors/sink/FtpFile.md
index 437a6d747f..14a5da01c6 100644
--- a/docs/zh/connectors/sink/FtpFile.md
+++ b/docs/zh/connectors/sink/FtpFile.md
@@ -297,7 +297,13 @@ Sink 插件的通用参数,请参考[Sink通用选项](../common-options/sink-
 
 现有数据处理方法:
 - DROP_DATA(删除数据):保留目录,删除数据文件。
-- APPEND_DATA(追加数据):保留目录和数据文件。
+- APPEND_DATA(追加数据):保留目录和数据文件。对于 FTP Sink,只有在作业配置中显式写出
+  `data_save_mode = "APPEND_DATA"` 时,才会把新数据追加到已有目标文件;如果省略该配置、
+  仅使用默认值,FTP Sink 会保持历史提交路径,不启用 FTP 字节级追加。字节级追加还要求目标
+  文件名在多次提交间保持稳定,即设置 `custom_filename = true` 且 `file_name_expression`
+  不随事务变化;使用默认表达式时每个检查点都会写出新文件,没有可追加的目标。FTP 追加为至少
+  一次语义:如果检查点在提交只完成一部分时被中止,该提交的数据行可能在目标文件中出现两次,
+  启用该模式前请确认这对目标场景可接受。
 - ERROR_WHEN_DATA_EXISTS(数据存在时报错):当存在数据文件时,报告错误。
 
 ### schema_evolution_enabled [boolean]
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/BaseFileSinkConfig.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/BaseFileSinkConfig.java
index 488cf83a89..cbfc2a0af7 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/BaseFileSinkConfig.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/BaseFileSinkConfig.java
@@ -18,6 +18,7 @@
 package org.apache.seatunnel.connectors.seatunnel.file.config;
 
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.DataSaveMode;
 import org.apache.seatunnel.common.utils.DateTimeUtils;
 import org.apache.seatunnel.common.utils.DateUtils;
 import org.apache.seatunnel.common.utils.TimeUtils;
@@ -45,6 +46,8 @@ public class BaseFileSinkConfig implements DelimiterConfig, 
Serializable {
     protected boolean createEmptyFileWhenNoData;
     protected FileFormat fileFormat;
     protected String filenameExtension;
+    protected DataSaveMode dataSaveMode;
+    protected boolean dataSaveModeExplicitlyConfigured;
     protected DateUtils.Formatter dateFormat;
     protected DateTimeUtils.Formatter datetimeFormat;
     protected TimeUtils.Formatter timeFormat;
@@ -66,6 +69,9 @@ public class BaseFileSinkConfig implements DelimiterConfig, 
Serializable {
         this.createEmptyFileWhenNoData =
                 
pluginConfig.get(FileBaseSinkOptions.CREATE_EMPTY_FILE_WHEN_NO_DATA);
         this.fileFormat = 
pluginConfig.get(FileBaseSinkOptions.FILE_FORMAT_TYPE);
+        this.dataSaveMode = 
pluginConfig.get(FileBaseSinkOptions.DATA_SAVE_MODE);
+        this.dataSaveModeExplicitlyConfigured =
+                
pluginConfig.getOptional(FileBaseSinkOptions.DATA_SAVE_MODE).isPresent();
         // if set, use user config, if not set, when format is csv, use "," 
otherwise use default
         // delimiter
         if 
(pluginConfig.getOptional(FileBaseSinkOptions.FIELD_DELIMITER).isPresent()) {
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxy.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxy.java
index b443d4749a..e1269482f2 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxy.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxy.java
@@ -33,6 +33,7 @@ import org.apache.hadoop.fs.FileSystem;
 import org.apache.hadoop.fs.LocatedFileStatus;
 import org.apache.hadoop.fs.Path;
 import org.apache.hadoop.fs.RemoteIterator;
+import org.apache.hadoop.io.IOUtils;
 import org.apache.hadoop.security.UserGroupInformation;
 
 import lombok.NonNull;
@@ -42,6 +43,7 @@ import lombok.extern.slf4j.Slf4j;
 import java.io.Closeable;
 import java.io.IOException;
 import java.io.Serializable;
+import java.nio.charset.StandardCharsets;
 import java.security.PrivilegedExceptionAction;
 import java.util.ArrayList;
 import java.util.List;
@@ -49,6 +51,9 @@ import java.util.List;
 @Slf4j
 public class HadoopFileSystemProxy implements Serializable, Closeable {
 
+    private static final int DEFAULT_BUFFER_SIZE = 4096;
+    private static final String APPEND_TARGET_LENGTH_MARKER_SUFFIX = 
".append-target-length";
+
     private transient UserGroupInformation userGroupInformation;
     private transient FileSystem fileSystem;
 
@@ -178,6 +183,115 @@ public class HadoopFileSystemProxy implements 
Serializable, Closeable {
                 });
     }
 
+    /**
+     * Appends a temporary file to an existing target file, or moves it when 
the target is absent.
+     *
+     * <p>The move path preserves the normal rename behavior for the first 
committed file. Later
+     * commits record the target length before appending. If a retry sees the 
expected post-append
+     * length, it only cleans the temporary file instead of appending the same 
bytes again.
+     *
+     * <p>A retry is therefore idempotent only when it replays the same 
aggregated commit info (for
+     * example a commit re-run after restore). Once abort() has deleted the 
transaction directory
+     * together with its markers, a later commit re-appends bytes that may 
already have been
+     * written: FTP append mode is at-least-once in that case, not 
exactly-once.
+     */
+    public void appendFile(@NonNull String sourceFilePath, @NonNull String 
targetFilePath)
+            throws IOException {
+        execute(
+                () -> {
+                    Path sourcePath = new Path(sourceFilePath);
+                    Path targetPath = new Path(targetFilePath);
+                    Path markerPath = new Path(sourceFilePath + 
APPEND_TARGET_LENGTH_MARKER_SUFFIX);
+                    FileSystem fileSystem = getFileSystem();
+                    if (!fileSystem.exists(sourcePath)) {
+                        deleteAppendMarkerIfExists(fileSystem, markerPath);
+                        log.warn(
+                                "append file:[{}] to [{}] already finished in 
the last commit, skip.",
+                                sourcePath,
+                                targetPath);
+                        return Void.class;
+                    }
+                    if (!fileSystem.exists(targetPath)) {
+                        renameFile(sourceFilePath, targetFilePath, false);
+                        return Void.class;
+                    }
+                    long sourceLength = 
fileSystem.getFileStatus(sourcePath).getLen();
+                    long initialTargetLength =
+                            getOrCreateAppendTargetLength(fileSystem, 
markerPath, targetPath);
+                    long expectedTargetLength = initialTargetLength + 
sourceLength;
+                    long currentTargetLength = 
fileSystem.getFileStatus(targetPath).getLen();
+                    if (currentTargetLength == expectedTargetLength) {
+                        cleanupCommittedAppend(fileSystem, sourcePath, 
markerPath);
+                        log.info(
+                                "append file:[{}] to [{}] already applied in 
the last commit, cleanup source file.",
+                                sourcePath,
+                                targetPath);
+                        return Void.class;
+                    }
+                    if (currentTargetLength != initialTargetLength) {
+                        throw new IOException(
+                                String.format(
+                                        "Cannot append file [%s] to [%s], 
target length changed from [%s] to [%s] before append.",
+                                        sourcePath,
+                                        targetPath,
+                                        initialTargetLength,
+                                        currentTargetLength));
+                    }
+                    try (FSDataInputStream inputStream = 
fileSystem.open(sourcePath);
+                            FSDataOutputStream outputStream =
+                                    fileSystem.append(targetPath, 
DEFAULT_BUFFER_SIZE, null)) {
+                        IOUtils.copyBytes(inputStream, outputStream, 
DEFAULT_BUFFER_SIZE, false);
+                    }
+                    long actualTargetLength = 
fileSystem.getFileStatus(targetPath).getLen();
+                    if (actualTargetLength != expectedTargetLength) {
+                        throw new IOException(
+                                String.format(
+                                        "Append file [%s] to [%s] finished 
with unexpected target length [%s], expected [%s].",
+                                        sourcePath,
+                                        targetPath,
+                                        actualTargetLength,
+                                        expectedTargetLength));
+                    }
+                    cleanupCommittedAppend(fileSystem, sourcePath, markerPath);
+                    log.info("append file:[{}] to [{}] finish", sourcePath, 
targetPath);
+                    return Void.class;
+                });
+    }
+
+    private long getOrCreateAppendTargetLength(
+            FileSystem fileSystem, Path markerPath, Path targetPath) throws 
IOException {
+        if (fileSystem.exists(markerPath)) {
+            try (FSDataInputStream inputStream = fileSystem.open(markerPath)) {
+                byte[] bytes = new byte[64];
+                int length = inputStream.read(bytes);
+                if (length <= 0) {
+                    throw new IOException("Append target length marker is 
empty: " + markerPath);
+                }
+                return Long.parseLong(new String(bytes, 0, length, 
StandardCharsets.UTF_8).trim());
+            }
+        }
+        long targetLength = fileSystem.getFileStatus(targetPath).getLen();
+        try (FSDataOutputStream outputStream = fileSystem.create(markerPath, 
false)) {
+            
outputStream.write(Long.toString(targetLength).getBytes(StandardCharsets.UTF_8));
+        }
+        return targetLength;
+    }
+
+    private void cleanupCommittedAppend(FileSystem fileSystem, Path 
sourcePath, Path markerPath)
+            throws IOException {
+        if (!fileSystem.delete(sourcePath, false)) {
+            throw CommonError.fileOperationFailed("SeaTunnel", "delete", 
sourcePath.toString());
+        }
+        deleteAppendMarkerIfExists(fileSystem, markerPath);
+    }
+
+    private void deleteAppendMarkerIfExists(FileSystem fileSystem, Path 
markerPath)
+            throws IOException {
+        if (fileSystem.exists(markerPath) && !fileSystem.delete(markerPath, 
false)) {
+            log.warn("Delete append marker [{}] failed, ignore this cleanup 
error.", markerPath);
+        }
+    }
+
     public void createDir(@NonNull String filePath) throws IOException {
         execute(
                 () -> {
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/BaseMultipleTableFileSink.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/BaseMultipleTableFileSink.java
index 075cdf84c2..e441b167a1 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/BaseMultipleTableFileSink.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/BaseMultipleTableFileSink.java
@@ -104,7 +104,9 @@ public abstract class BaseMultipleTableFileSink
     @Override
     public Optional<SinkAggregatedCommitter<FileCommitInfo, 
FileAggregatedCommitInfo>>
             createAggregatedCommitter() {
-        return Optional.of(new FileSinkAggregatedCommitter(hadoopConf));
+        boolean appendData =
+                FileSinkAggregatedCommitter.shouldAppendData(hadoopConf, 
fileSinkConfig);
+        return Optional.of(new FileSinkAggregatedCommitter(hadoopConf, 
appendData));
     }
 
     @Override
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitter.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitter.java
index 2d0a7efd0e..f9982151b6 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitter.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitter.java
@@ -17,9 +17,11 @@
 
 package org.apache.seatunnel.connectors.seatunnel.file.sink.commit;
 
+import org.apache.seatunnel.api.sink.DataSaveMode;
 import org.apache.seatunnel.api.sink.SinkAggregatedCommitter;
 import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
 import 
org.apache.seatunnel.connectors.seatunnel.file.hadoop.HadoopFileSystemProxy;
+import 
org.apache.seatunnel.connectors.seatunnel.file.sink.config.FileSinkConfig;
 
 import org.apache.hadoop.fs.Path;
 
@@ -38,11 +40,34 @@ public class FileSinkAggregatedCommitter
         implements SinkAggregatedCommitter<FileCommitInfo, 
FileAggregatedCommitInfo> {
     protected HadoopFileSystemProxy hadoopFileSystemProxy;
     private final HadoopConf hadoopConf;
+
+    /**
+     * True only for FTP sinks configured with APPEND_DATA. In this mode 
commits append staged bytes
+     * to an existing target and abort must not rename the target back to the 
transaction directory,
+     * because the target may already contain user data that predates this 
checkpoint.
+     *
+     * <p>Abort also deletes the transaction directory together with the 
append-length markers, so a
+     * partially applied commit can no longer be detected after restore and 
its rows are appended
+     * again: after an aborted commit, append mode is at-least-once, not 
exactly-once.
+     */
+    private final boolean appendData;
+
     private final Set<String> pendingUuidDirectories = new LinkedHashSet<>();
     private final Set<String> pendingJobDirectories = new LinkedHashSet<>();
 
     public FileSinkAggregatedCommitter(HadoopConf hadoopConf) {
+        this(hadoopConf, false);
+    }
+
+    public FileSinkAggregatedCommitter(HadoopConf hadoopConf, boolean 
appendData) {
         this.hadoopConf = hadoopConf;
+        this.appendData = appendData;
+    }
+
+    public static boolean shouldAppendData(HadoopConf hadoopConf, 
FileSinkConfig fileSinkConfig) {
+        return fileSinkConfig.isDataSaveModeExplicitlyConfigured()
+                && 
DataSaveMode.APPEND_DATA.equals(fileSinkConfig.getDataSaveMode())
+                && "ftp".equalsIgnoreCase(hadoopConf.getSchema());
     }
 
     @Override
@@ -61,9 +86,13 @@ public class FileSinkAggregatedCommitter
                                 
aggregatedCommitInfo.getTransactionMap().entrySet()) {
                             for (Map.Entry<String, String> mvFileEntry :
                                     entry.getValue().entrySet()) {
-                                // first rename temp file
-                                hadoopFileSystemProxy.renameFile(
-                                        mvFileEntry.getKey(), 
mvFileEntry.getValue(), true);
+                                if (appendData) {
+                                    hadoopFileSystemProxy.appendFile(
+                                            mvFileEntry.getKey(), 
mvFileEntry.getValue());
+                                } else {
+                                    hadoopFileSystemProxy.renameFile(
+                                            mvFileEntry.getKey(), 
mvFileEntry.getValue(), true);
+                                }
                             }
                             String transactionDir = entry.getKey();
                             // Data files are already committed after rename; 
tmp cleanup is
@@ -135,13 +164,16 @@ public class FileSinkAggregatedCommitter
                     try {
                         for (Map.Entry<String, LinkedHashMap<String, String>> 
entry :
                                 
aggregatedCommitInfo.getTransactionMap().entrySet()) {
-                            // rollback the file
-                            for (Map.Entry<String, String> mvFileEntry :
-                                    entry.getValue().entrySet()) {
-                                if 
(hadoopFileSystemProxy.fileExist(mvFileEntry.getValue())
-                                        && 
!hadoopFileSystemProxy.fileExist(mvFileEntry.getKey())) {
-                                    hadoopFileSystemProxy.renameFile(
-                                            mvFileEntry.getValue(), 
mvFileEntry.getKey(), true);
+                            if (!appendData) {
+                                // rollback the file
+                                for (Map.Entry<String, String> mvFileEntry :
+                                        entry.getValue().entrySet()) {
+                                    if 
(hadoopFileSystemProxy.fileExist(mvFileEntry.getValue())
+                                            && 
!hadoopFileSystemProxy.fileExist(
+                                                    mvFileEntry.getKey())) {
+                                        hadoopFileSystemProxy.renameFile(
+                                                mvFileEntry.getValue(), 
mvFileEntry.getKey(), true);
+                                    }
                                 }
                             }
                             // delete the transaction dir
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitterTest.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitterTest.java
index ab7b885376..41799b5287 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitterTest.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/sink/commit/FileSinkAggregatedCommitterTest.java
@@ -17,8 +17,10 @@
 
 package org.apache.seatunnel.connectors.seatunnel.file.sink.commit;
 
+import org.apache.seatunnel.api.sink.DataSaveMode;
 import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
 import 
org.apache.seatunnel.connectors.seatunnel.file.hadoop.HadoopFileSystemProxy;
+import 
org.apache.seatunnel.connectors.seatunnel.file.sink.config.FileSinkConfig;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -42,6 +44,10 @@ class FileSinkAggregatedCommitterTest {
             super(new HadoopConf("hdfs://dummy"));
         }
 
+        TestableCommitter(boolean appendData) {
+            super(new HadoopConf("ftp://dummy";), appendData);
+        }
+
         void setFileSystemProxy(HadoopFileSystemProxy proxy) {
             this.hadoopFileSystemProxy = proxy;
         }
@@ -97,12 +103,65 @@ class FileSinkAggregatedCommitterTest {
         Mockito.verify(fs).deleteFile(TRANSACTION_DIR);
     }
 
+    @Test
+    void shouldAppendTemporaryFileWhenAppendDataIsEnabled() throws Exception {
+        HadoopFileSystemProxy fs = Mockito.mock(HadoopFileSystemProxy.class);
+        TestableCommitter committer = newCommitter(fs, true);
+
+        List<FileAggregatedCommitInfo> errors =
+                
committer.commit(Collections.singletonList(newCommitInfo(true)));
+
+        Assertions.assertTrue(errors.isEmpty());
+        Mockito.verify(fs).appendFile(TEMP_FILE, TARGET_FILE);
+        Mockito.verify(fs, Mockito.never()).renameFile(TEMP_FILE, TARGET_FILE, 
true);
+    }
+
+    @Test
+    void shouldNotRollbackExistingTargetWhenAppendDataIsEnabled() throws 
Exception {
+        HadoopFileSystemProxy fs = Mockito.mock(HadoopFileSystemProxy.class);
+        TestableCommitter committer = newCommitter(fs, true);
+
+        committer.abort(Collections.singletonList(newCommitInfo(true)));
+
+        Mockito.verify(fs, Mockito.never())
+                .renameFile(Mockito.anyString(), Mockito.anyString(), 
Mockito.anyBoolean());
+        Mockito.verify(fs).deleteFile(TRANSACTION_DIR);
+    }
+
+    @Test
+    void shouldEnableAppendDataOnlyForFtpAppendDataMode() {
+        FileSinkConfig fileSinkConfig = Mockito.mock(FileSinkConfig.class);
+        
Mockito.when(fileSinkConfig.getDataSaveMode()).thenReturn(DataSaveMode.APPEND_DATA);
+        
Mockito.when(fileSinkConfig.isDataSaveModeExplicitlyConfigured()).thenReturn(true);
+
+        Assertions.assertTrue(
+                
FileSinkAggregatedCommitter.shouldAppendData(newHadoopConf("ftp"), 
fileSinkConfig));
+        Assertions.assertFalse(
+                FileSinkAggregatedCommitter.shouldAppendData(
+                        newHadoopConf("hdfs"), fileSinkConfig));
+
+        
Mockito.when(fileSinkConfig.isDataSaveModeExplicitlyConfigured()).thenReturn(false);
+        Assertions.assertFalse(
+                
FileSinkAggregatedCommitter.shouldAppendData(newHadoopConf("ftp"), 
fileSinkConfig));
+
+        
Mockito.when(fileSinkConfig.isDataSaveModeExplicitlyConfigured()).thenReturn(true);
+        
Mockito.when(fileSinkConfig.getDataSaveMode()).thenReturn(DataSaveMode.DROP_DATA);
+        Assertions.assertFalse(
+                
FileSinkAggregatedCommitter.shouldAppendData(newHadoopConf("ftp"), 
fileSinkConfig));
+    }
+
     private static TestableCommitter newCommitter(HadoopFileSystemProxy fs) {
         TestableCommitter committer = new TestableCommitter();
         committer.setFileSystemProxy(fs);
         return committer;
     }
 
+    private static TestableCommitter newCommitter(HadoopFileSystemProxy fs, 
boolean appendData) {
+        TestableCommitter committer = new TestableCommitter(appendData);
+        committer.setFileSystemProxy(fs);
+        return committer;
+    }
+
     private static FileAggregatedCommitInfo newCommitInfo(boolean 
withFileMove) {
         LinkedHashMap<String, String> fileMoves = new LinkedHashMap<>();
         if (withFileMove) {
@@ -112,4 +171,13 @@ class FileSinkAggregatedCommitterTest {
         transactionMap.put(TRANSACTION_DIR, fileMoves);
         return new FileAggregatedCommitInfo(transactionMap, new 
LinkedHashMap<>());
     }
+
+    private static HadoopConf newHadoopConf(String schema) {
+        return new HadoopConf(schema + "://dummy") {
+            @Override
+            public String getSchema() {
+                return schema;
+            }
+        };
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/FileSinkConfigTest.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/FileSinkConfigTest.java
index b1317c5ff7..a22e7750b7 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/FileSinkConfigTest.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/writer/FileSinkConfigTest.java
@@ -21,6 +21,7 @@ import org.apache.seatunnel.shade.com.typesafe.config.Config;
 import org.apache.seatunnel.shade.com.typesafe.config.ConfigFactory;
 
 import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.sink.DataSaveMode;
 import org.apache.seatunnel.api.table.type.BasicType;
 import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
 import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
@@ -85,4 +86,30 @@ public class FileSinkConfigTest {
         Assertions.assertEquals(
                 sinkColumnsIndexInRow.size(), 
seaTunnelRowTypeInfo.getFieldNames().length);
     }
+
+    @Test
+    public void testDataSaveModeExplicitlyConfigured() {
+        SeaTunnelRowType rowType = newRowType();
+        FileSinkConfig defaultConfig =
+                new FileSinkConfig(
+                        ReadonlyConfig.fromConfig(
+                                ConfigFactory.parseString("path = 
\"/data/test\"")),
+                        rowType);
+        FileSinkConfig explicitConfig =
+                new FileSinkConfig(
+                        ReadonlyConfig.fromConfig(
+                                ConfigFactory.parseString(
+                                        "path = \"/data/test\"\ndata_save_mode 
= \"APPEND_DATA\"")),
+                        rowType);
+
+        Assertions.assertEquals(DataSaveMode.APPEND_DATA, 
defaultConfig.getDataSaveMode());
+        
Assertions.assertFalse(defaultConfig.isDataSaveModeExplicitlyConfigured());
+        
Assertions.assertTrue(explicitConfig.isDataSaveModeExplicitlyConfigured());
+    }
+
+    private static SeaTunnelRowType newRowType() {
+        return new SeaTunnelRowType(
+                new String[] {"data", "ts"},
+                new SeaTunnelDataType[] {BasicType.STRING_TYPE, 
BasicType.STRING_TYPE});
+    }
 }
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 7cf322fd87..e708a4cfc2 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
@@ -401,11 +401,74 @@ public class SeaTunnelFTPFileSystem extends FileSystem 
implements StreamingFileS
         return fos;
     }
 
-    /** This optional operation is not yet supported. */
+    /**
+     * Appends bytes to an existing file through the FTP APPE command.
+     *
+     * <p>A stream obtained via this call must be closed before using other 
APIs of this class or
+     * else the invocation will block.
+     */
     @Override
-    public FSDataOutputStream append(Path f, int bufferSize, Progressable 
progress)
+    public FSDataOutputStream append(Path file, int bufferSize, Progressable 
progress)
             throws IOException {
-        throw new IOException("Not supported");
+        final FTPClient client = connect();
+        Path workDir = new Path(client.printWorkingDirectory());
+        Path absolute = makeAbsolute(workDir, file);
+        FileStatus status;
+        try {
+            status = getFileStatus(client, absolute);
+        } catch (FileNotFoundException e) {
+            disconnect(client);
+            throw e;
+        }
+        if (status.isDirectory()) {
+            disconnect(client);
+            throw new FileNotFoundException("Path " + file + " is a 
directory.");
+        }
+
+        Path parent = absolute.getParent();
+        client.allocate(bufferSize);
+        client.changeWorkingDirectory(parent.toUri().getPath());
+        FSDataOutputStream fos =
+                new 
FSDataOutputStream(client.appendFileStream(file.getName()), statistics) {
+                    @Override
+                    public void close() throws IOException {
+                        IOException closeException = null;
+                        try {
+                            super.close();
+                            if (!client.isConnected()) {
+                                throw new FTPException("Client not connected");
+                            }
+                            boolean cmdCompleted = 
client.completePendingCommand();
+                            if (!cmdCompleted) {
+                                throw new FTPException(
+                                        "Could not complete transfer, Reply 
Code - "
+                                                + client.getReplyCode());
+                            }
+                        } catch (IOException | FTPException e) {
+                            closeException =
+                                    e instanceof IOException
+                                            ? (IOException) e
+                                            : new IOException(e.getMessage(), 
e);
+                        } finally {
+                            // disconnect() is deliberately lenient: it 
returns when the client is
+                            // already disconnected and only logs logout 
failures, so releasing the
+                            // connection here can never mask closeException.
+                            disconnect(client);
+                        }
+                        if (closeException != null) {
+                            throw closeException;
+                        }
+                    }
+                };
+        if (!FTPReply.isPositivePreliminary(client.getReplyCode())) {
+            try {
+                fos.close();
+            } catch (IOException | FTPException e) {
+                LOG.warn("Close rejected FTP append stream failed, ignore 
cleanup error.", e);
+            }
+            throw new IOException("Unable to append file: " + file + ", 
Aborting");
+        }
+        return fos;
     }
 
     /**
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 6c75ca9bd1..1169f1c40e 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
@@ -151,6 +151,21 @@ public class SeaTunnelFTPFileSystemTest {
         assertFalse(ftpFileSystem.exists(testFile));
     }
 
+    @Test
+    public void testAppendToExistingFile() throws IOException {
+        Path testFile = new Path(HOME_DIR + "/test.txt");
+
+        try (FSDataOutputStream out = ftpFileSystem.append(testFile, 1024, 
null)) {
+            out.write(" appended".getBytes(StandardCharsets.UTF_8));
+        }
+
+        try (FSDataInputStream in = ftpFileSystem.open(testFile, 1024)) {
+            byte[] buffer = new byte["Test content appended".length()];
+            in.readFully(buffer);
+            assertEquals("Test content appended", new String(buffer, 
StandardCharsets.UTF_8));
+        }
+    }
+
     @Test
     public void testListStatus() throws IOException {
         // Create test directory structure

Reply via email to