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 99e53aad86 [Feature][Connector-V2] Add BosFile source and sink
connector (#11952)
99e53aad86 is described below
commit 99e53aad863b7991daf3aafffaf2e1b69882346d
Author: 雷炯 <[email protected]>
AuthorDate: Sun Aug 30 16:10:11 2026 +0000
[Feature][Connector-V2] Add BosFile source and sink connector (#11952)
Co-authored-by: Cursor <[email protected]>
---
.gitignore | 3 +
config/plugin_config | 1 +
docs/en/connectors/changelog/connector-file-bos.md | 8 +
.../connectors/common-options/sink-write-modes.md | 3 +-
docs/en/connectors/sink/BosFile.md | 158 +++++++++++++++++
docs/en/connectors/source/BosFile.md | 193 +++++++++++++++++++++
docs/zh/connectors/changelog/connector-file-bos.md | 8 +
.../connectors/common-options/sink-write-modes.md | 3 +-
docs/zh/connectors/sink/BosFile.md | 129 ++++++++++++++
docs/zh/connectors/source/BosFile.md | 136 +++++++++++++++
plugin-mapping.properties | 2 +
.../seatunnel/file/config/FileSystemType.java | 3 +-
.../connector-file-bos/lib/README.md | 18 ++
.../connector-file/connector-file-bos/pom.xml | 53 ++++++
.../seatunnel/file/bos/config/BosConf.java | 70 ++++++++
.../file/bos/config/BosFileBaseOptions.java | 39 +++++
.../file/bos/config/BosFileSinkOptions.java} | 10 +-
.../file/bos/config/BosFileSourceOptions.java} | 10 +-
.../seatunnel/file/bos/sink/BosFileSink.java | 42 +++++
.../file/bos/sink/BosFileSinkFactory.java | 120 +++++++++++++
.../seatunnel/file/bos/source/BosFileSource.java} | 33 ++--
.../file/bos/source/BosFileSourceFactory.java | 122 +++++++++++++
.../services/org.apache.hadoop.fs.FileSystem | 16 ++
.../seatunnel/file/bos/BosFileFactoryTest.java} | 21 ++-
seatunnel-connectors-v2/connector-file/pom.xml | 1 +
seatunnel-connectors-v2/connector-hive/pom.xml | 5 +
.../seatunnel/hive/storage/BOSStorage.java | 89 ++++++++++
.../seatunnel/hive/storage/StorageFactory.java | 2 +
.../seatunnel/hive/storage/StorageType.java | 1 +
.../seatunnel/hive/storage/BosStorageTest.java | 80 +++++++++
.../seatunnel/hive/storage/StorageFactoryTest.java | 1 +
.../src/test/resources/bos/core-site.xml | 37 ++++
seatunnel-dist/pom.xml | 7 +
.../connector-file-bos-e2e/pom.xml | 48 +++++
.../e2e/connector/file/bos/BosFileIT.java | 71 ++++++++
.../test/resources/excel/bos_excel_to_assert.conf | 118 +++++++++++++
.../test/resources/excel/fake_to_bos_excel.conf | 84 +++++++++
.../resources/json/bos_file_json_to_assert.conf | 116 +++++++++++++
.../test/resources/json/fake_to_bos_file_json.conf | 85 +++++++++
.../test/resources/orc/bos_file_orc_to_assert.conf | 82 +++++++++
.../test/resources/orc/fake_to_bos_file_orc.conf | 86 +++++++++
.../parquet/bos_file_parquet_to_assert.conf | 82 +++++++++
.../parquet/fake_to_bos_file_parquet.conf | 86 +++++++++
.../resources/text/bos_file_text_to_assert.conf | 116 +++++++++++++
.../test/resources/text/fake_to_bos_file_text.conf | 86 +++++++++
seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml | 1 +
46 files changed, 2442 insertions(+), 43 deletions(-)
diff --git a/.gitignore b/.gitignore
index dfb9eefda8..b1eba9b18d 100644
--- a/.gitignore
+++ b/.gitignore
@@ -18,6 +18,9 @@ target/
/logs
logs.zip
+# Local BOS integration test configs (may reference credentials)
+config/bos-local-test/
+
# Intellij Idea files
.idea/
*.iml
diff --git a/config/plugin_config b/config/plugin_config
index e2928d04bd..875e405d46 100644
--- a/config/plugin_config
+++ b/config/plugin_config
@@ -47,6 +47,7 @@ connector-file-jindo-oss
connector-file-s3
connector-file-sftp
connector-file-obs
+connector-file-bos
connector-fluss
connector-google-sheets
connector-google-firestore
diff --git a/docs/en/connectors/changelog/connector-file-bos.md
b/docs/en/connectors/changelog/connector-file-bos.md
new file mode 100644
index 0000000000..b3a765114c
--- /dev/null
+++ b/docs/en/connectors/changelog/connector-file-bos.md
@@ -0,0 +1,8 @@
+<details><summary> Change Log </summary>
+
+| Change | Commit | Version |
+| --- | --- | --- |
+| [Improve][Connector-V2] Add Hive BOSStorage and align BosFile e2e/docs with
CosFile | https://github.com/apache/seatunnel/pull/11952 | dev |
+| [Feature][Connector-V2] Add BosFile source and sink for Baidu Object Storage
| https://github.com/apache/seatunnel/pull/11952 | dev |
+
+</details>
diff --git a/docs/en/connectors/common-options/sink-write-modes.md
b/docs/en/connectors/common-options/sink-write-modes.md
index 06b2d871f3..1ba79e5ac1 100644
--- a/docs/en/connectors/common-options/sink-write-modes.md
+++ b/docs/en/connectors/common-options/sink-write-modes.md
@@ -100,6 +100,7 @@ Current connector support differs by file connector:
| OssFile | Yes | Handles existing OSS paths and objects through the file sink
save mode flow. |
| ObsFile | No | The current sink option rule does not expose
`schema_save_mode` or `data_save_mode`. |
| CosFile | No | The current sink option rule does not expose
`schema_save_mode` or `data_save_mode`. |
+| BosFile | No | The current sink option rule does not expose
`schema_save_mode` or `data_save_mode`. |
If a file connector page does not list `schema_save_mode` or `data_save_mode`,
do not assume the option is accepted by that connector.
@@ -174,7 +175,7 @@ Check whether the sink uses `query`. JDBC custom query mode
does not apply save
### A file sink rejects `data_save_mode`
-Check the specific connector option table. `S3File`, `OssFile`, `HdfsFile`,
`FtpFile`, `SftpFile`, and `LocalFile` expose file save mode options. `ObsFile`
and `CosFile` currently do not.
+Check the specific connector option table. `S3File`, `OssFile`, `HdfsFile`,
`FtpFile`, `SftpFile`, and `LocalFile` expose file save mode options.
`ObsFile`, `CosFile`, and `BosFile` currently do not.
### I only want to create the target table
diff --git a/docs/en/connectors/sink/BosFile.md
b/docs/en/connectors/sink/BosFile.md
new file mode 100644
index 0000000000..d88f61b950
--- /dev/null
+++ b/docs/en/connectors/sink/BosFile.md
@@ -0,0 +1,158 @@
+import ChangeLog from '../changelog/connector-file-bos.md';
+
+# BosFile
+
+> BOS file sink connector
+
+## Support Those Engines
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## Description
+
+Output data to Baidu Cloud BOS (Baidu Object Storage) via the BOS HDFS SDK.
+
+:::tip
+
+If you use Spark/Flink, in order to use this connector you must ensure your
Spark/Flink cluster already integrated Hadoop. The tested Hadoop version is 2.x.
+
+If you use SeaTunnel Engine, Hadoop jars are bundled under
`${SEATUNNEL_HOME}/lib`.
+
+To use this connector you need to put `bos-hdfs-sdk` (>= 1.0.4-community) into
`${SEATUNNEL_HOME}/lib`. Download:
[bos-hdfs-sdk-1.0.4-community.jar.zip](https://sdk.bce.baidu.com/console-sdk/bos-hdfs-sdk-1.0.4-community.jar.zip).
+
+:::
+
+## Key Features
+
+- [x]
[multimodal](../../introduction/concepts/connector-v2-features.md#multimodal)
+
+ Use binary file format to read and write files in any format, such as
videos, pictures, etc. In short, any files can be synchronized to the target
place.
+
+- [x] [exactly-once](../../introduction/concepts/connector-v2-features.md)
+
+ By default, we use 2PC commit to ensure `exactly-once`
+
+- [ ] [cdc](../../introduction/concepts/connector-v2-features.md)
+- [x] [support multiple table
write](../../introduction/concepts/connector-v2-features.md)
+- [ ] [timer flush](../../introduction/concepts/connector-v2-features.md)
+
+- [x] file format type
+ - [x] text
+ - [x] csv
+ - [x] parquet
+ - [x] orc
+ - [x] json
+ - [x] excel
+ - [x] xml
+ - [x] binary
+ - [x] canal_json
+ - [x] debezium_json
+ - [x] maxwell_json
+
+## Options
+
+| Name | Type | Required | Default
| Description
|
+|---------------------------------------|---------|----------|--------------------------------------------|---------------------------------------------------------------------------------------------------------------------------------------------------------------------------------|
+| path | string | yes | -
| The target directory the sink writes to inside the
bucket.
|
+| tmp_path | string | no | /tmp/seatunnel
| The result file will write to a tmp path first and
then use `mv` to submit tmp dir to target dir. Needs a BOS dir.
|
+| bucket | string | yes | -
| The BOS bucket address, for example
`bos://my-bucket`.
|
+| access_key | string | yes | -
| The Baidu Cloud BOS access key.
|
+| secret_key | string | yes | -
| The Baidu Cloud BOS secret key.
|
+| endpoint | string | yes | -
| The BOS endpoint, for example
`http://bj.bcebos.com`.
|
+| custom_filename | boolean | no | false
| Whether you need custom the filename.
|
+| file_name_expression | string | no |
"${transactionId}" | Only used when custom_filename is
true.
|
+| filename_time_format | string | no | "yyyy.MM.dd"
| Only used when custom_filename is true.
|
+| file_format_type | string | no | "csv"
| File format type, supported: `text`, `csv`,
`parquet`, `orc`, `json`, `excel`, `xml`, `binary`, `canal_json`,
`debezium_json`, `maxwell_json`. |
+| filename_extension | string | no | -
| Override the default file name extensions with
custom file name extensions. E.g. `.xml`, `.json`, `dat`, `.customtype`
|
+| field_delimiter | string | no | '\001' for text
and ',' for csv | Only used when file_format_type is text and csv.
|
+| row_delimiter | string | no | "\n"
| Only used when file_format_type is `text`, `csv`
and `json`.
|
+| have_partition | boolean | no | false
| Whether you need processing partitions.
|
+| partition_by | array | no | -
| Only used when have_partition is true.
|
+| partition_dir_expression | string | no |
"${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/" | Only used when have_partition is
true.
|
+| is_partition_field_write_in_file | boolean | no | false
| Only used when have_partition is true.
|
+| sink_columns | array | no |
| When this parameter is empty, all fields are sink
columns.
|
+| is_enable_transaction | boolean | no | true
| If `true`, data will not be lost or duplicated
when written to the target directory. When `true`, `${transactionId}_` is
automatically prefixed to the file name. |
+| batch_size | int | no | 1000000
| The maximum number of rows in a file. For
SeaTunnel Engine the file row count is jointly decided by `batch_size` and
`checkpoint.interval`. |
+| compress_codec | string | no | none
| The compress codec of files. Excel does not
support any compression format.
|
+| xml_root_tag | string | no | RECORDS
| Only used when file_format is xml.
|
+| xml_row_tag | string | no | RECORD
| Only used when file_format is xml.
|
+| xml_use_attr_format | boolean | no | -
| Only used when file_format is xml.
|
+| single_file_mode | boolean | no | false
| Each parallelism will only output one file. When
this parameter is turned on, batch_size will not take effect. The output file
name does not have a file block suffix. |
+| create_empty_file_when_no_data | boolean | no | false
| When there is no data synchronization upstream,
the corresponding data files are still generated.
|
+| parquet_avro_write_timestamp_as_int96 | boolean | no | false
| Only used when file_format is parquet.
|
+| parquet_avro_write_fixed_as_int96 | array | no | -
| Only used when file_format is parquet.
|
+| encoding | string | no | "UTF-8"
| Only used when file_format_type is
json,text,csv,xml.
|
+| common-options | object | no | -
| Sink plugin common parameters, please refer to
[Sink Common Options](../common-options/sink-common-options.md) for details.
|
+
+## Example
+
+For text file format with `have_partition`, `custom_filename` and
`sink_columns`:
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+sink {
+ BosFile {
+ path = "/sink"
+ bucket = "bos://sink-bucket"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ file_format_type = "text"
+ field_delimiter = "\t"
+ row_delimiter = "\n"
+ have_partition = true
+ partition_by = ["age"]
+ partition_dir_expression = "${k0}=${v0}"
+ is_partition_field_write_in_file = true
+ custom_filename = true
+ file_name_expression = "${transactionId}_${now}"
+ filename_time_format = "yyyy.MM.dd"
+ sink_columns = ["name", "age"]
+ is_enable_transaction = true
+ }
+}
+```
+
+For parquet file format:
+
+```hocon
+sink {
+ BosFile {
+ path = "/sink"
+ bucket = "bos://sink-bucket"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ file_format_type = "parquet"
+ is_enable_transaction = true
+ }
+}
+```
+
+Simple text sink:
+
+```hocon
+sink {
+ BosFile {
+ bucket = "bos://sink-bucket"
+ path = "/warehouse/table/"
+ file_format_type = "text"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ row_delimiter = "\n"
+ field_delimiter = ","
+ is_enable_transaction = true
+ }
+}
+```
+
+## Changelog
+
+<ChangeLog />
diff --git a/docs/en/connectors/source/BosFile.md
b/docs/en/connectors/source/BosFile.md
new file mode 100644
index 0000000000..bf149e6065
--- /dev/null
+++ b/docs/en/connectors/source/BosFile.md
@@ -0,0 +1,193 @@
+import ChangeLog from '../changelog/connector-file-bos.md';
+
+# BosFile
+
+> BOS file source connector
+
+## Support Those Engines
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## Key features
+
+- [x] [batch](../../introduction/concepts/connector-v2-features.md)
+- [ ] [stream](../../introduction/concepts/connector-v2-features.md)
+- [x]
[multimodal](../../introduction/concepts/connector-v2-features.md#multimodal)
+
+ Use binary file format to read and write files in any format, such as
videos, pictures, etc. In short, any files can be synchronized to the target
place.
+
+- [x] [exactly-once](../../introduction/concepts/connector-v2-features.md)
+
+ Read all the data in a split in a pollNext call. What splits are read will
be saved in snapshot.
+
+- [x] [column projection](../../introduction/concepts/connector-v2-features.md)
+- [x] [parallelism](../../introduction/concepts/connector-v2-features.md)
+- [ ] [support user-defined
split](../../introduction/concepts/connector-v2-features.md)
+- [x] file format type
+ - [x] text
+ - [x] csv
+ - [x] parquet
+ - [x] orc
+ - [x] json
+ - [x] excel
+ - [x] xml
+ - [x] binary
+ - [x] markdown
+ - [x] pdf
+
+## Description
+
+Read data from Baidu Cloud BOS (Baidu Object Storage) via the BOS HDFS SDK.
+
+:::tip
+
+If you use Spark/Flink, in order to use this connector you must ensure your
Spark/Flink cluster already integrated Hadoop. The tested Hadoop version is 2.x.
+
+If you use SeaTunnel Engine, Hadoop jars are bundled under
`${SEATUNNEL_HOME}/lib`.
+
+To use this connector you need to put `bos-hdfs-sdk` (>= 1.0.4-community) into
`${SEATUNNEL_HOME}/lib`. Download:
[bos-hdfs-sdk-1.0.4-community.jar.zip](https://sdk.bce.baidu.com/console-sdk/bos-hdfs-sdk-1.0.4-community.jar.zip).
See `connector-file-bos/lib/README.md` for details.
+
+:::
+
+## Options
+
+| name | type | required | default value
|
+|----------------------------|---------|----------|-----------------------------|
+| path | string | yes | -
|
+| file_format_type | string | yes | -
|
+| bucket | string | yes | -
|
+| access_key | string | yes | -
|
+| secret_key | string | yes | -
|
+| endpoint | string | yes | -
|
+| read_columns | list | no | -
|
+| delimiter/field_delimiter | string | no | \001 for text and , for
csv |
+| row_delimiter | string | no | \n
|
+| parse_partition_from_path | boolean | no | true
|
+| skip_header_row_number | long | no | 0
|
+| date_format | string | no | yyyy-MM-dd
|
+| datetime_format | string | no | yyyy-MM-dd HH:mm:ss
|
+| time_format | string | no | HH:mm:ss
|
+| schema | config | no | -
|
+| sheet_name | string | no | -
|
+| excel_engine | string | no | POI
|
+| poi_excel_max_file_size | long | no | 52428800
|
+| xml_row_tag | string | no | -
|
+| xml_use_attr_format | boolean | no | -
|
+| csv_use_header_line | boolean | no | false
|
+| file_filter_pattern | string | no | -
|
+| filename_extension | string | no | -
|
+| compress_codec | string | no | none
|
+| archive_compress_codec | string | no | none
|
+| encoding | string | no | UTF-8
|
+| binary_chunk_size | int | no | 1024
|
+| binary_complete_file_mode | boolean | no | false
|
+| common-options | | no | -
|
+| file_filter_modified_start | string | no | -
|
+| file_filter_modified_end | string | no | -
|
+| quote_char | string | no | "
|
+| escape_char | string | no | -
|
+| recursive_file_scan | boolean | no | true
|
+| sort_files_by_modification_time | boolean | no | false
|
+
+### path [string]
+
+The source file path under the bucket.
+
+### bucket [string]
+
+The BOS bucket address, for example `bos://my-bucket`.
+
+### access_key [string]
+
+The Baidu Cloud BOS access key.
+
+### secret_key [string]
+
+The Baidu Cloud BOS secret key.
+
+### endpoint [string]
+
+The BOS endpoint, for example `http://bj.bcebos.com`.
+
+### file_format_type [string]
+
+Supported file types: `text`, `csv`, `parquet`, `orc`, `json`, `excel`, `xml`,
`binary`, `markdown`, `pdf`.
+
+### common options
+
+Source plugin common parameters, please refer to [Source Common
Options](../common-options/source-common-options.md) for details.
+
+## Example
+
+```hocon
+source {
+ BosFile {
+ bucket = "bos://source-bucket"
+ path = "/warehouse/table/"
+ file_format_type = "orc"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ }
+}
+```
+
+### Transfer Binary File
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "BATCH"
+}
+
+source {
+ BosFile {
+ bucket = "bos://source-bucket"
+ path = "/read/binary/"
+ file_format_type = "binary"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ binary_chunk_size = 2048
+ }
+}
+
+sink {
+ BosFile {
+ bucket = "bos://sink-bucket"
+ path = "/write/binary/"
+ file_format_type = "binary"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ }
+}
+```
+
+### Filter File
+
+```hocon
+source {
+ BosFile {
+ bucket = "bos://source-bucket"
+ path = "/read/data/"
+ file_format_type = "text"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ file_filter_pattern = "abc[DX]*.*"
+ schema {
+ fields {
+ id = int
+ name = string
+ }
+ }
+ }
+}
+```
+
+## Changelog
+
+<ChangeLog />
diff --git a/docs/zh/connectors/changelog/connector-file-bos.md
b/docs/zh/connectors/changelog/connector-file-bos.md
new file mode 100644
index 0000000000..bb6921f843
--- /dev/null
+++ b/docs/zh/connectors/changelog/connector-file-bos.md
@@ -0,0 +1,8 @@
+<details><summary> 变更日志 </summary>
+
+| 变更 | Commit | 版本 |
+| --- | --- | --- |
+| [Improve][Connector-V2] 新增 Hive BOSStorage,对齐 BosFile e2e/文档至 CosFile |
https://github.com/apache/seatunnel/pull/11952 | dev |
+| [Feature][Connector-V2] 新增 BosFile Source/Sink 连接器 |
https://github.com/apache/seatunnel/pull/11952 | dev |
+
+</details>
diff --git a/docs/zh/connectors/common-options/sink-write-modes.md
b/docs/zh/connectors/common-options/sink-write-modes.md
index b486cd9e4b..09f70e0c8b 100644
--- a/docs/zh/connectors/common-options/sink-write-modes.md
+++ b/docs/zh/connectors/common-options/sink-write-modes.md
@@ -100,6 +100,7 @@ File Sink 写的是文件,因此不使用 `generate_sink_sql`、`query` 或数
| OssFile | 是 | 通过 File Sink save mode 流程处理 OSS 路径和对象。 |
| ObsFile | 否 | 当前 sink option rule 没有暴露 `schema_save_mode` 或
`data_save_mode`。 |
| CosFile | 否 | 当前 sink option rule 没有暴露 `schema_save_mode` 或
`data_save_mode`。 |
+| BosFile | 否 | 当前 sink option rule 没有暴露 `schema_save_mode` 或
`data_save_mode`。 |
如果某个文件 connector 页面没有列出 `schema_save_mode` 或 `data_save_mode`,不要默认认为该
connector 可以接收这些参数。
@@ -174,7 +175,7 @@ sink {
### File Sink 不接受 `data_save_mode`
-检查具体 connector
参数表。`S3File`、`OssFile`、`HdfsFile`、`FtpFile`、`SftpFile`、`LocalFile` 暴露文件 save
mode 参数;`ObsFile` 和 `CosFile` 当前未暴露。
+检查具体 connector
参数表。`S3File`、`OssFile`、`HdfsFile`、`FtpFile`、`SftpFile`、`LocalFile` 暴露文件 save
mode 参数;`ObsFile`、`CosFile` 和 `BosFile` 当前未暴露。
### 我只想建表,不想抽取数据
diff --git a/docs/zh/connectors/sink/BosFile.md
b/docs/zh/connectors/sink/BosFile.md
new file mode 100644
index 0000000000..f904891984
--- /dev/null
+++ b/docs/zh/connectors/sink/BosFile.md
@@ -0,0 +1,129 @@
+import ChangeLog from '../changelog/connector-file-bos.md';
+
+# BosFile
+
+> BOS 文件 Sink 连接器
+
+## 支持引擎
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## 描述
+
+通过 BOS HDFS SDK 将数据写入百度智能云 BOS。
+
+:::tip
+
+使用 Spark/Flink 时,需确保集群已集成 Hadoop 2.x。
+
+使用 SeaTunnel Engine 时,Hadoop 相关 jar 已包含在 `${SEATUNNEL_HOME}/lib` 中。
+
+使用本连接器需将 `bos-hdfs-sdk`(>= 1.0.4-community)放入
`${SEATUNNEL_HOME}/lib`。下载:[bos-hdfs-sdk-1.0.4-community.jar.zip](https://sdk.bce.baidu.com/console-sdk/bos-hdfs-sdk-1.0.4-community.jar.zip)。
+
+:::
+
+## 主要特性
+
+- [x] [多模态](../../introduction/concepts/connector-v2-features.md#multimodal)
+
+ 使用 binary 格式读写任意类型文件,可将任意文件同步到目标位置。
+
+- [x] [精确一次](../../introduction/concepts/connector-v2-features.md)
+
+ 默认通过 2PC 提交保证 exactly-once
+
+- [ ] [cdc](../../introduction/concepts/connector-v2-features.md)
+- [x] [支持多表写入](../../introduction/concepts/connector-v2-features.md)
+- [ ] [定时 flush](../../introduction/concepts/connector-v2-features.md)
+
+- [x] 文件格式
+ - [x] text
+ - [x] csv
+ - [x] parquet
+ - [x] orc
+ - [x] json
+ - [x] excel
+ - [x] xml
+ - [x] binary
+ - [x] canal_json
+ - [x] debezium_json
+ - [x] maxwell_json
+
+## 配置项
+
+| 名称 | 类型 | 必填 | 默认值
| 描述 |
+|---------------------------------------|---------|------|--------------------------------------------|------|
+| path | string | 是 | -
| bucket 内写入目录 |
+| tmp_path | string | 否 | /tmp/seatunnel
| 先写入临时目录,再通过 mv 提交到目标目录 |
+| bucket | string | 是 | -
| BOS bucket,例如 `bos://my-bucket` |
+| access_key | string | 是 | -
| BOS Access Key |
+| secret_key | string | 是 | -
| BOS Secret Key |
+| endpoint | string | 是 | -
| BOS Endpoint,例如 `http://bj.bcebos.com` |
+| custom_filename | boolean | 否 | false
| 是否自定义文件名 |
+| file_name_expression | string | 否 | "${transactionId}"
| custom_filename 为 true 时使用 |
+| filename_time_format | string | 否 | "yyyy.MM.dd"
| custom_filename 为 true 时使用 |
+| file_format_type | string | 否 | "csv"
| 支持 text、csv、parquet、orc、json、excel、xml、binary 等 |
+| filename_extension | string | 否 | -
| 自定义文件扩展名 |
+| field_delimiter | string | 否 | text 为 \001,csv 为 ,
| text/csv 格式使用 |
+| row_delimiter | string | 否 | "\n"
| text/csv/json 格式使用 |
+| have_partition | boolean | 否 | false
| 是否按分区写入 |
+| partition_by | array | 否 | -
| have_partition 为 true 时使用 |
+| partition_dir_expression | string | 否 |
"${k0}=${v0}/${k1}=${v1}/.../${kn}=${vn}/" | have_partition 为 true 时使用 |
+| is_partition_field_write_in_file | boolean | 否 | false
| have_partition 为 true 时使用 |
+| sink_columns | array | 否 |
| 为空时写入全部字段 |
+| is_enable_transaction | boolean | 否 | true
| 开启后保证写入不丢不重,文件名自动加 `${transactionId}_` 前缀 |
+| batch_size | int | 否 | 1000000
| 单文件最大行数 |
+| compress_codec | string | 否 | none
| 压缩格式 |
+| xml_root_tag | string | 否 | RECORDS
| xml 格式使用 |
+| xml_row_tag | string | 否 | RECORD
| xml 格式使用 |
+| xml_use_attr_format | boolean | 否 | -
| xml 格式使用 |
+| single_file_mode | boolean | 否 | false
| 每个并行度只输出一个文件 |
+| create_empty_file_when_no_data | boolean | 否 | false
| 上游无数据时仍生成空文件 |
+| parquet_avro_write_timestamp_as_int96 | boolean | 否 | false
| parquet 格式使用 |
+| parquet_avro_write_fixed_as_int96 | array | 否 | -
| parquet 格式使用 |
+| encoding | string | 否 | "UTF-8"
| json/text/csv/xml 格式使用 |
+
+## 示例
+
+分区写入 text 文件:
+
+```hocon
+sink {
+ BosFile {
+ path = "/sink"
+ bucket = "bos://sink-bucket"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ file_format_type = "text"
+ have_partition = true
+ partition_by = ["age"]
+ partition_dir_expression = "${k0}=${v0}"
+ is_enable_transaction = true
+ }
+}
+```
+
+简单 text 写入:
+
+```hocon
+sink {
+ BosFile {
+ bucket = "bos://sink-bucket"
+ path = "/warehouse/table/"
+ file_format_type = "text"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ row_delimiter = "\n"
+ field_delimiter = ","
+ is_enable_transaction = true
+ }
+}
+```
+
+## 变更日志
+
+<ChangeLog />
diff --git a/docs/zh/connectors/source/BosFile.md
b/docs/zh/connectors/source/BosFile.md
new file mode 100644
index 0000000000..a4f8ee6b0b
--- /dev/null
+++ b/docs/zh/connectors/source/BosFile.md
@@ -0,0 +1,136 @@
+import ChangeLog from '../changelog/connector-file-bos.md';
+
+# BosFile
+
+> BOS 文件 Source 连接器
+
+## 支持引擎
+
+> Spark<br/>
+> Flink<br/>
+> SeaTunnel Zeta<br/>
+
+## 关键特性
+
+- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
+- [ ] [流处理](../../introduction/concepts/connector-v2-features.md)
+- [x] [多模态](../../introduction/concepts/connector-v2-features.md#multimodal)
+
+ 使用 binary 格式读写任意类型文件(视频、图片等),可将任意文件同步到目标位置。
+
+- [x] [精确一次](../../introduction/concepts/connector-v2-features.md)
+
+ 一次 pollNext 读取 split 内全部数据,已读 split 会保存在 checkpoint 快照中。
+
+- [x] [列投影](../../introduction/concepts/connector-v2-features.md)
+- [x] [并行度](../../introduction/concepts/connector-v2-features.md)
+- [ ] [支持用户自定义 split](../../introduction/concepts/connector-v2-features.md)
+- [x] 文件格式
+ - [x] text
+ - [x] csv
+ - [x] parquet
+ - [x] orc
+ - [x] json
+ - [x] excel
+ - [x] xml
+ - [x] binary
+ - [x] markdown
+ - [x] pdf
+
+## 描述
+
+通过 BOS HDFS SDK 从百度智能云 BOS 读取文件数据。
+
+:::tip
+
+使用 Spark/Flink 时,需确保集群已集成 Hadoop 2.x。
+
+使用 SeaTunnel Engine 时,Hadoop 相关 jar 已包含在 `${SEATUNNEL_HOME}/lib` 中。
+
+使用本连接器需将 `bos-hdfs-sdk`(>= 1.0.4-community)放入
`${SEATUNNEL_HOME}/lib`。下载:[bos-hdfs-sdk-1.0.4-community.jar.zip](https://sdk.bce.baidu.com/console-sdk/bos-hdfs-sdk-1.0.4-community.jar.zip)。详见
`connector-file-bos/lib/README.md`。
+
+:::
+
+## 配置项
+
+| 名称 | 类型 | 必填 | 默认值 |
+|----------------------------|---------|------|-----------------------------|
+| path | string | 是 | - |
+| file_format_type | string | 是 | - |
+| bucket | string | 是 | - |
+| access_key | string | 是 | - |
+| secret_key | string | 是 | - |
+| endpoint | string | 是 | - |
+| read_columns | list | 否 | - |
+| delimiter/field_delimiter | string | 否 | text 为 \001,csv 为 , |
+| row_delimiter | string | 否 | \n |
+| parse_partition_from_path | boolean | 否 | true |
+| skip_header_row_number | long | 否 | 0 |
+| date_format | string | 否 | yyyy-MM-dd |
+| datetime_format | string | 否 | yyyy-MM-dd HH:mm:ss |
+| time_format | string | 否 | HH:mm:ss |
+| schema | config | 否 | - |
+| sheet_name | string | 否 | - |
+| excel_engine | string | 否 | POI |
+| poi_excel_max_file_size | long | 否 | 52428800 |
+| xml_row_tag | string | 否 | - |
+| xml_use_attr_format | boolean | 否 | - |
+| csv_use_header_line | boolean | 否 | false |
+| file_filter_pattern | string | 否 | - |
+| filename_extension | string | 否 | - |
+| compress_codec | string | 否 | none |
+| archive_compress_codec | string | 否 | none |
+| encoding | string | 否 | UTF-8 |
+| binary_chunk_size | int | 否 | 1024 |
+| binary_complete_file_mode | boolean | 否 | false |
+| file_filter_modified_start | string | 否 | - |
+| file_filter_modified_end | string | 否 | - |
+| quote_char | string | 否 | " |
+| escape_char | string | 否 | - |
+| recursive_file_scan | boolean | 否 | true |
+| sort_files_by_modification_time | boolean | 否 | false
|
+
+## 示例
+
+```hocon
+source {
+ BosFile {
+ bucket = "bos://source-bucket"
+ path = "/warehouse/table/"
+ file_format_type = "orc"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ }
+}
+```
+
+### 二进制文件同步
+
+```hocon
+source {
+ BosFile {
+ bucket = "bos://source-bucket"
+ path = "/read/binary/"
+ file_format_type = "binary"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ }
+}
+
+sink {
+ BosFile {
+ bucket = "bos://sink-bucket"
+ path = "/write/binary/"
+ file_format_type = "binary"
+ access_key = "your-access-key"
+ secret_key = "your-secret-key"
+ endpoint = "http://bj.bcebos.com"
+ }
+}
+```
+
+## 变更日志
+
+<ChangeLog />
diff --git a/plugin-mapping.properties b/plugin-mapping.properties
index dbafe0a8a3..841086a450 100644
--- a/plugin-mapping.properties
+++ b/plugin-mapping.properties
@@ -147,6 +147,8 @@ seatunnel.source.Oracle-CDC = connector-cdc-oracle
seatunnel.sink.Pulsar = connector-pulsar
seatunnel.source.ObsFile = connector-file-obs
seatunnel.sink.ObsFile = connector-file-obs
+seatunnel.source.BosFile = connector-file-bos
+seatunnel.sink.BosFile = connector-file-bos
seatunnel.source.Milvus = connector-milvus
seatunnel.sink.Milvus = connector-milvus
seatunnel.sink.ActiveMQ = connector-activemq
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/FileSystemType.java
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/FileSystemType.java
index c7a5c26aec..19c67151c0 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/FileSystemType.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/FileSystemType.java
@@ -28,7 +28,8 @@ public enum FileSystemType implements Serializable {
FTP("FtpFile"),
SFTP("SftpFile"),
S3("S3File"),
- OBS("ObsFile");
+ OBS("ObsFile"),
+ BOS("BosFile");
private final String fileSystemPluginName;
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-bos/lib/README.md
b/seatunnel-connectors-v2/connector-file/connector-file-bos/lib/README.md
new file mode 100644
index 0000000000..f4f999d2e8
--- /dev/null
+++ b/seatunnel-connectors-v2/connector-file/connector-file-bos/lib/README.md
@@ -0,0 +1,18 @@
+# BOS HDFS SDK (runtime dependency)
+
+`bos-hdfs-sdk` is **not published to Maven Central**. SeaTunnel does not
declare it as a Maven
+dependency; users must install the jar at **runtime** only.
+
+## Download
+
+Download `bos-hdfs-sdk-1.0.4-community.jar` from:
+
+https://sdk.bce.baidu.com/console-sdk/bos-hdfs-sdk-1.0.4-community.jar.zip
+
+## Install
+
+Copy the jar into `${SEATUNNEL_HOME}/lib` before running BosFile jobs (same as
documented in
+`docs/en/connectors/source/BosFile.md`).
+
+> Use **1.0.4+**. SDK 1.0.3 always calls `headBucket` during FileSystem init
and ignores
+> `fs.bos.bucket.hierarchy=false`, which fails when the AK/SK lacks
`HeadBucket` permission.
diff --git a/seatunnel-connectors-v2/connector-file/connector-file-bos/pom.xml
b/seatunnel-connectors-v2/connector-file/connector-file-bos/pom.xml
new file mode 100644
index 0000000000..607ebbd904
--- /dev/null
+++ b/seatunnel-connectors-v2/connector-file/connector-file-bos/pom.xml
@@ -0,0 +1,53 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+
+ 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.
+
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
+ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
+ <modelVersion>4.0.0</modelVersion>
+ <parent>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-file</artifactId>
+ <version>${revision}</version>
+ </parent>
+
+ <artifactId>connector-file-bos</artifactId>
+ <name>SeaTunnel : Connectors V2 : File : Bos</name>
+
+ <dependencies>
+
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-file-base</artifactId>
+ <version>${project.version}</version>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.flink</groupId>
+ <artifactId>flink-shaded-hadoop-2</artifactId>
+ <scope>provided</scope>
+ <exclusions>
+ <exclusion>
+ <groupId>org.apache.avro</groupId>
+ <artifactId>avro</artifactId>
+ </exclusion>
+ </exclusions>
+ </dependency>
+
+ </dependencies>
+
+</project>
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosConf.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosConf.java
new file mode 100644
index 0000000000..6bdbed2111
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosConf.java
@@ -0,0 +1,70 @@
+/*
+ * 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.bos.config;
+
+import org.apache.seatunnel.shade.com.typesafe.config.Config;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
+
+import java.util.HashMap;
+
+public class BosConf extends HadoopConf {
+ private static final String HDFS_IMPL =
"org.apache.hadoop.fs.bos.BaiduBosFileSystem";
+ private static final String SCHEMA = "bos";
+ private static final String ACCESS_KEY = "fs.bos.access.key";
+ private static final String SECRET_KEY = "fs.bos.secret.access.key";
+ private static final String ENDPOINT = "fs.bos.endpoint";
+ private static final String BUCKET_HIERARCHY = "fs.bos.bucket.hierarchy";
+
+ @Override
+ public String getFsHdfsImpl() {
+ return HDFS_IMPL;
+ }
+
+ @Override
+ public String getSchema() {
+ return SCHEMA;
+ }
+
+ public BosConf(String hdfsNameKey) {
+ super(hdfsNameKey);
+ }
+
+ public static HadoopConf buildWithConfig(Config config) {
+ HadoopConf hadoopConf = new
BosConf(config.getString(BosFileBaseOptions.BUCKET.key()));
+ HashMap<String, String> bosOptions = new HashMap<>();
+ bosOptions.put(ACCESS_KEY,
config.getString(BosFileBaseOptions.ACCESS_KEY.key()));
+ bosOptions.put(SECRET_KEY,
config.getString(BosFileBaseOptions.SECRET_KEY.key()));
+ bosOptions.put(ENDPOINT,
config.getString(BosFileBaseOptions.ENDPOINT.key()));
+ bosOptions.put(BUCKET_HIERARCHY, "false");
+ hadoopConf.setExtraOptions(bosOptions);
+ return hadoopConf;
+ }
+
+ public static HadoopConf buildWithReadonlyConfig(ReadonlyConfig
readonlyConfig) {
+ HadoopConf hadoopConf = new
BosConf(readonlyConfig.get(BosFileBaseOptions.BUCKET));
+ HashMap<String, String> bosOptions = new HashMap<>();
+ bosOptions.put(ACCESS_KEY,
readonlyConfig.get(BosFileBaseOptions.ACCESS_KEY));
+ bosOptions.put(SECRET_KEY,
readonlyConfig.get(BosFileBaseOptions.SECRET_KEY));
+ bosOptions.put(ENDPOINT,
readonlyConfig.get(BosFileBaseOptions.ENDPOINT));
+ bosOptions.put(BUCKET_HIERARCHY, "false");
+ hadoopConf.setExtraOptions(bosOptions);
+ return hadoopConf;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileBaseOptions.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileBaseOptions.java
new file mode 100644
index 0000000000..ebf4198aa5
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileBaseOptions.java
@@ -0,0 +1,39 @@
+/*
+ * 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.bos.config;
+
+import org.apache.seatunnel.api.configuration.Option;
+import org.apache.seatunnel.api.configuration.Options;
+import
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions;
+
+public class BosFileBaseOptions extends FileBaseSourceOptions {
+ public static final Option<String> ACCESS_KEY =
+ Options.key("access_key")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("BOS bucket access key");
+ public static final Option<String> SECRET_KEY =
+ Options.key("secret_key")
+ .stringType()
+ .noDefaultValue()
+ .withDescription("BOS bucket secret key");
+ public static final Option<String> ENDPOINT =
+
Options.key("endpoint").stringType().noDefaultValue().withDescription("BOS
endpoint");
+ public static final Option<String> BUCKET =
+
Options.key("bucket").stringType().noDefaultValue().withDescription("BOS
bucket");
+}
diff --git
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileSinkOptions.java
similarity index 85%
copy from
seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
copy to
seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileSinkOptions.java
index b8b195a87e..946091c6e3 100644
---
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileSinkOptions.java
@@ -15,12 +15,6 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.hive.storage;
+package org.apache.seatunnel.connectors.seatunnel.file.bos.config;
-public enum StorageType {
- S3,
- OSS,
- COS,
- FILE,
- HDFS
-}
+public class BosFileSinkOptions extends BosFileBaseOptions {}
diff --git
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileSourceOptions.java
similarity index 85%
copy from
seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
copy to
seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileSourceOptions.java
index b8b195a87e..2a5751abf4 100644
---
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/config/BosFileSourceOptions.java
@@ -15,12 +15,6 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.hive.storage;
+package org.apache.seatunnel.connectors.seatunnel.file.bos.config;
-public enum StorageType {
- S3,
- OSS,
- COS,
- FILE,
- HDFS
-}
+public class BosFileSourceOptions extends BosFileBaseOptions {}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/sink/BosFileSink.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/sink/BosFileSink.java
new file mode 100644
index 0000000000..4799ebd1ee
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/sink/BosFileSink.java
@@ -0,0 +1,42 @@
+/*
+ * 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.bos.sink;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosConf;
+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.sink.BaseFileSink;
+
+public class BosFileSink extends BaseFileSink {
+
+ public BosFileSink(ReadonlyConfig pluginConfig, CatalogTable catalogTable)
{
+ super(pluginConfig, catalogTable);
+ }
+
+ @Override
+ protected HadoopConf initHadoopConf() {
+ return BosConf.buildWithReadonlyConfig(pluginConfig);
+ }
+
+ @Override
+ public String getPluginName() {
+ return FileSystemType.BOS.getFileSystemPluginName();
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/sink/BosFileSinkFactory.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/sink/BosFileSinkFactory.java
new file mode 100644
index 0000000000..9de750a0a8
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/sink/BosFileSinkFactory.java
@@ -0,0 +1,120 @@
+/*
+ * 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.bos.sink;
+
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.table.connector.TableSink;
+import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.TableSinkFactory;
+import org.apache.seatunnel.api.table.factory.TableSinkFactoryContext;
+import
org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosFileSinkOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseOptions;
+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 com.google.auto.service.AutoService;
+
+import java.util.Arrays;
+
+@AutoService(Factory.class)
+public class BosFileSinkFactory implements TableSinkFactory {
+ @Override
+ public String factoryIdentifier() {
+ return FileSystemType.BOS.getFileSystemPluginName();
+ }
+
+ @Override
+ public OptionRule optionRule() {
+ return OptionRule.builder()
+ .required(FileBaseOptions.FILE_PATH)
+ .required(BosFileSinkOptions.BUCKET)
+ .required(BosFileSinkOptions.ACCESS_KEY)
+ .required(BosFileSinkOptions.SECRET_KEY)
+ .required(BosFileSinkOptions.ENDPOINT)
+ .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(FileBaseSinkOptions.FILENAME_EXTENSION)
+ .optional(FileBaseSinkOptions.TMP_PATH)
+ .build();
+ }
+
+ @Override
+ public TableSink createSink(TableSinkFactoryContext context) {
+ return () -> new BosFileSink(context.getOptions(),
context.getCatalogTable());
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/FileSystemType.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/source/BosFileSource.java
similarity index 50%
copy from
seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/FileSystemType.java
copy to
seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/source/BosFileSource.java
index c7a5c26aec..647292d5c5 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/config/FileSystemType.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/source/BosFileSource.java
@@ -15,28 +15,27 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.file.config;
+package org.apache.seatunnel.connectors.seatunnel.file.bos.source;
-import java.io.Serializable;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosConf;
+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.source.BaseFileSource;
-public enum FileSystemType implements Serializable {
- HDFS("HdfsFile"),
- LOCAL("LocalFile"),
- OSS("OssFile"),
- OSS_JINDO("OssJindoFile"),
- COS("CosFile"),
- FTP("FtpFile"),
- SFTP("SftpFile"),
- S3("S3File"),
- OBS("ObsFile");
+public class BosFileSource extends BaseFileSource {
- private final String fileSystemPluginName;
+ public BosFileSource(ReadonlyConfig pluginConfig) {
+ super(pluginConfig);
+ }
- FileSystemType(String fileSystemPluginName) {
- this.fileSystemPluginName = fileSystemPluginName;
+ @Override
+ protected HadoopConf initHadoopConf() {
+ return BosConf.buildWithReadonlyConfig(pluginConfig);
}
- public String getFileSystemPluginName() {
- return fileSystemPluginName;
+ @Override
+ public String getPluginName() {
+ return FileSystemType.BOS.getFileSystemPluginName();
}
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/source/BosFileSourceFactory.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/source/BosFileSourceFactory.java
new file mode 100644
index 0000000000..9da0434b51
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/bos/source/BosFileSourceFactory.java
@@ -0,0 +1,122 @@
+/*
+ * 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.bos.source;
+
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+import org.apache.seatunnel.api.source.SeaTunnelSource;
+import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.connector.TableSource;
+import org.apache.seatunnel.api.table.factory.Factory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactory;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
+import
org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosFileSourceOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseOptions;
+import
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileFormat;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileSystemType;
+
+import com.google.auto.service.AutoService;
+
+import java.io.Serializable;
+import java.util.Arrays;
+
+@AutoService(Factory.class)
+public class BosFileSourceFactory implements TableSourceFactory {
+ @Override
+ public String factoryIdentifier() {
+ return FileSystemType.BOS.getFileSystemPluginName();
+ }
+
+ @Override
+ public OptionRule optionRule() {
+ return OptionRule.builder()
+ .required(FileBaseOptions.FILE_PATH)
+ .required(BosFileSourceOptions.BUCKET)
+ .required(BosFileSourceOptions.ACCESS_KEY)
+ .required(BosFileSourceOptions.SECRET_KEY)
+ .required(BosFileSourceOptions.ENDPOINT)
+ .required(FileBaseSourceOptions.FILE_FORMAT_TYPE)
+ .conditional(
+ FileBaseSourceOptions.FILE_FORMAT_TYPE,
+ FileFormat.TEXT,
+ FileBaseSourceOptions.ROW_DELIMITER,
+ FileBaseSourceOptions.FIELD_DELIMITER,
+ FileBaseSourceOptions.SKIP_HEADER_ROW_NUMBER)
+ .conditional(
+ FileBaseSourceOptions.FILE_FORMAT_TYPE,
+ FileFormat.XML,
+ FileBaseSourceOptions.XML_ROW_TAG,
+ FileBaseSourceOptions.XML_USE_ATTR_FORMAT)
+ .conditional(
+ FileBaseSourceOptions.FILE_FORMAT_TYPE,
+ FileFormat.CSV,
+ FileBaseSourceOptions.SKIP_HEADER_ROW_NUMBER)
+ .conditional(
+ FileBaseSourceOptions.FILE_FORMAT_TYPE,
+ Arrays.asList(
+ FileFormat.TEXT,
+ FileFormat.JSON,
+ FileFormat.EXCEL,
+ FileFormat.CSV,
+ FileFormat.XML),
+ ConnectorCommonOptions.SCHEMA)
+ .conditional(
+ FileBaseSourceOptions.FILE_FORMAT_TYPE,
+ Arrays.asList(
+ FileFormat.TEXT, FileFormat.JSON,
FileFormat.CSV, FileFormat.XML),
+ FileBaseSourceOptions.ENCODING)
+ .optional(FileBaseSourceOptions.PARSE_PARTITION_FROM_PATH)
+ .optional(FileBaseSourceOptions.DATE_FORMAT_LEGACY)
+ .optional(FileBaseSourceOptions.DATETIME_FORMAT_LEGACY)
+ .optional(FileBaseSourceOptions.TIME_FORMAT_LEGACY)
+ .optional(FileBaseSourceOptions.FILE_FILTER_PATTERN)
+ .optional(FileBaseSourceOptions.COMPRESS_CODEC)
+ .optional(FileBaseSourceOptions.ARCHIVE_COMPRESS_CODEC)
+ .optional(FileBaseSourceOptions.NULL_FORMAT)
+ .optional(FileBaseSourceOptions.FILENAME_EXTENSION)
+ .optional(FileBaseSourceOptions.READ_COLUMNS)
+ .optional(
+ FileBaseSourceOptions.SHEET_NAME,
+ FileBaseSourceOptions.EXCEL_ENGINE,
+ FileBaseSourceOptions.POI_EXCEL_MAX_FILE_SIZE)
+ .conditional(
+ FileBaseSourceOptions.FILE_FORMAT_TYPE,
+ FileFormat.MARKDOWN,
+ FileBaseSourceOptions.MARKDOWN_RAG_METADATA_ENABLED)
+ .conditional(
+ FileBaseSourceOptions.FILE_FORMAT_TYPE,
+ FileFormat.PDF,
+ FileBaseSourceOptions.PDF_RAG_METADATA_ENABLED)
+ .optional(FileBaseSourceOptions.QUOTE_CHAR)
+ .optional(FileBaseSourceOptions.ESCAPE_CHAR)
+ .optional(FileBaseSourceOptions.RECURSIVE_FILE_SCAN)
+ .build();
+ }
+
+ @Override
+ public Class<? extends SeaTunnelSource> getSourceClass() {
+ return BosFileSource.class;
+ }
+
+ @Override
+ public <T, SplitT extends SourceSplit, StateT extends Serializable>
+ TableSource<T, SplitT, StateT>
createSource(TableSourceFactoryContext context) {
+ return () -> (SeaTunnelSource<T, SplitT, StateT>) new
BosFileSource(context.getOptions());
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/resources/META-INF/services/org.apache.hadoop.fs.FileSystem
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/resources/META-INF/services/org.apache.hadoop.fs.FileSystem
new file mode 100644
index 0000000000..b3cd556dfb
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/main/resources/META-INF/services/org.apache.hadoop.fs.FileSystem
@@ -0,0 +1,16 @@
+# 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.
+
+org.apache.hadoop.fs.bos.BaiduBosFileSystem
diff --git
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/bos/BosFileFactoryTest.java
similarity index 60%
copy from
seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
copy to
seatunnel-connectors-v2/connector-file/connector-file-bos/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/bos/BosFileFactoryTest.java
index b8b195a87e..9071d46c97 100644
---
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-bos/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/bos/BosFileFactoryTest.java
@@ -15,12 +15,19 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.hive.storage;
+package org.apache.seatunnel.connectors.seatunnel.file.bos;
-public enum StorageType {
- S3,
- OSS,
- COS,
- FILE,
- HDFS
+import
org.apache.seatunnel.connectors.seatunnel.file.bos.sink.BosFileSinkFactory;
+import
org.apache.seatunnel.connectors.seatunnel.file.bos.source.BosFileSourceFactory;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+public class BosFileFactoryTest {
+
+ @Test
+ void optionRule() {
+ Assertions.assertNotNull((new BosFileSourceFactory()).optionRule());
+ Assertions.assertNotNull((new BosFileSinkFactory()).optionRule());
+ }
}
diff --git a/seatunnel-connectors-v2/connector-file/pom.xml
b/seatunnel-connectors-v2/connector-file/pom.xml
index efb32ab444..44ee201599 100644
--- a/seatunnel-connectors-v2/connector-file/pom.xml
+++ b/seatunnel-connectors-v2/connector-file/pom.xml
@@ -41,6 +41,7 @@
<module>connector-file-obs</module>
<module>connector-file-jindo-oss</module>
<module>connector-file-cos</module>
+ <module>connector-file-bos</module>
</modules>
<properties>
diff --git a/seatunnel-connectors-v2/connector-hive/pom.xml
b/seatunnel-connectors-v2/connector-hive/pom.xml
index 877ea3dd4d..b333c91573 100644
--- a/seatunnel-connectors-v2/connector-hive/pom.xml
+++ b/seatunnel-connectors-v2/connector-hive/pom.xml
@@ -65,6 +65,11 @@
<artifactId>connector-file-cos</artifactId>
<version>${project.version}</version>
</dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-file-bos</artifactId>
+ <version>${project.version}</version>
+ </dependency>
<dependency>
<groupId>org.apache.seatunnel</groupId>
<artifactId>seatunnel-shade-hadoop3-uber</artifactId>
diff --git
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/BOSStorage.java
b/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/BOSStorage.java
new file mode 100644
index 0000000000..5da041cc4e
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/BOSStorage.java
@@ -0,0 +1,89 @@
+/*
+ * 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.hive.storage;
+
+import org.apache.seatunnel.shade.com.typesafe.config.Config;
+import org.apache.seatunnel.shade.com.typesafe.config.ConfigValueFactory;
+import org.apache.seatunnel.shade.org.apache.commons.lang3.StringUtils;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosConf;
+import
org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosFileBaseOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
+
+import org.apache.hadoop.conf.Configuration;
+
+import java.util.Map;
+
+/**
+ * Builds Hadoop configuration for Hive tables stored on BOS ({@code bos://}).
+ *
+ * <p>Loads bucket and BOS credentials from {@code hive-site.xml}/{@code
core-site.xml} or {@code
+ * hive.hadoop.conf}, then delegates to {@link BosConf} for {@link HadoopConf}
used by Hive
+ * source/sink file access.
+ */
+public class BOSStorage extends AbstractStorage {
+
+ private static final String FS_BOS_ACCESS_KEY = "fs.bos.access.key";
+ private static final String FS_BOS_SECRET_KEY = "fs.bos.secret.access.key";
+ private static final String FS_BOS_ENDPOINT = "fs.bos.endpoint";
+
+ @Override
+ public HadoopConf buildHadoopConfWithReadOnlyConfig(ReadonlyConfig
readonlyConfig) {
+ Configuration configuration = loadHiveBaseHadoopConfig(readonlyConfig);
+ Config config = fillBucket(readonlyConfig, configuration);
+ config =
+ config.withValue(
+ BosFileBaseOptions.ACCESS_KEY.key(),
+ ConfigValueFactory.fromAnyRef(
+ getConfigurationValue(
+ configuration,
+ BosFileBaseOptions.ACCESS_KEY.key(),
+ FS_BOS_ACCESS_KEY)));
+ config =
+ config.withValue(
+ BosFileBaseOptions.SECRET_KEY.key(),
+ ConfigValueFactory.fromAnyRef(
+ getConfigurationValue(
+ configuration,
+ BosFileBaseOptions.SECRET_KEY.key(),
+ FS_BOS_SECRET_KEY)));
+ config =
+ config.withValue(
+ BosFileBaseOptions.ENDPOINT.key(),
+ ConfigValueFactory.fromAnyRef(
+ getConfigurationValue(
+ configuration,
+ BosFileBaseOptions.ENDPOINT.key(),
+ FS_BOS_ENDPOINT)));
+ HadoopConf hadoopConf = BosConf.buildWithConfig(config);
+ Map<String, String> propsInConfiguration =
+ configuration.getPropsWithPrefix(StringUtils.EMPTY);
+ hadoopConf.setExtraOptions(propsInConfiguration);
+ return hadoopConf;
+ }
+
+ private static String getConfigurationValue(
+ Configuration configuration, String optionKey, String hadoopKey) {
+ String value = configuration.get(optionKey);
+ if (StringUtils.isBlank(value)) {
+ value = configuration.get(hadoopKey);
+ }
+ return value;
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactory.java
b/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactory.java
index 926e7d70b9..6a7cb9dd4c 100644
---
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactory.java
+++
b/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactory.java
@@ -26,6 +26,8 @@ public class StorageFactory {
return new OSSStorage();
} else if
(hiveSdLocation.startsWith(StorageType.COS.name().toLowerCase())) {
return new COSStorage();
+ } else if
(hiveSdLocation.startsWith(StorageType.BOS.name().toLowerCase())) {
+ return new BOSStorage();
} else if
(hiveSdLocation.startsWith(StorageType.FILE.name().toLowerCase())) {
// Currently used in e2e, When Hive uses local files as storage,
"file:" needs to be
// replaced with "file:/" to avoid being recognized as HDFS
storage.
diff --git
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
b/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
index b8b195a87e..6b241d75d9 100644
---
a/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
+++
b/seatunnel-connectors-v2/connector-hive/src/main/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageType.java
@@ -21,6 +21,7 @@ public enum StorageType {
S3,
OSS,
COS,
+ BOS,
FILE,
HDFS
}
diff --git
a/seatunnel-connectors-v2/connector-hive/src/test/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/BosStorageTest.java
b/seatunnel-connectors-v2/connector-hive/src/test/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/BosStorageTest.java
new file mode 100644
index 0000000000..7d82b22c72
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-hive/src/test/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/BosStorageTest.java
@@ -0,0 +1,80 @@
+/*
+ * 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.hive.storage;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosConf;
+import
org.apache.seatunnel.connectors.seatunnel.file.bos.config.BosFileBaseOptions;
+import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Test;
+
+import java.net.URISyntaxException;
+import java.net.URL;
+import java.nio.file.Paths;
+import java.util.HashMap;
+
+public class BosStorageTest {
+
+ private static final ReadonlyConfig BOS =
+ ReadonlyConfig.fromMap(
+ new HashMap<String, Object>() {
+ {
+ put(
+ "hive.hadoop.conf",
+ new HashMap<String, String>() {
+ {
+ put("bucket", "bos://my_bucket");
+
put(BosFileBaseOptions.ACCESS_KEY.key(), "test");
+
put(BosFileBaseOptions.SECRET_KEY.key(), "test");
+ put(
+
BosFileBaseOptions.ENDPOINT.key(),
+ "http://bj.bcebos.com");
+ }
+ });
+ }
+ });
+
+ @Test
+ void fillBucketInHadoopConf() {
+ BOSStorage bosStorage = new BOSStorage();
+ HadoopConf bosConf = bosStorage.buildHadoopConfWithReadOnlyConfig(BOS);
+ assertHadoopConf(bosConf);
+ }
+
+ @Test
+ void fillBucketInHadoopConfPath() throws URISyntaxException {
+ URL resource = BosStorageTest.class.getResource("/bos");
+ String filePath = Paths.get(resource.toURI()).toString();
+ HashMap<String, Object> map = new HashMap<>();
+ map.put("hive.hadoop.conf-path", filePath);
+ map.putAll(BOS.toMap());
+ ReadonlyConfig readonlyConfig = ReadonlyConfig.fromMap(map);
+ BOSStorage bosStorage = new BOSStorage();
+ HadoopConf hadoopConf =
bosStorage.buildHadoopConfWithReadOnlyConfig(readonlyConfig);
+ assertHadoopConf(hadoopConf);
+ }
+
+ private static void assertHadoopConf(HadoopConf bosConf) {
+ Assertions.assertTrue(bosConf instanceof BosConf);
+ Assertions.assertEquals(bosConf.getSchema(), "bos");
+ Assertions.assertEquals(
+ bosConf.getFsHdfsImpl(),
"org.apache.hadoop.fs.bos.BaiduBosFileSystem");
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-hive/src/test/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactoryTest.java
b/seatunnel-connectors-v2/connector-hive/src/test/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactoryTest.java
index cd4d99bfaa..9d225827bf 100644
---
a/seatunnel-connectors-v2/connector-hive/src/test/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactoryTest.java
+++
b/seatunnel-connectors-v2/connector-hive/src/test/java/org/apache/seatunnel/connectors/seatunnel/hive/storage/StorageFactoryTest.java
@@ -34,6 +34,7 @@ public class StorageFactoryTest {
put("s3a://path/to/", S3Storage.class);
put("oss://path/to/", OSSStorage.class);
put("cosn://path/to/", COSStorage.class);
+ put("bos://path/to/", BOSStorage.class);
}
};
diff --git
a/seatunnel-connectors-v2/connector-hive/src/test/resources/bos/core-site.xml
b/seatunnel-connectors-v2/connector-hive/src/test/resources/bos/core-site.xml
new file mode 100644
index 0000000000..4111e50fc9
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-hive/src/test/resources/bos/core-site.xml
@@ -0,0 +1,37 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+ 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.
+-->
+<configuration>
+ <property>
+ <name>fs.defaultFS</name>
+ <value>bos://mybucket</value>
+ </property>
+ <property>
+ <name>fs.bos.impl</name>
+ <value>org.apache.hadoop.fs.bos.BaiduBosFileSystem</value>
+ </property>
+ <property>
+ <name>access_key</name>
+ <value>your-bos-access-key</value>
+ </property>
+ <property>
+ <name>secret_key</name>
+ <value>your-bos-secret-key</value>
+ </property>
+ <property>
+ <name>endpoint</name>
+ <value>http://bj.bcebos.com</value>
+ </property>
+</configuration>
diff --git a/seatunnel-dist/pom.xml b/seatunnel-dist/pom.xml
index 0821306ef4..d5e16abebe 100644
--- a/seatunnel-dist/pom.xml
+++ b/seatunnel-dist/pom.xml
@@ -767,6 +767,13 @@
<scope>provided</scope>
</dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-file-bos</artifactId>
+ <version>${project.version}</version>
+ <scope>provided</scope>
+ </dependency>
+
<dependency>
<groupId>org.apache.seatunnel</groupId>
<artifactId>connector-paimon</artifactId>
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/pom.xml
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/pom.xml
new file mode 100644
index 0000000000..5d97640b3e
--- /dev/null
+++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/pom.xml
@@ -0,0 +1,48 @@
+<?xml version="1.0" encoding="UTF-8"?>
+<!--
+ 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.
+-->
+<project xmlns="http://maven.apache.org/POM/4.0.0"
xmlns:xsi="http://www.w3.org/2001/XMLSchema-instance"
+ xsi:schemaLocation="http://maven.apache.org/POM/4.0.0
http://maven.apache.org/xsd/maven-4.0.0.xsd">
+ <modelVersion>4.0.0</modelVersion>
+ <parent>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>seatunnel-connector-v2-e2e</artifactId>
+ <version>${revision}</version>
+ </parent>
+
+ <artifactId>connector-file-bos-e2e</artifactId>
+ <name>SeaTunnel : E2E : Connector V2 : File Bos</name>
+
+ <dependencies>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-fake</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-file-bos</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
+ <dependency>
+ <groupId>org.apache.seatunnel</groupId>
+ <artifactId>connector-assert</artifactId>
+ <version>${project.version}</version>
+ <scope>test</scope>
+ </dependency>
+ </dependencies>
+</project>
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/bos/BosFileIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/bos/BosFileIT.java
new file mode 100644
index 0000000000..51a772aa73
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/bos/BosFileIT.java
@@ -0,0 +1,71 @@
+/*
+ * 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.e2e.connector.file.bos;
+
+import org.apache.seatunnel.e2e.common.TestSuiteBase;
+import org.apache.seatunnel.e2e.common.container.TestContainer;
+
+import org.junit.jupiter.api.Assertions;
+import org.junit.jupiter.api.Disabled;
+import org.junit.jupiter.api.TestTemplate;
+import org.testcontainers.containers.Container;
+
+import java.io.IOException;
+
+@Disabled
+public class BosFileIT extends TestSuiteBase {
+
+ @TestTemplate
+ public void testBosFileWriteAndRead(TestContainer container)
+ throws IOException, InterruptedException {
+ Container.ExecResult excelWriteResult =
+ container.executeJob("/excel/fake_to_bos_excel.conf");
+ Assertions.assertEquals(0, excelWriteResult.getExitCode(),
excelWriteResult.getStderr());
+ Container.ExecResult excelReadResult =
+ container.executeJob("/excel/bos_excel_to_assert.conf");
+ Assertions.assertEquals(0, excelReadResult.getExitCode(),
excelReadResult.getStderr());
+
+ Container.ExecResult textWriteResult =
+ container.executeJob("/text/fake_to_bos_file_text.conf");
+ Assertions.assertEquals(0, textWriteResult.getExitCode());
+ Container.ExecResult textReadResult =
+ container.executeJob("/text/bos_file_text_to_assert.conf");
+ Assertions.assertEquals(0, textReadResult.getExitCode());
+
+ Container.ExecResult jsonWriteResult =
+ container.executeJob("/json/fake_to_bos_file_json.conf");
+ Assertions.assertEquals(0, jsonWriteResult.getExitCode());
+ Container.ExecResult jsonReadResult =
+ container.executeJob("/json/bos_file_json_to_assert.conf");
+ Assertions.assertEquals(0, jsonReadResult.getExitCode());
+
+ Container.ExecResult orcWriteResult =
+ container.executeJob("/orc/fake_to_bos_file_orc.conf");
+ Assertions.assertEquals(0, orcWriteResult.getExitCode());
+ Container.ExecResult orcReadResult =
+ container.executeJob("/orc/bos_file_orc_to_assert.conf");
+ Assertions.assertEquals(0, orcReadResult.getExitCode());
+
+ Container.ExecResult parquetWriteResult =
+ container.executeJob("/parquet/fake_to_bos_file_parquet.conf");
+ Assertions.assertEquals(0, parquetWriteResult.getExitCode());
+ Container.ExecResult parquetReadResult =
+
container.executeJob("/parquet/bos_file_parquet_to_assert.conf");
+ Assertions.assertEquals(0, parquetReadResult.getExitCode());
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/excel/bos_excel_to_assert.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/excel/bos_excel_to_assert.conf
new file mode 100644
index 0000000000..49d2bb292d
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/excel/bos_excel_to_assert.conf
@@ -0,0 +1,118 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ BosFile {
+ path = "/read/excel"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ plugin_output = "fake"
+ file_format_type = excel
+ field_delimiter = ;
+ skip_header_row_number = 1
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ }
+ }
+ }
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MAX_ROW
+ rule_value = 5
+ }
+ ],
+ field_rules = [
+ {
+ field_name = c_string
+ field_type = string
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_boolean
+ field_type = boolean
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_double
+ field_type = double
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ }
+ ]
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/excel/fake_to_bos_excel.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/excel/fake_to_bos_excel.conf
new file mode 100644
index 0000000000..58d54922b0
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/excel/fake_to_bos_excel.conf
@@ -0,0 +1,84 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ FakeSource {
+ plugin_output = "fake"
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ }
+ }
+ }
+ }
+}
+
+sink {
+ BosFile {
+ path="/sink/execl"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ partition_dir_expression = "${k0}=${v0}"
+ is_partition_field_write_in_file = true
+ file_name_expression = "${transactionId}_${now}"
+ file_format_type = "excel"
+ filename_time_format = "yyyy.MM.dd"
+ is_enable_transaction = true
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/json/bos_file_json_to_assert.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/json/bos_file_json_to_assert.conf
new file mode 100644
index 0000000000..23cfab4e94
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/json/bos_file_json_to_assert.conf
@@ -0,0 +1,116 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ BosFile {
+ path = "/read/json"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ file_format_type = "json"
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ C_MAP = "map<string, string>"
+ C_ARRAY = "array<int>"
+ C_STRING = string
+ C_BOOLEAN = boolean
+ C_TINYINT = tinyint
+ C_SMALLINT = smallint
+ C_INT = int
+ C_BIGINT = bigint
+ C_FLOAT = float
+ C_DOUBLE = double
+ C_BYTES = bytes
+ C_DATE = date
+ C_DECIMAL = "decimal(38, 18)"
+ C_TIMESTAMP = timestamp
+ }
+ }
+ }
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MAX_ROW
+ rule_value = 5
+ }
+ ],
+ field_rules = [
+ {
+ field_name = c_string
+ field_type = string
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_boolean
+ field_type = boolean
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_double
+ field_type = double
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ }
+ ]
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/json/fake_to_bos_file_json.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/json/fake_to_bos_file_json.conf
new file mode 100644
index 0000000000..9a8c3e14b5
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/json/fake_to_bos_file_json.conf
@@ -0,0 +1,85 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ FakeSource {
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ }
+ }
+ }
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ BosFile {
+ path="/sink/json"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ row_delimiter = "\n"
+ partition_dir_expression = "${k0}=${v0}"
+ is_partition_field_write_in_file = true
+ file_name_expression = "${transactionId}_${now}"
+ file_format_type = "json"
+ filename_time_format = "yyyy.MM.dd"
+ is_enable_transaction = true
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/orc/bos_file_orc_to_assert.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/orc/bos_file_orc_to_assert.conf
new file mode 100644
index 0000000000..e3ae38cfbe
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/orc/bos_file_orc_to_assert.conf
@@ -0,0 +1,82 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ BosFile {
+ path = "/read/orc"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ file_format_type = "orc"
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MAX_ROW
+ rule_value = 5
+ }
+ ],
+ field_rules = [
+ {
+ field_name = c_string
+ field_type = string
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_boolean
+ field_type = boolean
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_double
+ field_type = double
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ }
+ ]
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/orc/fake_to_bos_file_orc.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/orc/fake_to_bos_file_orc.conf
new file mode 100644
index 0000000000..4d934280ee
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/orc/fake_to_bos_file_orc.conf
@@ -0,0 +1,86 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ FakeSource {
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ }
+ }
+ }
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ BosFile {
+ path="/sink/orc"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ row_delimiter = "\n"
+ partition_dir_expression = "${k0}=${v0}"
+ is_partition_field_write_in_file = true
+ file_name_expression = "${transactionId}_${now}"
+ file_format_type = "orc"
+ filename_time_format = "yyyy.MM.dd"
+ is_enable_transaction = true
+ compress_codec = "zlib"
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/parquet/bos_file_parquet_to_assert.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/parquet/bos_file_parquet_to_assert.conf
new file mode 100644
index 0000000000..8b962f5f1a
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/parquet/bos_file_parquet_to_assert.conf
@@ -0,0 +1,82 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ BosFile {
+ path = "/read/parquet"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ file_format_type = "parquet"
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MAX_ROW
+ rule_value = 5
+ }
+ ],
+ field_rules = [
+ {
+ field_name = c_string
+ field_type = string
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_boolean
+ field_type = boolean
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_double
+ field_type = double
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ }
+ ]
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/parquet/fake_to_bos_file_parquet.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/parquet/fake_to_bos_file_parquet.conf
new file mode 100644
index 0000000000..018277680b
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/parquet/fake_to_bos_file_parquet.conf
@@ -0,0 +1,86 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ FakeSource {
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ }
+ }
+ }
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ BosFile {
+ path="/sink/parquet"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ row_delimiter = "\n"
+ partition_dir_expression = "${k0}=${v0}"
+ is_partition_field_write_in_file = true
+ file_name_expression = "${transactionId}_${now}"
+ file_format_type = "parquet"
+ filename_time_format = "yyyy.MM.dd"
+ is_enable_transaction = true
+ compress_codec = "gzip"
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/text/bos_file_text_to_assert.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/text/bos_file_text_to_assert.conf
new file mode 100644
index 0000000000..0cd8513e72
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/text/bos_file_text_to_assert.conf
@@ -0,0 +1,116 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ BosFile {
+ path = "/read/text"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ file_format_type = "text"
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ }
+ }
+ }
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ Assert {
+ rules {
+ row_rules = [
+ {
+ rule_type = MAX_ROW
+ rule_value = 5
+ }
+ ],
+ field_rules = [
+ {
+ field_name = c_string
+ field_type = string
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_boolean
+ field_type = boolean
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ },
+ {
+ field_name = c_double
+ field_type = double
+ field_value = [
+ {
+ rule_type = NOT_NULL
+ }
+ ]
+ }
+ ]
+ }
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/text/fake_to_bos_file_text.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/text/fake_to_bos_file_text.conf
new file mode 100644
index 0000000000..546f6b7dee
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-bos-e2e/src/test/resources/text/fake_to_bos_file_text.conf
@@ -0,0 +1,86 @@
+#
+# 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"
+
+ # You can set spark configuration here
+ spark.app.name = "SeaTunnel"
+ spark.executor.instances = 2
+ spark.executor.cores = 1
+ spark.executor.memory = "1g"
+ spark.master = local
+}
+
+source {
+ FakeSource {
+ schema = {
+ fields {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ c_row = {
+ c_map = "map<string, string>"
+ c_array = "array<int>"
+ c_string = string
+ c_boolean = boolean
+ c_tinyint = tinyint
+ c_smallint = smallint
+ c_int = int
+ c_bigint = bigint
+ c_float = float
+ c_double = double
+ c_bytes = bytes
+ c_date = date
+ c_decimal = "decimal(38, 18)"
+ c_timestamp = timestamp
+ }
+ }
+ }
+ plugin_output = "fake"
+ }
+}
+
+sink {
+ BosFile {
+ path="/sink/text"
+ bucket = "bos://seatunnel-test"
+ access_key = "dummy"
+ secret_key = "dummy"
+ endpoint = "http://bj.bcebos.com"
+ row_delimiter = "\n"
+ partition_dir_expression = "${k0}=${v0}"
+ is_partition_field_write_in_file = true
+ file_name_expression = "${transactionId}_${now}"
+ file_format_type = "text"
+ filename_time_format = "yyyy.MM.dd"
+ is_enable_transaction = true
+ compress_codec = "lzo"
+ }
+}
diff --git a/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
b/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
index cbbb7012a2..ab91ca5eba 100644
--- a/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
+++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/pom.xml
@@ -40,6 +40,7 @@
<module>connector-amazonsqs-e2e</module>
<module>connector-file-local-e2e</module>
<module>connector-file-cos-e2e</module>
+ <module>connector-file-bos-e2e</module>
<module>connector-file-hadoop-e2e</module>
<module>connector-file-sftp-e2e</module>
<module>connector-file-oss-e2e</module>