This is an automated email from the ASF dual-hosted git repository.
github-merge-queue[bot] 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 72acda5704 [Feature][Connector-V2] Add Google Cloud Storage file sink
(#12146)
72acda5704 is described below
commit 72acda57042efa3408681dc2c93663910b06b3e8
Author: Goutam Adwant <[email protected]>
AuthorDate: Mon Sep 7 06:04:19 2026 +0000
[Feature][Connector-V2] Add Google Cloud Storage file sink (#12146)
Signed-off-by: goutamadwant <[email protected]>
Signed-off-by: Goutam Adwant <[email protected]>
---
docs/en/connectors/sink/GcsFile.md | 172 +++++++++++++++++++
docs/zh/connectors/sink/GcsFile.md | 165 ++++++++++++++++++
plugin-mapping.properties | 1 +
.../file/hadoop/HadoopFileSystemProxy.java | 11 ++
.../file/hadoop/HadoopFileSystemProxyTest.java | 50 ++++++
.../seatunnel/file/gcs/catalog/GcsFileCatalog.java | 44 +++++
.../file/gcs/catalog/GcsFileCatalogFactory.java | 55 ++++++
.../file/gcs/config/GcsFileSinkOptions.java | 20 +++
.../seatunnel/file/gcs/sink/GcsFileSink.java | 48 ++++++
.../file/gcs/sink/GcsFileSinkFactory.java | 158 +++++++++++++++++
.../seatunnel/file/gcs/GcsFileSinkFactoryTest.java | 188 +++++++++++++++++++++
.../file/gcs/catalog/GcsFileCatalogTest.java | 106 ++++++++++++
.../e2e/connector/file/gcs/GcsFileIT.java | 14 ++
.../src/test/resources/gcs/gcs_file_to_gcs.conf | 67 ++++++++
.../src/test/resources/gcs/gcs_sink_to_assert.conf | 63 +++++++
15 files changed, 1162 insertions(+)
diff --git a/docs/en/connectors/sink/GcsFile.md
b/docs/en/connectors/sink/GcsFile.md
new file mode 100644
index 0000000000..e7fc89b1b9
--- /dev/null
+++ b/docs/en/connectors/sink/GcsFile.md
@@ -0,0 +1,172 @@
+import ChangeLog from '../changelog/connector-file-gcs.md';
+
+# GcsFile
+
+> Google Cloud Storage file sink connector
+
+## Support Those Engines
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## Key Features
+
+- [x] [batch](../../introduction/concepts/connector-v2-features.md)
+- [x] [exactly-once](../../introduction/concepts/connector-v2-features.md)
+- [x] [parallelism](../../introduction/concepts/connector-v2-features.md)
+- [x] multiple table sink
+- [x] partitioned output
+- [x] save modes
+- [x] file formats: `text`, `csv`, `parquet`, `orc`, `json`, `excel`, `xml`,
`binary`, `canal_json`, `debezium_json`, and `maxwell_json`
+
+## Description
+
+Writes data to Google Cloud Storage through the Google Cloud Storage connector
for Hadoop. The
+connector reuses SeaTunnel's file sink implementation for serialization,
partitioning, file
+rotation, checkpoint commits, multiple-table jobs, and save-mode handling.
+
+Set `bucket` to a bucket URI such as `gs://my-bucket`. Set `path` and
`tmp_path` to prefixes inside
+that bucket, such as `/warehouse/orders` and `/tmp/seatunnel/orders`. Do not
include an object path
+in `bucket`.
+
+When transactions are enabled, writers first create files under `tmp_path` and
publish them to
+`path` during commit. The configured identity therefore needs read, create,
list, and delete
+permissions for both prefixes, plus object move permission when using the
driver's native move
+operation. Transactional commits do not provide atomic visibility of an entire
+directory to external readers: the Hadoop GCS connector can publish files
using copy and delete
+operations. Disabling `is_enable_transaction` disables the file sink's
transactional commit path.
+
+## Dependency
+
+The connector uses
`com.google.cloud.bigdataoss:gcs-connector:hadoop3-2.2.33:shaded`, which is
+Apache License 2.0 software and targets Java 8. The shaded GCS Hadoop library
is packaged in the
+`connector-file-gcs` connector JAR. Spark and Flink deployments must provide a
compatible Hadoop 3
+runtime on every driver and worker.
+
+## Authentication
+
+The connector supports these authentication modes:
+
+1. **Application Default Credentials (ADC):** omit `service_account_key_file`.
The Hadoop GCS
+ connector discovers credentials from `GOOGLE_APPLICATION_CREDENTIALS` or
the service account
+ attached to the Google Cloud runtime.
+2. **Service-account JSON file:** set `service_account_key_file` to a local
path that exists at the
+ same location on every node that writes GCS.
+
+The explicit `service_account_key_file` option takes precedence over the
corresponding entry in
+`hadoop_gcs_properties`. Do not store service-account JSON content in the job
configuration or in
+`hadoop_gcs_properties`.
+
+## Sink Options
+
+| Name | Type | Required | Default | Description |
+|------|------|----------|---------|-------------|
+| path | string | yes | - | Destination prefix inside `bucket`, for example
`/warehouse/orders`. Supports `${database_name}`, `${schema_name}`, and
`${table_name}` placeholders. |
+| bucket | string | yes | - | GCS bucket URI, for example `gs://my-bucket`. |
+| service_account_key_file | string | no | - | Service-account JSON key file
on every worker. When omitted, ADC is used. |
+| hadoop_gcs_properties | map | no | - | Additional `fs.gs.*` Hadoop
properties. Explicit connector options take precedence. |
+| tmp_path | string | no | `/tmp/seatunnel` | Temporary prefix inside the
configured bucket used before transactional commit. Use a distinct prefix from
`path`. |
+| file_format_type | string | no | `csv` | Output file format. |
+| schema_save_mode | enum | no | `CREATE_SCHEMA_WHEN_NOT_EXIST` | How the
destination prefix is prepared before writing. |
+| data_save_mode | enum | no | `APPEND_DATA` | Whether existing objects are
kept, deleted, or rejected before writing. |
+| custom_filename | boolean | no | `false` | Enables `file_name_expression`. |
+| file_name_expression | string | conditional | `${transactionId}` | Output
filename expression when `custom_filename=true`. Supports `${transactionId}`,
`${now}`, and `${uuid}`. |
+| filename_time_format | string | no | `yyyy.MM.dd` | Time format used by
`${now}`. |
+| filename_extension | string | no | format-specific | Overrides the output
filename extension. |
+| have_partition | boolean | no | `false` | Writes rows into partition
directories. |
+| partition_by | list | conditional | - | Partition columns when
`have_partition=true`. |
+| partition_dir_expression | string | no |
`${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/` | Partition directory expression. |
+| is_partition_field_write_in_file | boolean | no | `false` | Whether
partition fields remain in output rows. |
+| sink_columns | list | no | all columns | Columns written to the output file,
in output order. |
+| is_enable_transaction | boolean | no | `true` | Stages files and publishes
them during commit. |
+| batch_size | int | no | `1000000` | Maximum rows in a split output file. |
+| single_file_mode | boolean | no | `false` | Produces one output file per
parallel sink task. Not supported with checkpointing or streaming mode. |
+| create_empty_file_when_no_data | boolean | no | `false` | Creates an empty
output file when no rows arrive. |
+| row_delimiter | string | conditional | `\n` | Row delimiter for text, CSV,
and JSON. |
+| field_delimiter | string | conditional | `\001` | Field delimiter for text. |
+| encoding | string | conditional | `UTF-8` | Encoding for text, CSV, JSON,
and XML. |
+| enable_header_write | boolean | conditional | `false` | Writes a header for
text or CSV output. |
+| compress_codec | string | no | `none` | Compression codec supported by the
selected format. |
+| common-options | | no | - | See [Sink Common
Options](../common-options/sink-common-options.md). |
+
+The sink also accepts the format-specific file sink options documented by the
corresponding
+SeaTunnel file formats, including Parquet INT96 and XML element options.
+
+When `is_partition_field_write_in_file=true`, set
`parse_partition_from_path=false` on a file
+source reading that output if its schema already includes those partition
columns.
+
+## Save Modes
+
+`schema_save_mode` controls destination-prefix creation:
+
+- `RECREATE_SCHEMA`: delete and recreate the prefix.
+- `CREATE_SCHEMA_WHEN_NOT_EXIST`: create the prefix only when absent.
+- `ERROR_WHEN_SCHEMA_NOT_EXIST`: fail when the prefix does not exist.
+- `IGNORE`: do not prepare the prefix.
+
+`data_save_mode` controls existing objects:
+
+- `APPEND_DATA`: keep existing objects and add new output files.
+- `DROP_DATA`: delete objects below `path` before writing.
+- `ERROR_WHEN_DATA_EXISTS`: fail when objects already exist below `path`.
+
+These operations apply only to `path` inside the configured bucket. GcsFile
rejects a bucket-root
+`path` when `RECREATE_SCHEMA` or `DROP_DATA` is selected; configure a
dedicated prefix instead.
+
+## Examples
+
+### Write Parquet With ADC
+
+```hocon
+sink {
+ GcsFile {
+ bucket = "gs://my-bucket"
+ path = "/warehouse/orders"
+ tmp_path = "/tmp/seatunnel/orders"
+ file_format_type = "parquet"
+ data_save_mode = "APPEND_DATA"
+ }
+}
+```
+
+### Write Partitioned CSV With a Service Account
+
+```hocon
+sink {
+ GcsFile {
+ bucket = "gs://my-bucket"
+ path = "/exports/customers"
+ tmp_path = "/tmp/seatunnel/customers"
+ file_format_type = "csv"
+ service_account_key_file = "/opt/seatunnel/keys/gcs-writer.json"
+ enable_header_write = true
+ have_partition = true
+ partition_by = ["region"]
+ hadoop_gcs_properties = {
+ "fs.gs.project.id" = "my-project"
+ }
+ }
+}
+```
+
+### Multiple-Table Output
+
+Use table placeholders to keep each upstream table under a separate
destination and temporary
+prefix:
+
+```hocon
+sink {
+ GcsFile {
+ bucket = "gs://my-bucket"
+ path = "/warehouse/${database_name}/${table_name}"
+ tmp_path = "/tmp/seatunnel/${database_name}/${table_name}"
+ file_format_type = "parquet"
+ data_save_mode = "APPEND_DATA"
+ }
+}
+```
+
+## Changelog
+
+<ChangeLog />
diff --git a/docs/zh/connectors/sink/GcsFile.md
b/docs/zh/connectors/sink/GcsFile.md
new file mode 100644
index 0000000000..08457493fe
--- /dev/null
+++ b/docs/zh/connectors/sink/GcsFile.md
@@ -0,0 +1,165 @@
+import ChangeLog from '../changelog/connector-file-gcs.md';
+
+# GcsFile
+
+> Google Cloud Storage 文件 Sink 连接器
+
+## 支持的引擎
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## 主要特性
+
+- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
+- [x] [精确一次](../../introduction/concepts/connector-v2-features.md)
+- [x] [并行度](../../introduction/concepts/connector-v2-features.md)
+- [x] 多表 Sink
+- [x] 分区输出
+- [x] 保存模式
+- [x]
文件格式:`text`、`csv`、`parquet`、`orc`、`json`、`excel`、`xml`、`binary`、`canal_json`、`debezium_json`
和 `maxwell_json`
+
+## 描述
+
+通过 Google Cloud Storage Hadoop 连接器向 GCS 写入文件。序列化、分区、文件轮转、
+Checkpoint 提交、多表任务和保存模式复用 SeaTunnel 现有的 File Sink 实现。
+
+`bucket` 必须是 `gs://my-bucket` 形式的存储桶 URI。`path` 和 `tmp_path` 是该存储桶内的
+前缀,例如 `/warehouse/orders` 和 `/tmp/seatunnel/orders`。不要把对象路径写入 `bucket`。
+
+启用事务时,Writer 首先在 `tmp_path` 下创建文件,并在提交阶段发布到 `path`。配置的身份
+需要对两个前缀具有读取、创建、列举和删除对象的权限;使用驱动的原生移动操作时还需要对象移动
+权限。事务提交不保证外部读取者能原子地看到
+整个目录:Hadoop GCS 连接器可能通过复制和删除操作发布文件。禁用 `is_enable_transaction`
+会禁用 File Sink 的事务提交路径。
+
+## 依赖
+
+本连接器使用 `com.google.cloud.bigdataoss:gcs-connector:hadoop3-2.2.33:shaded`。该依赖采用
+Apache License 2.0,并以 Java 8 为目标版本。shaded GCS Hadoop 库已打包到
+`connector-file-gcs` 连接器 JAR 中。Spark 和 Flink 部署必须在驱动节点和所有工作节点
+提供兼容的 Hadoop 3 运行环境。
+
+## 认证
+
+连接器支持以下认证方式:
+
+1. **应用默认凭据(ADC)**:不配置 `service_account_key_file`。Hadoop GCS 连接器会从
+ `GOOGLE_APPLICATION_CREDENTIALS` 或 Google Cloud 运行环境绑定的服务账号获取凭据。
+2. **服务账号 JSON 文件**:配置 `service_account_key_file`。该本地路径必须以相同位置
+ 存在于每个写入 GCS 的节点上。
+
+显式配置的 `service_account_key_file` 优先于 `hadoop_gcs_properties` 中的同名 Hadoop
+属性。不要把服务账号 JSON 内容写入任务配置或 `hadoop_gcs_properties`。
+
+## Sink 配置项
+
+| 名称 | 类型 | 是否必填 | 默认值 | 描述 |
+|------|------|----------|--------|------|
+| path | string | 是 | - | `bucket` 内的目标前缀,例如 `/warehouse/orders`。支持
`${database_name}`、`${schema_name}` 和 `${table_name}` 占位符。 |
+| bucket | string | 是 | - | GCS 存储桶 URI,例如 `gs://my-bucket`。 |
+| service_account_key_file | string | 否 | - | 每个工作节点上的服务账号 JSON 文件。省略时使用 ADC。 |
+| hadoop_gcs_properties | map | 否 | - | 额外的 `fs.gs.*` Hadoop 属性。显式连接器配置优先。 |
+| tmp_path | string | 否 | `/tmp/seatunnel` | 事务提交前使用的存储桶内临时前缀,应与 `path` 不同。 |
+| file_format_type | string | 否 | `csv` | 输出文件格式。 |
+| schema_save_mode | enum | 否 | `CREATE_SCHEMA_WHEN_NOT_EXIST` | 写入前如何准备目标前缀。 |
+| data_save_mode | enum | 否 | `APPEND_DATA` | 写入前保留、删除或拒绝已有对象。 |
+| custom_filename | boolean | 否 | `false` | 是否启用 `file_name_expression`。 |
+| file_name_expression | string | 条件必填 | `${transactionId}` |
`custom_filename=true` 时的文件名表达式,支持 `${transactionId}`、`${now}` 和 `${uuid}`。 |
+| filename_time_format | string | 否 | `yyyy.MM.dd` | `${now}` 使用的时间格式。 |
+| filename_extension | string | 否 | 按格式决定 | 自定义输出文件扩展名。 |
+| have_partition | boolean | 否 | `false` | 是否写入分区目录。 |
+| partition_by | list | 条件必填 | - | `have_partition=true` 时使用的分区字段。 |
+| partition_dir_expression | string | 否 |
`${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/` | 分区目录表达式。 |
+| is_partition_field_write_in_file | boolean | 否 | `false` | 是否在输出行中保留分区字段。 |
+| sink_columns | list | 否 | 全部字段 | 写入输出文件的字段及顺序。 |
+| is_enable_transaction | boolean | 否 | `true` | 在提交阶段发布临时文件。 |
+| batch_size | int | 否 | `1000000` | 单个切分输出文件的最大行数。 |
+| single_file_mode | boolean | 否 | `false` | 每个 Sink 并行任务输出一个文件,不支持开启
Checkpoint 或流模式。 |
+| create_empty_file_when_no_data | boolean | 否 | `false` | 没有输入数据时创建空文件。 |
+| row_delimiter | string | 条件必填 | `\n` | text、CSV 和 JSON 的行分隔符。 |
+| field_delimiter | string | 条件必填 | `\001` | text 的字段分隔符。 |
+| encoding | string | 条件必填 | `UTF-8` | text、CSV、JSON 和 XML 的编码。 |
+| enable_header_write | boolean | 条件必填 | `false` | 为 text 或 CSV 输出写入表头。 |
+| compress_codec | string | 否 | `none` | 所选格式支持的压缩编码。 |
+| common-options | | 否 | - | 参见 [Sink
通用配置](../common-options/sink-common-options.md)。 |
+
+Sink 也支持对应 SeaTunnel 文件格式定义的格式专用配置,包括 Parquet INT96 和 XML 元素配置。
+
+当 `is_partition_field_write_in_file=true` 时,如果读取这些输出文件的 File Source 的 Schema
+已包含分区字段,应设置 `parse_partition_from_path=false`,避免重复添加分区字段。
+
+## 保存模式
+
+`schema_save_mode` 控制目标前缀的创建:
+
+- `RECREATE_SCHEMA`:删除并重新创建前缀。
+- `CREATE_SCHEMA_WHEN_NOT_EXIST`:仅当前缀不存在时创建。
+- `ERROR_WHEN_SCHEMA_NOT_EXIST`:前缀不存在时失败。
+- `IGNORE`:不处理前缀。
+
+`data_save_mode` 控制已有对象:
+
+- `APPEND_DATA`:保留已有对象并添加新的输出文件。
+- `DROP_DATA`:写入前删除 `path` 下的对象。
+- `ERROR_WHEN_DATA_EXISTS`:`path` 下已有对象时失败。
+
+这些操作仅应用于所配置存储桶中的 `path`。选择 `RECREATE_SCHEMA` 或 `DROP_DATA` 时,
+GcsFile 会拒绝存储桶根路径;请改用专用前缀。
+
+## 示例
+
+### 使用 ADC 写入 Parquet
+
+```hocon
+sink {
+ GcsFile {
+ bucket = "gs://my-bucket"
+ path = "/warehouse/orders"
+ tmp_path = "/tmp/seatunnel/orders"
+ file_format_type = "parquet"
+ data_save_mode = "APPEND_DATA"
+ }
+}
+```
+
+### 使用服务账号写入分区 CSV
+
+```hocon
+sink {
+ GcsFile {
+ bucket = "gs://my-bucket"
+ path = "/exports/customers"
+ tmp_path = "/tmp/seatunnel/customers"
+ file_format_type = "csv"
+ service_account_key_file = "/opt/seatunnel/keys/gcs-writer.json"
+ enable_header_write = true
+ have_partition = true
+ partition_by = ["region"]
+ hadoop_gcs_properties = {
+ "fs.gs.project.id" = "my-project"
+ }
+ }
+}
+```
+
+### 多表输出
+
+使用表占位符把每个上游表写入独立的目标和临时前缀:
+
+```hocon
+sink {
+ GcsFile {
+ bucket = "gs://my-bucket"
+ path = "/warehouse/${database_name}/${table_name}"
+ tmp_path = "/tmp/seatunnel/${database_name}/${table_name}"
+ file_format_type = "parquet"
+ data_save_mode = "APPEND_DATA"
+ }
+}
+```
+
+## Changelog
+
+<ChangeLog />
diff --git a/plugin-mapping.properties b/plugin-mapping.properties
index 7c07bb5fed..6baf75df38 100644
--- a/plugin-mapping.properties
+++ b/plugin-mapping.properties
@@ -84,6 +84,7 @@ seatunnel.source.InfluxDB = connector-influxdb
seatunnel.source.S3File = connector-file-s3
seatunnel.sink.S3File = connector-file-s3
seatunnel.source.GcsFile = connector-file-gcs
+seatunnel.sink.GcsFile = connector-file-gcs
seatunnel.source.AmazonDynamodb = connector-amazondynamodb
seatunnel.sink.AmazonDynamodb = connector-amazondynamodb
seatunnel.source.AmazonDocumentDB = connector-amazondocumentdb
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 bc84f43f52..69eaf478b4 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
@@ -215,6 +215,17 @@ public class HadoopFileSystemProxy implements
Serializable, Closeable {
});
}
+ /** Checks for a file without collecting the directory listing in memory.
*/
+ public boolean hasAnyFile(@NonNull String path, boolean recursive) throws
IOException {
+ return execute(
+ () -> {
+ Path fileName = new Path(path);
+ FileSystem fileSystem = getFileSystem();
+ return fileSystem.exists(fileName)
+ && fileSystem.listFiles(fileName,
recursive).hasNext();
+ });
+ }
+
public List<Path> getAllSubFiles(@NonNull String filePath) throws
IOException {
return execute(
() -> {
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxyTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxyTest.java
index f4401bdadc..833bc2e6de 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxyTest.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/hadoop/HadoopFileSystemProxyTest.java
@@ -19,13 +19,17 @@ package
org.apache.seatunnel.connectors.seatunnel.file.hadoop;
import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
+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.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
import org.junit.jupiter.api.condition.DisabledOnOs;
import org.junit.jupiter.api.condition.OS;
import org.junit.jupiter.api.io.TempDir;
+import org.mockito.Mockito;
import java.io.IOException;
import java.nio.charset.StandardCharsets;
@@ -66,6 +70,45 @@ class HadoopFileSystemProxyTest {
}
}
+ @Test
+ void testHasAnyFileUsesLazyListingAndRecursiveFlag() throws Exception {
+ Path warehouse = new Path("gs://test-bucket/warehouse");
+ FileSystem fileSystem = Mockito.mock(FileSystem.class);
+ @SuppressWarnings("unchecked")
+ RemoteIterator<LocatedFileStatus> directFiles =
Mockito.mock(RemoteIterator.class);
+ @SuppressWarnings("unchecked")
+ RemoteIterator<LocatedFileStatus> recursiveFiles =
Mockito.mock(RemoteIterator.class);
+ Mockito.when(fileSystem.exists(warehouse)).thenReturn(true);
+ Mockito.when(fileSystem.listFiles(warehouse,
false)).thenReturn(directFiles);
+ Mockito.when(fileSystem.listFiles(warehouse,
true)).thenReturn(recursiveFiles);
+ Mockito.when(directFiles.hasNext()).thenReturn(false);
+ Mockito.when(recursiveFiles.hasNext()).thenReturn(true);
+
+ try (HadoopFileSystemProxy proxy = newProxy(fileSystem)) {
+ Assertions.assertFalse(proxy.hasAnyFile(warehouse.toString(),
false));
+ Assertions.assertTrue(proxy.hasAnyFile(warehouse.toString(),
true));
+ }
+
+ Mockito.verify(fileSystem).listFiles(warehouse, false);
+ Mockito.verify(fileSystem).listFiles(warehouse, true);
+ Mockito.verify(directFiles, Mockito.never()).next();
+ Mockito.verify(recursiveFiles, Mockito.never()).next();
+ }
+
+ @Test
+ void testHasAnyFileSkipsListingWhenPathDoesNotExist() throws Exception {
+ Path missingPath = new Path("gs://test-bucket/missing");
+ FileSystem fileSystem = Mockito.mock(FileSystem.class);
+ Mockito.when(fileSystem.exists(missingPath)).thenReturn(false);
+
+ try (HadoopFileSystemProxy proxy = newProxy(fileSystem)) {
+ Assertions.assertFalse(proxy.hasAnyFile(missingPath.toString(),
true));
+ }
+
+ Mockito.verify(fileSystem, Mockito.never())
+ .listFiles(Mockito.any(Path.class), Mockito.anyBoolean());
+ }
+
@Test
void testRenameTreatsExistingTargetAsCompletedRetry() throws Exception {
HadoopFileSystemProxy proxy = new HadoopFileSystemProxy(new
HadoopConf("file:///"));
@@ -101,4 +144,11 @@ class HadoopFileSystemProxyTest {
proxy.close();
}
}
+
+ private static HadoopFileSystemProxy newProxy(FileSystem fileSystem) {
+ HadoopFileSystemProxy proxy =
+ Mockito.mock(HadoopFileSystemProxy.class,
Mockito.CALLS_REAL_METHODS);
+ Mockito.doReturn(fileSystem).when(proxy).getFileSystem();
+ return proxy;
+ }
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalog.java
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalog.java
new file mode 100644
index 0000000000..9c8a7a30cf
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalog.java
@@ -0,0 +1,44 @@
+/*
+ * 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.file.gcs.catalog;
+
+import org.apache.seatunnel.api.table.catalog.TablePath;
+import
org.apache.seatunnel.connectors.seatunnel.file.catalog.AbstractFileCatalog;
+import
org.apache.seatunnel.connectors.seatunnel.file.hadoop.HadoopFileSystemProxy;
+
+import lombok.SneakyThrows;
+
+/** File catalog used by GCS sink save-mode handling. */
+public class GcsFileCatalog extends AbstractFileCatalog {
+
+ private final HadoopFileSystemProxy hadoopFileSystemProxy;
+ private final String filePath;
+
+ GcsFileCatalog(
+ HadoopFileSystemProxy hadoopFileSystemProxy, String filePath,
String catalogName) {
+ super(hadoopFileSystemProxy, filePath, catalogName);
+ this.hadoopFileSystemProxy = hadoopFileSystemProxy;
+ this.filePath = filePath;
+ }
+
+ @SneakyThrows
+ @Override
+ public boolean isExistsData(TablePath tablePath) {
+ return hadoopFileSystemProxy.hasAnyFile(filePath, true);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalogFactory.java
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalogFactory.java
new file mode 100644
index 0000000000..4b935a7b56
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalogFactory.java
@@ -0,0 +1,55 @@
+/*
+ * 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.file.gcs.catalog;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.table.catalog.Catalog;
+import org.apache.seatunnel.api.table.factory.CatalogFactory;
+import org.apache.seatunnel.api.table.factory.Factory;
+import
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSinkOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileSystemType;
+import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
+import org.apache.seatunnel.connectors.seatunnel.file.gcs.config.GcsHadoopConf;
+import
org.apache.seatunnel.connectors.seatunnel.file.hadoop.HadoopFileSystemProxy;
+
+import com.google.auto.service.AutoService;
+
+/** Creates the GCS file catalog used for sink save modes. */
+@AutoService(Factory.class)
+public class GcsFileCatalogFactory implements CatalogFactory {
+
+ @Override
+ public Catalog createCatalog(String catalogName, ReadonlyConfig options) {
+ HadoopConf hadoopConf = GcsHadoopConf.buildWithReadonlyConfig(options);
+ return new GcsFileCatalog(
+ new HadoopFileSystemProxy(hadoopConf),
+ options.get(FileBaseSinkOptions.FILE_PATH),
+ factoryIdentifier());
+ }
+
+ @Override
+ public String factoryIdentifier() {
+ return FileSystemType.GCS.getFileSystemPluginName();
+ }
+
+ @Override
+ public OptionRule optionRule() {
+ return OptionRule.builder().build();
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/config/GcsFileSinkOptions.java
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/config/GcsFileSinkOptions.java
new file mode 100644
index 0000000000..1f62d98784
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/config/GcsFileSinkOptions.java
@@ -0,0 +1,20 @@
+/*
+ * 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.file.gcs.config;
+
+public class GcsFileSinkOptions extends GcsFileBaseOptions {}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/sink/GcsFileSink.java
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/sink/GcsFileSink.java
new file mode 100644
index 0000000000..10bd24b5b1
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/sink/GcsFileSink.java
@@ -0,0 +1,48 @@
+/*
+ * 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.file.gcs.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileSystemType;
+import org.apache.seatunnel.connectors.seatunnel.file.gcs.config.GcsHadoopConf;
+import
org.apache.seatunnel.connectors.seatunnel.file.sink.BaseMultipleTableFileSink;
+
+import java.util.Optional;
+
+/** File sink implementation for objects stored in Google Cloud Storage. */
+public class GcsFileSink extends BaseMultipleTableFileSink {
+
+ private final CatalogTable catalogTable;
+
+ /** Creates a GCS file sink for one resolved catalog table. */
+ public GcsFileSink(ReadonlyConfig readonlyConfig, CatalogTable
catalogTable) {
+ super(GcsHadoopConf.buildWithReadonlyConfig(readonlyConfig),
readonlyConfig, catalogTable);
+ this.catalogTable = catalogTable;
+ }
+
+ @Override
+ public String getPluginName() {
+ return FileSystemType.GCS.getFileSystemPluginName();
+ }
+
+ @Override
+ public Optional<CatalogTable> getWriteCatalogTable() {
+ return Optional.ofNullable(catalogTable);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/sink/GcsFileSinkFactory.java
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/sink/GcsFileSinkFactory.java
new file mode 100644
index 0000000000..b822d6c238
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/sink/GcsFileSinkFactory.java
@@ -0,0 +1,158 @@
+/*
+ * 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.file.gcs.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.options.SinkConnectorCommonOptions;
+import org.apache.seatunnel.api.sink.DataSaveMode;
+import org.apache.seatunnel.api.sink.SchemaSaveMode;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.connector.TableSink;
+import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSinkOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileFormat;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileSystemType;
+import
org.apache.seatunnel.connectors.seatunnel.file.factory.BaseMultipleTableFileSinkFactory;
+import
org.apache.seatunnel.connectors.seatunnel.file.gcs.config.GcsFileSinkOptions;
+import
org.apache.seatunnel.connectors.seatunnel.file.sink.commit.FileAggregatedCommitInfo;
+import
org.apache.seatunnel.connectors.seatunnel.file.sink.commit.FileCommitInfo;
+import org.apache.seatunnel.connectors.seatunnel.file.sink.state.FileSinkState;
+
+import org.apache.hadoop.fs.Path;
+
+import com.google.auto.service.AutoService;
+
+import java.util.Arrays;
+
+/** Creates Google Cloud Storage file sinks from table factory configuration.
*/
+@AutoService(Factory.class)
+public class GcsFileSinkFactory extends BaseMultipleTableFileSinkFactory {
+
+ @Override
+ public String factoryIdentifier() {
+ return FileSystemType.GCS.getFileSystemPluginName();
+ }
+
+ @Override
+ public OptionRule optionRule() {
+ return OptionRule.builder()
+ .required(FileBaseSinkOptions.FILE_PATH)
+ .required(GcsFileSinkOptions.BUCKET)
+ .optional(
+ GcsFileSinkOptions.SERVICE_ACCOUNT_KEY_FILE,
+ GcsFileSinkOptions.GCS_PROPERTIES)
+ .optional(FileBaseSinkOptions.SCHEMA_SAVE_MODE)
+ .optional(FileBaseSinkOptions.DATA_SAVE_MODE)
+ .optional(FileBaseSinkOptions.FILE_FORMAT_TYPE)
+ .conditional(
+ FileBaseSinkOptions.FILE_FORMAT_TYPE,
+ FileFormat.TEXT,
+ FileBaseSinkOptions.ROW_DELIMITER,
+ FileBaseSinkOptions.FIELD_DELIMITER,
+ FileBaseSinkOptions.TXT_COMPRESS,
+ FileBaseSinkOptions.ENABLE_HEADER_WRITE)
+ .conditional(
+ FileBaseSinkOptions.FILE_FORMAT_TYPE,
+ FileFormat.CSV,
+ FileBaseSinkOptions.ROW_DELIMITER,
+ FileBaseSinkOptions.TXT_COMPRESS,
+ FileBaseSinkOptions.ENABLE_HEADER_WRITE)
+ .conditional(
+ FileBaseSinkOptions.FILE_FORMAT_TYPE,
+ FileFormat.JSON,
+ FileBaseSinkOptions.ROW_DELIMITER,
+ FileBaseSinkOptions.TXT_COMPRESS)
+ .conditional(
+ FileBaseSinkOptions.FILE_FORMAT_TYPE,
+ FileFormat.ORC,
+ FileBaseSinkOptions.ORC_COMPRESS)
+ .conditional(
+ FileBaseSinkOptions.FILE_FORMAT_TYPE,
+ FileFormat.PARQUET,
+ FileBaseSinkOptions.PARQUET_COMPRESS,
+ FileBaseSinkOptions.PARQUET_AVRO_WRITE_FIXED_AS_INT96,
+
FileBaseSinkOptions.PARQUET_AVRO_WRITE_TIMESTAMP_AS_INT96)
+ .conditional(
+ FileBaseSinkOptions.FILE_FORMAT_TYPE,
+ FileFormat.XML,
+ FileBaseSinkOptions.XML_USE_ATTR_FORMAT,
+ FileBaseSinkOptions.XML_ROOT_TAG,
+ FileBaseSinkOptions.XML_ROW_TAG)
+ .optional(FileBaseSinkOptions.CUSTOM_FILENAME)
+ .conditional(
+ FileBaseSinkOptions.CUSTOM_FILENAME,
+ true,
+ FileBaseSinkOptions.FILE_NAME_EXPRESSION,
+ FileBaseSinkOptions.FILENAME_TIME_FORMAT)
+ .optional(FileBaseSinkOptions.HAVE_PARTITION)
+ .conditional(
+ FileBaseSinkOptions.HAVE_PARTITION,
+ true,
+ FileBaseSinkOptions.PARTITION_BY,
+ FileBaseSinkOptions.PARTITION_DIR_EXPRESSION,
+ FileBaseSinkOptions.IS_PARTITION_FIELD_WRITE_IN_FILE)
+ .conditional(
+ FileBaseSinkOptions.FILE_FORMAT_TYPE,
+ Arrays.asList(
+ FileFormat.TEXT, FileFormat.JSON,
FileFormat.CSV, FileFormat.XML),
+ FileBaseSinkOptions.ENCODING)
+ .optional(FileBaseSinkOptions.SINK_COLUMNS)
+ .optional(FileBaseSinkOptions.IS_ENABLE_TRANSACTION)
+ .optional(FileBaseSinkOptions.DATE_FORMAT_LEGACY)
+ .optional(FileBaseSinkOptions.DATETIME_FORMAT_LEGACY)
+ .optional(FileBaseSinkOptions.TIME_FORMAT_LEGACY)
+ .optional(FileBaseSinkOptions.SINGLE_FILE_MODE)
+ .optional(FileBaseSinkOptions.BATCH_SIZE)
+ .optional(FileBaseSinkOptions.CREATE_EMPTY_FILE_WHEN_NO_DATA)
+ .optional(SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA)
+ .optional(FileBaseSinkOptions.FILENAME_EXTENSION)
+ .optional(FileBaseSinkOptions.TMP_PATH)
+ .build();
+ }
+
+ @Override
+ public TableSink<SeaTunnelRow, FileSinkState, FileCommitInfo,
FileAggregatedCommitInfo>
+ createSink(TableSinkFactoryContext context) {
+ ReadonlyConfig readonlyConfig = context.getOptions();
+ validateDestructivePath(readonlyConfig);
+ CatalogTable catalogTable = context.getCatalogTable();
+ return () -> new GcsFileSink(readonlyConfig, catalogTable);
+ }
+
+ private static void validateDestructivePath(ReadonlyConfig readonlyConfig)
{
+ SchemaSaveMode schemaSaveMode =
readonlyConfig.get(FileBaseSinkOptions.SCHEMA_SAVE_MODE);
+ DataSaveMode dataSaveMode =
readonlyConfig.get(FileBaseSinkOptions.DATA_SAVE_MODE);
+ if (schemaSaveMode != SchemaSaveMode.RECREATE_SCHEMA
+ && dataSaveMode != DataSaveMode.DROP_DATA) {
+ return;
+ }
+
+ String filePath = readonlyConfig.get(FileBaseSinkOptions.FILE_PATH);
+ String normalizedPath = new
Path(filePath).toUri().normalize().getPath();
+ if (normalizedPath == null || normalizedPath.isEmpty() ||
"/".equals(normalizedPath)) {
+ throw new IllegalArgumentException(
+ String.format(
+ "GcsFile path must point to a prefix below the
bucket root when "
+ + "schema_save_mode is %s or
data_save_mode is %s",
+ schemaSaveMode, dataSaveMode));
+ }
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/GcsFileSinkFactoryTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/GcsFileSinkFactoryTest.java
new file mode 100644
index 0000000000..1a7f42fdc4
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/GcsFileSinkFactoryTest.java
@@ -0,0 +1,188 @@
+/*
+ * 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.file.gcs;
+
+import org.apache.seatunnel.api.configuration.Option;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.options.SinkConnectorCommonOptions;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.api.table.factory.CatalogFactory;
+import org.apache.seatunnel.api.table.factory.TableSinkFactory;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSinkOptions;
+import
org.apache.seatunnel.connectors.seatunnel.file.gcs.catalog.GcsFileCatalogFactory;
+import
org.apache.seatunnel.connectors.seatunnel.file.gcs.config.GcsFileSinkOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.gcs.config.GcsHadoopConf;
+import
org.apache.seatunnel.connectors.seatunnel.file.gcs.sink.GcsFileSinkFactory;
+
+import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
+
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.Map;
+
+import static
org.apache.seatunnel.api.table.factory.FactoryUtil.discoverFactory;
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class GcsFileSinkFactoryTest {
+
+ @Test
+ void shouldIdentifyGcsFileSinkAndCatalog() {
+ assertEquals("GcsFile", new GcsFileSinkFactory().factoryIdentifier());
+ assertEquals("GcsFile", new
GcsFileCatalogFactory().factoryIdentifier());
+ }
+
+ @Test
+ void shouldRegisterSinkAndCatalogFactories() {
+ ClassLoader classLoader =
Thread.currentThread().getContextClassLoader();
+
+ assertEquals(
+ GcsFileSinkFactory.class,
+ discoverFactory(classLoader, TableSinkFactory.class,
"GcsFile").getClass());
+ assertEquals(
+ GcsFileCatalogFactory.class,
+ discoverFactory(classLoader, CatalogFactory.class,
"GcsFile").getClass());
+ }
+
+ @Test
+ void shouldRequirePathAndBucket() {
+ OptionRule optionRule = new GcsFileSinkFactory().optionRule();
+ Map<String, Object> config = sinkConfig();
+
+ assertDoesNotThrow(() -> validate(config, optionRule));
+
+ config.remove(FileBaseSinkOptions.FILE_PATH.key());
+ assertThrows(OptionValidationException.class, () -> validate(config,
optionRule));
+
+ config.put(FileBaseSinkOptions.FILE_PATH.key(), "/warehouse/orders");
+ config.remove(GcsFileSinkOptions.BUCKET.key());
+ assertThrows(OptionValidationException.class, () -> validate(config,
optionRule));
+ }
+
+ @Test
+ void shouldExposeAuthenticationSaveModeAndMultiTableOptions() {
+ OptionRule optionRule = new GcsFileSinkFactory().optionRule();
+
+ assertTrue(optionRuleContains(optionRule,
GcsFileSinkOptions.SERVICE_ACCOUNT_KEY_FILE));
+ assertTrue(optionRuleContains(optionRule,
GcsFileSinkOptions.GCS_PROPERTIES));
+ assertTrue(optionRuleContains(optionRule,
FileBaseSinkOptions.SCHEMA_SAVE_MODE));
+ assertTrue(optionRuleContains(optionRule,
FileBaseSinkOptions.DATA_SAVE_MODE));
+ assertTrue(
+ optionRuleContains(
+ optionRule,
SinkConnectorCommonOptions.MULTI_TABLE_SINK_REPLICA));
+ }
+
+ @Test
+ void shouldRejectInvalidBucketForSinkConfiguration() {
+ Map<String, Object> config = sinkConfig();
+ config.put(GcsFileSinkOptions.BUCKET.key(), "test-bucket");
+
+ assertThrows(
+ IllegalArgumentException.class,
+ () ->
GcsHadoopConf.buildWithReadonlyConfig(ReadonlyConfig.fromMap(config)));
+ }
+
+ @ParameterizedTest
+ @ValueSource(
+ strings = {
+ "/",
+ "/./",
+ "/warehouse/..",
+ "/warehouse/../",
+ "gs://test-bucket",
+ "gs://test-bucket/",
+ "gs://test-bucket/warehouse/.."
+ })
+ void shouldRejectNormalizedBucketRootForDropData(String path) {
+ Map<String, Object> config = sinkConfig();
+ config.put(FileBaseSinkOptions.FILE_PATH.key(), path);
+ config.put(FileBaseSinkOptions.DATA_SAVE_MODE.key(), "DROP_DATA");
+
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> new
GcsFileSinkFactory().createSink(sinkContext(config)));
+ }
+
+ @Test
+ void shouldRejectBucketRootForRecreateSchemaAfterPlaceholderResolution() {
+ Map<String, Object> config = sinkConfig();
+ config.put(FileBaseSinkOptions.FILE_PATH.key(), "/${table_name}/..");
+ config.put(FileBaseSinkOptions.SCHEMA_SAVE_MODE.key(),
"RECREATE_SCHEMA");
+
+ IllegalArgumentException error =
+ assertThrows(
+ IllegalArgumentException.class,
+ () -> new
GcsFileSinkFactory().createSink(sinkContext(config)));
+
+ assertTrue(error.getMessage().contains("bucket root"));
+ }
+
+ @Test
+ void shouldAllowBucketRootForNonDestructiveModes() {
+ Map<String, Object> config = sinkConfig();
+ config.put(FileBaseSinkOptions.FILE_PATH.key(), "/");
+
+ assertDoesNotThrow(() -> new
GcsFileSinkFactory().createSink(sinkContext(config)));
+ }
+
+ @Test
+ void shouldAllowNormalizedPrefixForDestructiveModes() {
+ Map<String, Object> config = sinkConfig();
+ config.put(FileBaseSinkOptions.FILE_PATH.key(),
"/warehouse/../orders");
+ config.put(FileBaseSinkOptions.SCHEMA_SAVE_MODE.key(),
"RECREATE_SCHEMA");
+ config.put(FileBaseSinkOptions.DATA_SAVE_MODE.key(), "DROP_DATA");
+
+ assertDoesNotThrow(() -> new
GcsFileSinkFactory().createSink(sinkContext(config)));
+ }
+
+ private static Map<String, Object> sinkConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put(FileBaseSinkOptions.FILE_PATH.key(), "/warehouse/orders");
+ config.put(GcsFileSinkOptions.BUCKET.key(), "gs://test-bucket");
+ config.put(FileBaseSinkOptions.FILE_FORMAT_TYPE.key(), "parquet");
+ return config;
+ }
+
+ private static boolean optionRuleContains(OptionRule optionRule, Option<?>
option) {
+ if (optionRule.getOptionalOptions().contains(option)) {
+ return true;
+ }
+ return optionRule.getRequiredOptions().stream()
+ .anyMatch(requiredOption ->
requiredOption.getOptions().contains(option));
+ }
+
+ private static void validate(Map<String, Object> config, OptionRule
optionRule) {
+
ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(optionRule);
+ }
+
+ private static TableSinkFactoryContext sinkContext(Map<String, Object>
config) {
+ return TableSinkFactoryContext.replacePlaceholderAndCreate(
+ CatalogTableUtil.buildSimpleTextTable(),
+ ReadonlyConfig.fromMap(config),
+ Thread.currentThread().getContextClassLoader(),
+ Collections.emptyList());
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalogTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalogTest.java
new file mode 100644
index 0000000000..9698c36303
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-gcs/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/gcs/catalog/GcsFileCatalogTest.java
@@ -0,0 +1,106 @@
+/*
+ * 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.file.gcs.catalog;
+
+import org.apache.seatunnel.api.sink.DataSaveMode;
+import org.apache.seatunnel.api.sink.DefaultSaveModeHandler;
+import org.apache.seatunnel.api.sink.SchemaSaveMode;
+import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
+import org.apache.seatunnel.common.exception.SeaTunnelRuntimeException;
+import
org.apache.seatunnel.connectors.seatunnel.file.hadoop.HadoopFileSystemProxy;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
+import org.mockito.Mockito;
+
+class GcsFileCatalogTest {
+
+ private static final String SINK_PATH =
"gs://test-bucket/warehouse/orders";
+
+ @Test
+ void shouldRejectExistingPartitionDataInErrorMode() throws Exception {
+ HadoopFileSystemProxy fileSystemProxy =
Mockito.mock(HadoopFileSystemProxy.class);
+ Mockito.when(fileSystemProxy.hasAnyFile(SINK_PATH,
true)).thenReturn(true);
+ GcsFileCatalog catalog = new GcsFileCatalog(fileSystemProxy,
SINK_PATH, "GcsFile");
+ DefaultSaveModeHandler handler =
+ new DefaultSaveModeHandler(
+ SchemaSaveMode.IGNORE,
+ DataSaveMode.ERROR_WHEN_DATA_EXISTS,
+ catalog,
+ CatalogTableUtil.buildSimpleTextTable(),
+ null);
+
+ Assertions.assertThrows(SeaTunnelRuntimeException.class,
handler::handleDataSaveMode);
+
+ Mockito.verify(fileSystemProxy).hasAnyFile(SINK_PATH, true);
+ Mockito.verifyNoMoreInteractions(fileSystemProxy);
+ }
+
+ @Test
+ void shouldPreserveExistingDataInAppendMode() {
+ HadoopFileSystemProxy fileSystemProxy =
Mockito.mock(HadoopFileSystemProxy.class);
+ GcsFileCatalog catalog = new GcsFileCatalog(fileSystemProxy,
SINK_PATH, "GcsFile");
+ DefaultSaveModeHandler handler =
+ new DefaultSaveModeHandler(
+ SchemaSaveMode.IGNORE,
+ DataSaveMode.APPEND_DATA,
+ catalog,
+ CatalogTableUtil.buildSimpleTextTable(),
+ null);
+
+ handler.handleDataSaveMode();
+
+ Mockito.verifyNoInteractions(fileSystemProxy);
+ }
+
+ @Test
+ void shouldNotReportDataWhenListingIsEmpty() throws Exception {
+ HadoopFileSystemProxy fileSystemProxy =
Mockito.mock(HadoopFileSystemProxy.class);
+ Mockito.when(fileSystemProxy.hasAnyFile(SINK_PATH,
true)).thenReturn(false);
+ GcsFileCatalog catalog = new GcsFileCatalog(fileSystemProxy,
SINK_PATH, "GcsFile");
+
+ Assertions.assertFalse(catalog.isExistsData(null));
+
+ Mockito.verify(fileSystemProxy).hasAnyFile(SINK_PATH, true);
+ }
+
+ @Test
+ void shouldTruncateOnlyTheConfiguredPrefix() throws Exception {
+ HadoopFileSystemProxy fileSystemProxy =
Mockito.mock(HadoopFileSystemProxy.class);
+ GcsFileCatalog catalog = new GcsFileCatalog(fileSystemProxy,
SINK_PATH, "GcsFile");
+
+ catalog.truncateTable(null, false);
+
+ InOrder operations = Mockito.inOrder(fileSystemProxy);
+ operations.verify(fileSystemProxy).deleteFile(SINK_PATH);
+ operations.verify(fileSystemProxy).createDir(SINK_PATH);
+ Mockito.verifyNoMoreInteractions(fileSystemProxy);
+ }
+
+ @Test
+ void shouldCheckForExistingDataRecursively() throws Exception {
+ HadoopFileSystemProxy fileSystemProxy =
Mockito.mock(HadoopFileSystemProxy.class);
+ Mockito.when(fileSystemProxy.hasAnyFile(SINK_PATH,
true)).thenReturn(true);
+ GcsFileCatalog catalog = new GcsFileCatalog(fileSystemProxy,
SINK_PATH, "GcsFile");
+
+ Assertions.assertTrue(catalog.isExistsData(null));
+
+ Mockito.verify(fileSystemProxy).hasAnyFile(SINK_PATH, true);
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/gcs/GcsFileIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/gcs/GcsFileIT.java
index 7bcfb40c1f..37b1daf8aa 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/gcs/GcsFileIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/gcs/GcsFileIT.java
@@ -57,9 +57,12 @@ public class GcsFileIT extends TestSuiteBase implements
TestResource {
@Override
public void startUp() {
DockerImageName image = DockerImageName.parse(FAKE_GCS_IMAGE);
+ // Preserve the directory-marker objects required by Hadoop mkdirs and
rename.
fakeGcs =
new GenericContainer<>(image)
.withCommand(
+ "-backend",
+ "memory",
"-scheme",
"http",
"-port",
@@ -101,4 +104,15 @@ public class GcsFileIT extends TestSuiteBase implements
TestResource {
Container.ExecResult result = container.executeJob(JOB_CONFIG);
Assertions.assertEquals(0, result.getExitCode(), result.getStderr());
}
+
+ @TestTemplate
+ public void testWritePartitionedJsonAndReplaceExistingData(TestContainer
container)
+ throws Exception {
+ for (int attempt = 0; attempt < 2; attempt++) {
+ Container.ExecResult write =
container.executeJob("/gcs/gcs_file_to_gcs.conf");
+ Assertions.assertEquals(0, write.getExitCode(), write.getStderr());
+ Container.ExecResult read =
container.executeJob("/gcs/gcs_sink_to_assert.conf");
+ Assertions.assertEquals(0, read.getExitCode(), read.getStderr());
+ }
+ }
}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/resources/gcs/gcs_file_to_gcs.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/resources/gcs/gcs_file_to_gcs.conf
new file mode 100644
index 0000000000..571d0506f3
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/resources/gcs/gcs_file_to_gcs.conf
@@ -0,0 +1,67 @@
+#
+# 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.
+#
+
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 1
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ GcsFile {
+ bucket = "gs://seatunnel-gcs-test"
+ path = "/input"
+ file_format_type = "json"
+ hadoop_gcs_properties = {
+ "fs.gs.storage.root.url" = "http://fake-gcs:4443/"
+ "fs.gs.auth.service.account.enable" = "false"
+ "fs.gs.auth.null.enable" = "true"
+ "fs.gs.project.id" = "seatunnel-test"
+ }
+ schema = {
+ fields {
+ id = int
+ name = string
+ }
+ }
+ }
+}
+
+sink {
+ GcsFile {
+ bucket = "gs://seatunnel-gcs-test"
+ path = "/sink-output"
+ tmp_path = "/sink-staging"
+ file_format_type = "json"
+ data_save_mode = "DROP_DATA"
+ have_partition = true
+ partition_by = ["id"]
+ is_partition_field_write_in_file = true
+ hadoop_gcs_properties = {
+ # fake-gcs-server does not implement the native GCS move endpoint.
+ "fs.gs.operation.move.enable" = "false"
+ "fs.gs.storage.root.url" = "http://fake-gcs:4443/"
+ "fs.gs.auth.service.account.enable" = "false"
+ "fs.gs.auth.null.enable" = "true"
+ "fs.gs.project.id" = "seatunnel-test"
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/resources/gcs/gcs_sink_to_assert.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/resources/gcs/gcs_sink_to_assert.conf
new file mode 100644
index 0000000000..632e9944d5
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-gcs-e2e/src/test/resources/gcs/gcs_sink_to_assert.conf
@@ -0,0 +1,63 @@
+#
+# 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.
+#
+
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 1
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ GcsFile {
+ bucket = "gs://seatunnel-gcs-test"
+ path = "/sink-output"
+ file_format_type = "json"
+ # Partition columns are already included in the JSON output.
+ parse_partition_from_path = false
+ hadoop_gcs_properties = {
+ "fs.gs.storage.root.url" = "http://fake-gcs:4443/"
+ "fs.gs.auth.service.account.enable" = "false"
+ "fs.gs.auth.null.enable" = "true"
+ "fs.gs.project.id" = "seatunnel-test"
+ }
+ schema = {
+ fields {
+ id = int
+ name = string
+ }
+ }
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ { rule_type = MIN_ROW, rule_value = 2 },
+ { rule_type = MAX_ROW, rule_value = 2 }
+ ]
+ field_rules = [
+ { field_name = id, field_type = int, field_value = [{ rule_type =
NOT_NULL }] },
+ { field_name = name, field_type = string, field_value = [{ rule_type =
NOT_NULL }] }
+ ]
+ }
+ }
+}