This is an automated email from the ASF dual-hosted git repository.
zhangshenghang 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 ec92180c8b [Feature][Connector-V2] Enable continuous discovery for S3
and OSS file sources (#11789)
ec92180c8b is described below
commit ec92180c8b476d13b3fe07c0b2913430e053906b
Author: Goutam Adwant <[email protected]>
AuthorDate: Fri Aug 14 05:44:03 2026 -0700
[Feature][Connector-V2] Enable continuous discovery for S3 and OSS file
sources (#11789)
Signed-off-by: goutamadwant <[email protected]>
---
docs/en/connectors/source/OssFile.md | 58 ++++++++++++-
docs/en/connectors/source/S3File.md | 60 ++++++++++++-
docs/zh/connectors/source/OssFile.md | 58 ++++++++++++-
docs/zh/connectors/source/S3File.md | 60 ++++++++++++-
.../file/oss/source/OssFileSourceFactory.java | 26 ++++++
.../file/oss/OssFileSourceFactoryTest.java | 91 ++++++++++++++++++++
.../file/s3/source/S3FileSourceFactory.java | 26 ++++++
.../seatunnel/file/s3/S3FileSourceFactoryTest.java | 99 ++++++++++++++++++++++
.../e2e/connector/file/s3/S3FileWithFilterIT.java | 74 ++++++++++++++++
.../seatunnel/e2e/connector/file/s3/S3Utils.java | 41 +++++++++
.../s3_file_binary_update_distcp_continuous.conf | 65 ++++++++++++++
11 files changed, 654 insertions(+), 4 deletions(-)
diff --git a/docs/en/connectors/source/OssFile.md
b/docs/en/connectors/source/OssFile.md
index 08275a88a7..cab894e741 100644
--- a/docs/en/connectors/source/OssFile.md
+++ b/docs/en/connectors/source/OssFile.md
@@ -24,7 +24,7 @@ import ChangeLog from '../changelog/connector-file-oss.md';
## Key features
- [x] [batch](../../introduction/concepts/connector-v2-features.md)
-- [ ] [stream](../../introduction/concepts/connector-v2-features.md)
+- [x] [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.
@@ -213,6 +213,20 @@ If you assign file type to `parquet` `orc`, schema option
not required, connecto
| null_format | string | no | - | Only
used when file_format_type is text. null_format to define which strings can be
represented as null. e.g: `\N`
|
| binary_chunk_size | int | no | 1024 | Only
used when file_format_type is binary. The chunk size (in bytes) for reading
binary files. Default is 1024 bytes. Larger values may improve performance for
large files but use more memory.
|
| binary_complete_file_mode | boolean | no | false | Only
used when file_format_type is binary. Whether to read the complete file as a
single chunk instead of splitting into chunks. When enabled, the entire file
content will be read into memory at once. Default is false.
|
+| discovery_mode | string | no | once | File
discovery mode. Supported values: `once` (default), `continuous`. Continuous
mode periodically scans the path and currently requires `sync_mode=update` and
`file_format_type=binary`. |
+| scan_interval | string | no | 10S |
Polling interval used when `discovery_mode=continuous`. Shorthand values such
as `10S` and ISO-8601 values such as `PT10S` are supported. |
+| start_mode | string | no | earliest |
Initial scan behavior for continuous discovery. `earliest` processes existing
files; `latest` processes only later additions or changes. |
+| sync_mode | string | no | full | File
sync mode. `update` compares source objects with `target_path` and reads only
new or changed objects. Update mode currently supports binary format only. |
+| target_path | string | no | - |
Required when `sync_mode=update`. Target base path used for comparison and
normally the same as the sink `path`. |
+| target_hadoop_conf | map | no | - |
Optional Hadoop configuration for the target filesystem when
`sync_mode=update`. |
+| update_strategy | string | no | distcp |
Comparison strategy used by update mode. Supported values are `distcp` and
`strict`. |
+| compare_mode | string | no | len_mtime |
Comparison mode used by update mode. Supported values are `len_mtime` and
`checksum`; checksum requires `update_strategy=strict`. |
+| update_compare_parallelism | int | no | 8 |
Maximum parallelism for target metadata lookups. Valid range is `1-64`. |
+| update_compare_bulk_threshold | int | no | 0 |
Positive values switch comparison to a directory listing when that many
candidates share a target parent. `0` disables automatic bulk listing. |
+| post_sync_action | string | no | none |
Optional action after a continuously discovered object is checkpointed.
Supported values are `none`, `delete`, and `backup`. |
+| backup_path | string | no | - |
Required when `post_sync_action=backup`. Backup destination must not overlap
with the source `path`. |
+| retention_max_age | string | no | - |
Optional maximum age for SeaTunnel backup objects under `backup_path`. |
+| retention_check_interval | string | no | 1H |
Retention scan interval when backup retention is configured. |
| file_filter_pattern | string | no | |
Filter pattern, which used for filtering files.
|
| common-options | config | no | - |
Source plugin common parameters, please refer to [Source Common
Options](../common-options/source-common-options.md) for details.
|
| file_filter_modified_start | string | no | - | File
modification time filter. The connector will filter some files base on the last
modification start time (include start time). The default data format is
`yyyy-MM-dd HH:mm:ss`.
|
@@ -385,6 +399,48 @@ When specified, the connector will fetch table schema from
the external metadata
For more information, please refer to [Metadata
SPI](../../introduction/concepts/metadata-spi.md).
+## Continuous Discovery
+
+`discovery_mode=continuous` keeps a streaming job running and polls OSS for
new or changed objects. This mode uses the existing file comparison path; it
does not consume OSS event notifications and does not emit object delete events
or changelog rows.
+
+Continuous discovery currently requires `file_format_type="binary"` and
`sync_mode="update"`. Set `target_path` to the same base path used by the sink
so the source can skip unchanged objects. The default `discovery_mode="once"`
preserves the existing bounded-source behavior.
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+}
+
+source {
+ OssFile {
+ path = "/watch/source"
+ bucket = "oss://seatunnel-test"
+ endpoint = "oss-cn-hangzhou.aliyuncs.com"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ access_secret = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+
+ discovery_mode = "continuous"
+ scan_interval = "10S"
+ start_mode = "earliest"
+ sync_mode = "update"
+ target_path = "/watch/target"
+ }
+}
+
+sink {
+ OssFile {
+ path = "/watch/target"
+ tmp_path = "/watch/tmp"
+ bucket = "oss://seatunnel-test"
+ endpoint = "oss-cn-hangzhou.aliyuncs.com"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ access_secret = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+ }
+}
+```
+
## How to Create a Oss Data Synchronization Jobs
The following example demonstrates how to create a data synchronization job
that reads data from Oss and prints it on the local client:
diff --git a/docs/en/connectors/source/S3File.md
b/docs/en/connectors/source/S3File.md
index c9a7003382..9e2b0e3456 100644
--- a/docs/en/connectors/source/S3File.md
+++ b/docs/en/connectors/source/S3File.md
@@ -13,7 +13,7 @@ import ChangeLog from '../changelog/connector-file-s3.md';
## Key Features
- [x] [batch](../../introduction/concepts/connector-v2-features.md)
-- [ ] [stream](../../introduction/concepts/connector-v2-features.md)
+- [x] [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.
@@ -224,6 +224,20 @@ If you assign file type to `parquet` `orc`, schema option
not required, connecto
| null_format | string | no | -
| Only used when file_format_type is text.
null_format to define which strings can be represented as null. e.g: `\N`
[...]
| binary_chunk_size | int | no | 1024
| Only used when file_format_type is binary.
The chunk size (in bytes) for reading binary files. Default is 1024 bytes.
Larger values may improve performance for large files but use more memory.
[...]
| binary_complete_file_mode | boolean | no | false
| Only used when file_format_type is binary.
Whether to read the complete file as a single chunk instead of splitting into
chunks. When enabled, the entire file content will be read into memory at once.
Default is false.
[...]
+| discovery_mode | string | no | once
| File discovery mode. Supported values: `once`
(default), `continuous`. When `continuous`, the source periodically scans the
path and processes new or changed files as an unbounded source. Continuous mode
currently requires `sync_mode=update` and `file_format_type=binary`.
|
+| scan_interval | string | no | 10S
| Polling interval used when
`discovery_mode=continuous`. Shorthand values such as `10S` and ISO-8601 values
such as `PT10S` are supported.
|
+| start_mode | string | no | earliest
| Initial scan behavior for continuous
discovery. `earliest` processes existing files; `latest` ignores files present
when the job starts and processes later additions or changes.
|
+| sync_mode | string | no | full
| File sync mode. `update` compares source
objects with `target_path` and reads only new or changed objects. Update mode
currently supports binary format only.
|
+| target_path | string | no | -
| Required when `sync_mode=update`. Target base
path used to compare objects by relative path. It should normally match the
sink `path`.
|
+| target_hadoop_conf | map | no | -
| Optional Hadoop configuration for the target
filesystem when `sync_mode=update`.
|
+| update_strategy | string | no | distcp
| Comparison strategy used by update mode.
Supported values are `distcp` and `strict`.
|
+| compare_mode | string | no | len_mtime
| Comparison mode used by update mode.
Supported values are `len_mtime` and `checksum`; checksum is valid only with
`update_strategy=strict`.
|
+| update_compare_parallelism | int | no | 8
| Maximum parallelism for target metadata
lookups. Valid range is `1-64`.
|
+| update_compare_bulk_threshold | int | no | 0
| Positive values switch comparison to a
directory listing when that many candidates share a target parent. `0` disables
automatic bulk listing.
|
+| post_sync_action | string | no | none
| Optional action after a continuously
discovered object is checkpointed. Supported values are `none`, `delete`, and
`backup`.
|
+| backup_path | string | no | -
| Required when `post_sync_action=backup`.
Backup destination must not overlap with the source `path`.
|
+| retention_max_age | string | no | -
| Optional maximum age for SeaTunnel backup
objects under `backup_path`.
|
+| retention_check_interval | string | no | 1H
| Retention scan interval when backup retention
is configured.
|
| file_filter_pattern | string | no |
| Filter pattern, which used for filtering
files.
[...]
| filename_extension | string | no | -
| Filter filename extension, which used for
filtering files with specific extension. Example: `csv` `.txt` `json` `.xml`.
[...]
| common-options | | no | -
| Source plugin common parameters, please refer
to [Source Common Options](../common-options/source-common-options.md) for
details.
[...]
@@ -443,6 +457,50 @@ For more information, please refer to [Metadata
SPI](../../introduction/concepts
Whether to scan subdirectories recursively.
If `false`, subdirectories will be ignored.
+## Continuous Discovery
+
+`discovery_mode=continuous` keeps a streaming job running and polls S3 for new
or changed objects. This mode uses the existing file comparison path; it does
not consume S3 event notifications and does not emit object delete events or
changelog rows.
+
+Continuous discovery currently requires `file_format_type="binary"` and
`sync_mode="update"`. Set `target_path` to the same base path used by the sink
so the source can skip unchanged objects. The default `discovery_mode="once"`
preserves the existing bounded-source behavior.
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+}
+
+source {
+ S3File {
+ path = "/watch/source"
+ bucket = "s3a://seatunnel-test"
+ fs.s3a.endpoint = "s3.amazonaws.com"
+ fs.s3a.aws.credentials.provider =
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ secret_key = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+
+ discovery_mode = "continuous"
+ scan_interval = "10S"
+ start_mode = "earliest"
+ sync_mode = "update"
+ target_path = "/watch/target"
+ }
+}
+
+sink {
+ S3File {
+ path = "/watch/target"
+ tmp_path = "/watch/tmp"
+ bucket = "s3a://seatunnel-test"
+ fs.s3a.endpoint = "s3.amazonaws.com"
+ fs.s3a.aws.credentials.provider =
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ secret_key = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+ }
+}
+```
+
## Example
1. In this example, We read data from s3 path
`s3a://seatunnel-test/seatunnel/text` and the file type is orc in this path.
diff --git a/docs/zh/connectors/source/OssFile.md
b/docs/zh/connectors/source/OssFile.md
index 4572d66edb..cd8d7f7531 100644
--- a/docs/zh/connectors/source/OssFile.md
+++ b/docs/zh/connectors/source/OssFile.md
@@ -28,7 +28,7 @@ import ChangeLog from '../changelog/connector-file-oss.md';
使用二进制文件格式读取和写入任何格式的文件,例如视频、图片等。简而言之,任何文件都可以同步到目标位置。
- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
-- [ ] [流处理](../../introduction/concepts/connector-v2-features.md)
+- [x] [流处理](../../introduction/concepts/connector-v2-features.md)
- [x] [精确一次](../../introduction/concepts/connector-v2-features.md)
在一次pollNext调用中读取分片中的所有数据。将读取的分片保存在快照中。
@@ -213,6 +213,20 @@ schema {
| null_format | string | 否 | - |
仅在file_format_type为text时使用。null_format用于定义哪些字符串可以表示为null。例如:`\N`
|
| binary_chunk_size | int | 否 | 1024 |
仅在file_format_type为binary时使用。读取二进制文件的块大小(以字节为单位)。默认为1024字节。较大的值可能会提高大文件的性能,但会使用更多内存。
|
| binary_complete_file_mode | boolean | 否 | false |
仅在file_format_type为binary时使用。是否将完整文件作为单个块读取,而不是分割成块。启用时,整个文件内容将一次性读入内存。默认为false。
|
+| discovery_mode | string | 否 | once | 文件发现模式,支持
`once`(默认)和 `continuous`。continuous 模式会定期扫描路径,目前需要同时配置 `sync_mode=update` 和
`file_format_type=binary`。 |
+| scan_interval | string | 否 | 10S |
`discovery_mode=continuous` 的轮询间隔,支持 `10S` 等简写和 `PT10S` 等 ISO-8601 格式。 |
+| start_mode | string | 否 | earliest |
持续发现的初始扫描方式。`earliest` 会处理已有文件,`latest` 仅处理后续新增或变更。 |
+| sync_mode | string | 否 | full |
文件同步模式。`update` 会将源对象与 `target_path` 对比,只读取新增或变更对象,目前仅支持 binary 格式。 |
+| target_path | string | 否 | - |
`sync_mode=update` 时必填,通常应与 sink 的 `path` 一致。 |
+| target_hadoop_conf | map | 否 | - |
`sync_mode=update` 时可选的目标文件系统 Hadoop 配置。 |
+| update_strategy | string | 否 | distcp | update
模式的对比策略,支持 `distcp` 和 `strict`。 |
+| compare_mode | string | 否 | len_mtime | update
模式的对比方式,支持 `len_mtime` 和 `checksum`;checksum 需要 `update_strategy=strict`。 |
+| update_compare_parallelism | int | 否 | 8 |
目标对象元数据查询的最大并发数,有效范围为 `1-64`。 |
+| update_compare_bulk_threshold | int | 否 | 0 |
同一目标父目录的候选数达到正数阈值时改用一次目录枚举;`0` 表示关闭自动批量枚举。 |
+| post_sync_action | string | 否 | none | 持续发现对象完成
checkpoint 后的可选动作,支持 `none`、`delete` 和 `backup`。 |
+| backup_path | string | 否 | - |
`post_sync_action=backup` 时必填,且备份路径不能与源 `path` 重叠。 |
+| retention_max_age | string | 否 | - |
`backup_path` 中 SeaTunnel 备份对象的可选最大保留时间。 |
+| retention_check_interval | string | 否 | 1H |
配置备份保留策略时的清理扫描间隔。 |
| file_filter_pattern | string | 否 | |
过滤模式,用于过滤文件。
|
| common-options | config | 否 | - |
数据源插件通用参数,请参考[数据源通用选项](../common-options/source-common-options.md)了解详情。
|
| file_filter_modified_start | string | 否 | - |
按照最后修改时间过滤文件。 要过滤的开始时间(包括改时间),时间格式是:`yyyy-MM-dd HH:mm:ss`
|
@@ -384,6 +398,48 @@ abc.*
更多信息请参考 [元数据 SPI](../../introduction/concepts/metadata-spi.md)。
+## 持续发现
+
+`discovery_mode=continuous` 会让流式作业保持运行,并定期轮询 OSS 中的新增或变更对象。该模式使用现有文件对比逻辑,不消费
OSS 事件通知,也不会为对象删除或覆盖生成 changelog 行。
+
+持续发现目前需要同时配置 `file_format_type="binary"` 和 `sync_mode="update"`。`target_path`
应与 sink 的基础 `path` 一致,以便源端跳过未变化对象。默认的 `discovery_mode="once"` 会保持原有有界读取行为。
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+}
+
+source {
+ OssFile {
+ path = "/watch/source"
+ bucket = "oss://seatunnel-test"
+ endpoint = "oss-cn-hangzhou.aliyuncs.com"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ access_secret = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+
+ discovery_mode = "continuous"
+ scan_interval = "10S"
+ start_mode = "earliest"
+ sync_mode = "update"
+ target_path = "/watch/target"
+ }
+}
+
+sink {
+ OssFile {
+ path = "/watch/target"
+ tmp_path = "/watch/tmp"
+ bucket = "oss://seatunnel-test"
+ endpoint = "oss-cn-hangzhou.aliyuncs.com"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ access_secret = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+ }
+}
+```
+
## 如何创建Oss数据同步作业
以下示例演示如何创建从Oss读取数据并在本地客户端打印的数据同步作业:
diff --git a/docs/zh/connectors/source/S3File.md
b/docs/zh/connectors/source/S3File.md
index 9fbb78d989..f3d8a40db8 100644
--- a/docs/zh/connectors/source/S3File.md
+++ b/docs/zh/connectors/source/S3File.md
@@ -17,7 +17,7 @@ import ChangeLog from '../changelog/connector-file-s3.md';
使用二进制文件格式读取和写入任何格式的文件,例如视频、图片等。简而言之,任何文件都可以同步到目标位置。
- [x] [批处理](../../introduction/concepts/connector-v2-features.md)
-- [ ] [流处理](../../introduction/concepts/connector-v2-features.md)
+- [x] [流处理](../../introduction/concepts/connector-v2-features.md)
- [x] [精确一次](../../introduction/concepts/connector-v2-features.md)
在一次pollNext调用中读取分片中的所有数据。将读取的分片保存在快照中。
@@ -223,6 +223,20 @@ schema {
| null_format | string | 否 | -
|
仅在file_format_type为text时使用。null_format用于定义哪些字符串可以表示为null。例如:`\N`
|
| binary_chunk_size | int | 否 | 1024
|
仅在file_format_type为binary时使用。读取二进制文件的块大小(以字节为单位)。默认为1024字节。较大的值可能会提高大文件的性能,但会使用更多内存。
|
| binary_complete_file_mode | boolean | 否 | false
|
仅在file_format_type为binary时使用。是否将完整文件作为单个块读取,而不是分割成块。启用时,整个文件内容将一次性读入内存。默认为false。
|
+| discovery_mode | string | 否 | once
| 文件发现模式,支持 `once`(默认)和 `continuous`。continuous
模式会定期扫描路径并处理新增或变更对象,目前需要同时配置 `sync_mode=update` 和 `file_format_type=binary`。 |
+| scan_interval | string | 否 | 10S
| `discovery_mode=continuous` 的轮询间隔,支持 `10S` 等简写和
`PT10S` 等 ISO-8601 格式。 |
+| start_mode | string | 否 | earliest
| 持续发现的初始扫描方式。`earliest` 会处理已有文件,`latest`
会忽略作业启动时已有文件,仅处理后续新增或变更。 |
+| sync_mode | string | 否 | full
| 文件同步模式。`update` 会将源对象与 `target_path`
对比,只读取新增或变更对象,目前仅支持 binary 格式。 |
+| target_path | string | 否 | -
| `sync_mode=update` 时必填,用于按相对路径进行对比,通常应与 sink 的
`path` 一致。 |
+| target_hadoop_conf | map | 否 | -
| `sync_mode=update` 时可选的目标文件系统 Hadoop 配置。 |
+| update_strategy | string | 否 | distcp
| update 模式的对比策略,支持 `distcp` 和 `strict`。 |
+| compare_mode | string | 否 | len_mtime
| update 模式的对比方式,支持 `len_mtime` 和
`checksum`;checksum 仅能与 `update_strategy=strict` 一起使用。 |
+| update_compare_parallelism | int | 否 | 8
| 目标对象元数据查询的最大并发数,有效范围为 `1-64`。 |
+| update_compare_bulk_threshold | int | 否 | 0
| 同一目标父目录的候选数达到正数阈值时改用一次目录枚举;`0` 表示关闭自动批量枚举。 |
+| post_sync_action | string | 否 | none
| 持续发现对象完成 checkpoint 后的可选动作,支持 `none`、`delete` 和
`backup`。 |
+| backup_path | string | 否 | -
| `post_sync_action=backup` 时必填,且备份路径不能与源 `path`
重叠。 |
+| retention_max_age | string | 否 | -
| `backup_path` 中 SeaTunnel 备份对象的可选最大保留时间。 |
+| retention_check_interval | string | 否 | 1H
| 配置备份保留策略时的清理扫描间隔。 |
| file_filter_pattern | string | 否 |
| 过滤模式,用于过滤文件。
|
| filename_extension | string | 否 | -
| 过滤文件名扩展名,用于过滤具有特定扩展名的文件。例如:`csv` `.txt` `json`
`.xml`。
|
| common-options | | 否 | -
|
数据源插件通用参数,请参考[数据源通用选项](../common-options/source-common-options.md)了解详情。
|
@@ -442,6 +456,50 @@ PDF 特有的解析行为如下:
更多信息请参考 [元数据 SPI](../../introduction/concepts/metadata-spi.md)。
+## 持续发现
+
+`discovery_mode=continuous` 会让流式作业保持运行,并定期轮询 S3 中的新增或变更对象。该模式使用现有文件对比逻辑,不消费 S3
事件通知,也不会为对象删除或覆盖生成 changelog 行。
+
+持续发现目前需要同时配置 `file_format_type="binary"` 和 `sync_mode="update"`。`target_path`
应与 sink 的基础 `path` 一致,以便源端跳过未变化对象。默认的 `discovery_mode="once"` 会保持原有有界读取行为。
+
+```hocon
+env {
+ parallelism = 1
+ job.mode = "STREAMING"
+}
+
+source {
+ S3File {
+ path = "/watch/source"
+ bucket = "s3a://seatunnel-test"
+ fs.s3a.endpoint = "s3.amazonaws.com"
+ fs.s3a.aws.credentials.provider =
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ secret_key = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+
+ discovery_mode = "continuous"
+ scan_interval = "10S"
+ start_mode = "earliest"
+ sync_mode = "update"
+ target_path = "/watch/target"
+ }
+}
+
+sink {
+ S3File {
+ path = "/watch/target"
+ tmp_path = "/watch/tmp"
+ bucket = "s3a://seatunnel-test"
+ fs.s3a.endpoint = "s3.amazonaws.com"
+ fs.s3a.aws.credentials.provider =
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
+ access_key = "xxxxxxxxxxxxxxxxx"
+ secret_key = "xxxxxxxxxxxxxxxxx"
+ file_format_type = "binary"
+ }
+}
+```
+
## 示例
1. 在此示例中,我们从s3路径`s3a://seatunnel-test/seatunnel/text`读取数据,此路径中的文件类型是orc。
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-oss/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/oss/source/OssFileSourceFactory.java
b/seatunnel-connectors-v2/connector-file/connector-file-oss/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/oss/source/OssFileSourceFactory.java
index 943b02c3a1..e664d1f82d 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-oss/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/oss/source/OssFileSourceFactory.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-oss/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/oss/source/OssFileSourceFactory.java
@@ -28,6 +28,8 @@ import
org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
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.FilePostSyncAction;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileSyncMode;
import org.apache.seatunnel.connectors.seatunnel.file.config.FileSystemType;
import
org.apache.seatunnel.connectors.seatunnel.file.oss.config.OssFileSourceOptions;
@@ -114,6 +116,30 @@ public class OssFileSourceFactory implements
TableSourceFactory {
.optional(FileBaseSourceOptions.QUOTE_CHAR)
.optional(FileBaseSourceOptions.ESCAPE_CHAR)
.optional(ConnectorCommonOptions.METALAKE_TYPE)
+ .optional(
+ FileBaseSourceOptions.DISCOVERY_MODE,
+ FileBaseSourceOptions.SCAN_INTERVAL,
+ FileBaseSourceOptions.START_MODE)
+ .optional(
+ FileBaseSourceOptions.SYNC_MODE,
+ FileBaseSourceOptions.TARGET_HADOOP_CONF,
+ FileBaseSourceOptions.UPDATE_STRATEGY,
+ FileBaseSourceOptions.COMPARE_MODE,
+ FileBaseSourceOptions.UPDATE_COMPARE_PARALLELISM,
+ FileBaseSourceOptions.UPDATE_COMPARE_BULK_THRESHOLD)
+ .optional(
+ FileBaseSourceOptions.POST_SYNC_ACTION,
+ FileBaseSourceOptions.BACKUP_PATH,
+ FileBaseSourceOptions.RETENTION_MAX_AGE,
+ FileBaseSourceOptions.RETENTION_CHECK_INTERVAL)
+ .conditional(
+ FileBaseSourceOptions.SYNC_MODE,
+ FileSyncMode.UPDATE,
+ FileBaseSourceOptions.TARGET_PATH)
+ .conditional(
+ FileBaseSourceOptions.POST_SYNC_ACTION,
+ FilePostSyncAction.BACKUP,
+ FileBaseSourceOptions.BACKUP_PATH)
.optional(FileBaseSourceOptions.RECURSIVE_FILE_SCAN)
.build();
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-oss/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/oss/OssFileSourceFactoryTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-oss/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/oss/OssFileSourceFactoryTest.java
new file mode 100644
index 0000000000..cff6d753b3
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-oss/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/oss/OssFileSourceFactoryTest.java
@@ -0,0 +1,91 @@
+/*
+ * 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.oss;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseOptions;
+import
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions;
+import
org.apache.seatunnel.connectors.seatunnel.file.oss.source.OssFileSourceFactory;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class OssFileSourceFactoryTest {
+
+ @Test
+ void shouldExposeContinuousDiscoveryOptions() {
+ OptionRule optionRule = new OssFileSourceFactory().optionRule();
+
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.DISCOVERY_MODE));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.SCAN_INTERVAL));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.START_MODE));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.SYNC_MODE));
+ assertTrue(
+
optionRule.getOptionalOptions().contains(FileBaseSourceOptions.TARGET_HADOOP_CONF));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.UPDATE_STRATEGY));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.COMPARE_MODE));
+ assertTrue(
+ optionRule
+ .getOptionalOptions()
+
.contains(FileBaseSourceOptions.UPDATE_COMPARE_PARALLELISM));
+ assertTrue(
+ optionRule
+ .getOptionalOptions()
+
.contains(FileBaseSourceOptions.UPDATE_COMPARE_BULK_THRESHOLD));
+ assertTrue(
+
optionRule.getOptionalOptions().contains(FileBaseSourceOptions.POST_SYNC_ACTION));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.BACKUP_PATH));
+ assertTrue(
+
optionRule.getOptionalOptions().contains(FileBaseSourceOptions.RETENTION_MAX_AGE));
+ assertTrue(
+ optionRule
+ .getOptionalOptions()
+
.contains(FileBaseSourceOptions.RETENTION_CHECK_INTERVAL));
+ }
+
+ @Test
+ void shouldRequireTargetPathForUpdateMode() {
+ OptionRule optionRule = new OssFileSourceFactory().optionRule();
+ Map<String, Object> config = sourceConfig();
+ config.put(FileBaseSourceOptions.SYNC_MODE.key(), "update");
+
+ assertThrows(OptionValidationException.class, () -> validate(config,
optionRule));
+
+ config.put(FileBaseSourceOptions.TARGET_PATH.key(), "/target");
+ assertDoesNotThrow(() -> validate(config, optionRule));
+ }
+
+ private static Map<String, Object> sourceConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put(FileBaseOptions.FILE_PATH.key(), "/source");
+ return config;
+ }
+
+ private static void validate(Map<String, Object> config, OptionRule
optionRule) {
+
ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(optionRule);
+ }
+}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java
b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java
index b30154c6ae..083d6ee211 100644
---
a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java
+++
b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/s3/source/S3FileSourceFactory.java
@@ -27,6 +27,8 @@ import
org.apache.seatunnel.api.table.factory.TableSourceFactory;
import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
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.FilePostSyncAction;
+import org.apache.seatunnel.connectors.seatunnel.file.config.FileSyncMode;
import org.apache.seatunnel.connectors.seatunnel.file.config.FileSystemType;
import
org.apache.seatunnel.connectors.seatunnel.file.s3.config.S3FileSourceOptions;
@@ -130,6 +132,30 @@ public class S3FileSourceFactory implements
TableSourceFactory {
.optional(FileBaseSourceOptions.QUOTE_CHAR)
.optional(FileBaseSourceOptions.ESCAPE_CHAR)
.optional(ConnectorCommonOptions.METALAKE_TYPE)
+ .optional(
+ FileBaseSourceOptions.DISCOVERY_MODE,
+ FileBaseSourceOptions.SCAN_INTERVAL,
+ FileBaseSourceOptions.START_MODE)
+ .optional(
+ FileBaseSourceOptions.SYNC_MODE,
+ FileBaseSourceOptions.TARGET_HADOOP_CONF,
+ FileBaseSourceOptions.UPDATE_STRATEGY,
+ FileBaseSourceOptions.COMPARE_MODE,
+ FileBaseSourceOptions.UPDATE_COMPARE_PARALLELISM,
+ FileBaseSourceOptions.UPDATE_COMPARE_BULK_THRESHOLD)
+ .optional(
+ FileBaseSourceOptions.POST_SYNC_ACTION,
+ FileBaseSourceOptions.BACKUP_PATH,
+ FileBaseSourceOptions.RETENTION_MAX_AGE,
+ FileBaseSourceOptions.RETENTION_CHECK_INTERVAL)
+ .conditional(
+ FileBaseSourceOptions.SYNC_MODE,
+ FileSyncMode.UPDATE,
+ FileBaseSourceOptions.TARGET_PATH)
+ .conditional(
+ FileBaseSourceOptions.POST_SYNC_ACTION,
+ FilePostSyncAction.BACKUP,
+ FileBaseSourceOptions.BACKUP_PATH)
.optional(FileBaseSourceOptions.RECURSIVE_FILE_SCAN)
.build();
}
diff --git
a/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/S3FileSourceFactoryTest.java
b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/S3FileSourceFactoryTest.java
new file mode 100644
index 0000000000..a60221fbc2
--- /dev/null
+++
b/seatunnel-connectors-v2/connector-file/connector-file-s3/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/s3/S3FileSourceFactoryTest.java
@@ -0,0 +1,99 @@
+/*
+ * 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.s3;
+
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions;
+import
org.apache.seatunnel.connectors.seatunnel.file.s3.config.S3FileSourceOptions;
+import
org.apache.seatunnel.connectors.seatunnel.file.s3.source.S3FileSourceFactory;
+
+import org.junit.jupiter.api.Test;
+
+import java.util.HashMap;
+import java.util.Map;
+
+import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+
+class S3FileSourceFactoryTest {
+
+ @Test
+ void shouldExposeContinuousDiscoveryOptions() {
+ OptionRule optionRule = new S3FileSourceFactory().optionRule();
+
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.DISCOVERY_MODE));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.SCAN_INTERVAL));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.START_MODE));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.SYNC_MODE));
+ assertTrue(
+
optionRule.getOptionalOptions().contains(FileBaseSourceOptions.TARGET_HADOOP_CONF));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.UPDATE_STRATEGY));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.COMPARE_MODE));
+ assertTrue(
+ optionRule
+ .getOptionalOptions()
+
.contains(FileBaseSourceOptions.UPDATE_COMPARE_PARALLELISM));
+ assertTrue(
+ optionRule
+ .getOptionalOptions()
+
.contains(FileBaseSourceOptions.UPDATE_COMPARE_BULK_THRESHOLD));
+ assertTrue(
+
optionRule.getOptionalOptions().contains(FileBaseSourceOptions.POST_SYNC_ACTION));
+
assertTrue(optionRule.getOptionalOptions().contains(FileBaseSourceOptions.BACKUP_PATH));
+ assertTrue(
+
optionRule.getOptionalOptions().contains(FileBaseSourceOptions.RETENTION_MAX_AGE));
+ assertTrue(
+ optionRule
+ .getOptionalOptions()
+
.contains(FileBaseSourceOptions.RETENTION_CHECK_INTERVAL));
+ }
+
+ @Test
+ void shouldRequireTargetPathForUpdateMode() {
+ OptionRule optionRule = new S3FileSourceFactory().optionRule();
+ Map<String, Object> config = sourceConfig();
+ config.put(FileBaseSourceOptions.SYNC_MODE.key(), "update");
+
+ assertThrows(OptionValidationException.class, () -> validate(config,
optionRule));
+
+ config.put(FileBaseSourceOptions.TARGET_PATH.key(), "/target");
+ assertDoesNotThrow(() -> validate(config, optionRule));
+ }
+
+ private static Map<String, Object> sourceConfig() {
+ Map<String, Object> config = new HashMap<>();
+ config.put(S3FileSourceOptions.FILE_PATH.key(), "/source");
+ config.put(S3FileSourceOptions.FILE_FORMAT_TYPE.key(), "binary");
+ config.put(S3FileSourceOptions.S3_BUCKET.key(), "s3a://bucket");
+ config.put(S3FileSourceOptions.FS_S3A_ENDPOINT.key(),
"http://localhost:9000");
+ config.put(
+ S3FileSourceOptions.S3A_AWS_CREDENTIALS_PROVIDER_CLASS.key(),
+ S3FileSourceOptions.SIMPLE_AWS_CREDENTIALS_PROVIDER);
+ config.put(S3FileSourceOptions.S3_ACCESS_KEY.key(), "access-key");
+ config.put(S3FileSourceOptions.S3_SECRET_KEY.key(), "secret-key");
+ return config;
+ }
+
+ private static void validate(Map<String, Object> config, OptionRule
optionRule) {
+
ConfigValidator.of(ReadonlyConfig.fromMap(config)).validate(optionRule);
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileWithFilterIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileWithFilterIT.java
index 7fc77d239b..3bd8c61757 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileWithFilterIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3FileWithFilterIT.java
@@ -19,7 +19,9 @@ package org.apache.seatunnel.e2e.connector.file.s3;
import org.apache.seatunnel.e2e.common.container.seatunnel.SeaTunnelContainer;
import org.apache.seatunnel.e2e.common.util.ContainerUtil;
+import org.apache.seatunnel.e2e.common.util.JobIdGenerator;
+import org.awaitility.Awaitility;
import org.junit.jupiter.api.AfterAll;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.BeforeAll;
@@ -35,6 +37,8 @@ import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
import java.nio.file.Paths;
+import java.util.concurrent.CompletableFuture;
+import java.util.concurrent.TimeUnit;
/**
* MinIO-based S3 E2E test suite for connector-file-s3, covering:
@@ -159,4 +163,74 @@ public class S3FileWithFilterIT extends SeaTunnelContainer
{
executeJob("/text/s3_file_text_enable_split_to_assert.conf");
Assertions.assertEquals(0, execResult.getExitCode());
}
+
+ @Test
+ public void testS3BinaryUpdateModeContinuousDiscovery()
+ throws IOException, InterruptedException {
+ S3Utils.deletePrefix("/continuous/");
+ S3Utils.uploadContent("/continuous/src/test1.bin", "abc");
+
+ String jobId = String.valueOf(JobIdGenerator.newJobId());
+ CompletableFuture<Container.ExecResult> jobFuture =
+ CompletableFuture.supplyAsync(
+ () -> {
+ try {
+ return executeJob(
+
"/binary/s3_file_binary_update_distcp_continuous.conf",
+ jobId);
+ } catch (Exception e) {
+ throw new RuntimeException(e);
+ }
+ });
+
+ try {
+ Awaitility.await()
+ .atMost(120, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> {
+ assertContinuousJobIsRunning(jobFuture);
+ Assertions.assertTrue(
+
S3Utils.objectExists("/continuous/dst/test1.bin"));
+ Assertions.assertEquals(
+ "abc",
S3Utils.readContent("/continuous/dst/test1.bin"));
+ });
+
+ S3Utils.uploadContent("/continuous/src/test2.bin", "def");
+ Awaitility.await()
+ .atMost(120, TimeUnit.SECONDS)
+ .untilAsserted(
+ () -> {
+ assertContinuousJobIsRunning(jobFuture);
+ Assertions.assertTrue(
+
S3Utils.objectExists("/continuous/dst/test2.bin"));
+ Assertions.assertEquals(
+ "def",
S3Utils.readContent("/continuous/dst/test2.bin"));
+ });
+ } finally {
+ Container.ExecResult cancelResult = cancelJob(jobId);
+ Assertions.assertEquals(0, cancelResult.getExitCode(),
cancelResult.getStderr());
+ }
+
+ try {
+ Container.ExecResult execResult = jobFuture.get(120,
TimeUnit.SECONDS);
+ Assertions.assertEquals(0, execResult.getExitCode(),
execResult.getStderr());
+ } catch (Exception e) {
+ throw new RuntimeException("Wait continuous S3 job exit failed.",
e);
+ } finally {
+ S3Utils.deletePrefix("/continuous/");
+ }
+ }
+
+ private static void assertContinuousJobIsRunning(
+ CompletableFuture<Container.ExecResult> jobFuture) {
+ if (!jobFuture.isDone()) {
+ return;
+ }
+ Container.ExecResult result = jobFuture.join();
+ Assertions.fail(
+ "Continuous S3 job exited before cancellation. exitCode="
+ + result.getExitCode()
+ + ", stderr="
+ + result.getStderr());
+ }
}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3Utils.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3Utils.java
index 463ef0be63..8acbda4b88 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3Utils.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/java/org/apache/seatunnel/e2e/connector/file/s3/S3Utils.java
@@ -29,10 +29,14 @@ import com.amazonaws.services.s3.AmazonS3;
import com.amazonaws.services.s3.AmazonS3ClientBuilder;
import com.amazonaws.services.s3.model.ObjectMetadata;
import com.amazonaws.services.s3.model.PutObjectRequest;
+import com.amazonaws.services.s3.model.S3Object;
import java.io.ByteArrayInputStream;
+import java.io.ByteArrayOutputStream;
import java.io.File;
+import java.io.IOException;
import java.io.InputStream;
+import java.nio.charset.StandardCharsets;
public class S3Utils implements AutoCloseable {
private static Logger logger = LoggerFactory.getLogger(S3Utils.class);
@@ -89,6 +93,43 @@ public class S3Utils implements AutoCloseable {
getS3Client().putObject(putObjectRequest);
}
+ public static void uploadContent(String targetFilePath, String content) {
+ byte[] contentBytes = content.getBytes(StandardCharsets.UTF_8);
+ ObjectMetadata metadata = new ObjectMetadata();
+ metadata.setContentLength(contentBytes.length);
+ getS3Client()
+ .putObject(
+ new PutObjectRequest(
+ BUCKET,
+ targetFilePath,
+ new ByteArrayInputStream(contentBytes),
+ metadata));
+ }
+
+ public static String readContent(String targetFilePath) throws IOException
{
+ try (S3Object object = getS3Client().getObject(BUCKET, targetFilePath);
+ InputStream inputStream = object.getObjectContent();
+ ByteArrayOutputStream outputStream = new
ByteArrayOutputStream()) {
+ byte[] buffer = new byte[1024];
+ int bytesRead;
+ while ((bytesRead = inputStream.read(buffer)) != -1) {
+ outputStream.write(buffer, 0, bytesRead);
+ }
+ return new String(outputStream.toByteArray(),
StandardCharsets.UTF_8);
+ }
+ }
+
+ public static boolean objectExists(String targetFilePath) {
+ return getS3Client().doesObjectExist(BUCKET, targetFilePath);
+ }
+
+ public static void deletePrefix(String prefix) {
+ getS3Client()
+ .listObjectsV2(BUCKET, prefix)
+ .getObjectSummaries()
+ .forEach(summary -> getS3Client().deleteObject(BUCKET,
summary.getKey()));
+ }
+
@Override
public void close() throws Exception {
if (s3Client != null) {
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/resources/binary/s3_file_binary_update_distcp_continuous.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/resources/binary/s3_file_binary_update_distcp_continuous.conf
new file mode 100644
index 0000000000..f77597bf26
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-file-s3-e2e/src/test/resources/binary/s3_file_binary_update_distcp_continuous.conf
@@ -0,0 +1,65 @@
+#
+# 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 = "STREAMING"
+ checkpoint.interval = 5000
+}
+
+source {
+ S3File {
+ path = "/continuous/src"
+ file_format_type = "binary"
+ bucket = "s3a://ws-package"
+ fs.s3a.endpoint = "http://s3:9000"
+ hadoop_s3_properties = {
+ "fs.s3a.path.style.access" = "true"
+ "fs.s3a.statistics.enable" = "false"
+ }
+ fs.s3a.aws.credentials.provider =
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
+ access_key = "minioadmin"
+ secret_key = "minioadmin"
+
+ discovery_mode = "continuous"
+ scan_interval = "1S"
+ start_mode = "earliest"
+
+ sync_mode = "update"
+ target_path = "/continuous/dst"
+ update_strategy = "distcp"
+ compare_mode = "len_mtime"
+ }
+}
+
+sink {
+ S3File {
+ path = "/continuous/dst"
+ tmp_path = "/continuous/tmp"
+ file_format_type = "binary"
+ data_save_mode = "APPEND_DATA"
+ bucket = "s3a://ws-package"
+ fs.s3a.endpoint = "http://s3:9000"
+ hadoop_s3_properties = {
+ "fs.s3a.path.style.access" = "true"
+ "fs.s3a.statistics.enable" = "false"
+ }
+ fs.s3a.aws.credentials.provider =
"org.apache.hadoop.fs.s3a.SimpleAWSCredentialsProvider"
+ access_key = "minioadmin"
+ secret_key = "minioadmin"
+ }
+}