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
