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"
+  }
+}

Reply via email to