This is an automated email from the ASF dual-hosted git repository.

davidzollo 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 2091c6e8a9 [Improve][Connector-V2] Route Markdown file source splits 
by document id (#10964)
2091c6e8a9 is described below

commit 2091c6e8a91ee0788cde9023204f7ac8929119c3
Author: yzeng1618 <[email protected]>
AuthorDate: Fri Jul 31 13:17:03 2026 +0800

    [Improve][Connector-V2] Route Markdown file source splits by document id 
(#10964)
    
    Co-authored-by: zengyi <[email protected]>
---
 docs/en/connectors/source/CosFile.md               |   2 +
 docs/en/connectors/source/FtpFile.md               |   2 +
 docs/en/connectors/source/HdfsFile.md              |   2 +
 docs/en/connectors/source/LocalFile.md             |   2 +
 docs/en/connectors/source/ObsFile.md               |   2 +
 docs/en/connectors/source/OssFile.md               |   2 +
 docs/en/connectors/source/OssJindoFile.md          |   2 +
 docs/en/connectors/source/S3File.md                |   2 +
 docs/en/connectors/source/SftpFile.md              |   2 +
 docs/zh/connectors/source/CosFile.md               |   2 +
 docs/zh/connectors/source/FtpFile.md               |   2 +
 docs/zh/connectors/source/HdfsFile.md              |   2 +
 docs/zh/connectors/source/LocalFile.md             |   2 +
 docs/zh/connectors/source/ObsFile.md               |   2 +
 docs/zh/connectors/source/OssFile.md               |   2 +
 docs/zh/connectors/source/OssJindoFile.md          |   2 +
 docs/zh/connectors/source/S3File.md                |   2 +
 docs/zh/connectors/source/SftpFile.md              |   2 +
 .../seatunnel/file/source/BaseFileSource.java      |  15 ++-
 .../file/source/BaseMultipleTableFileSource.java   |  22 ++-
 .../file/source/FileSourceDocumentRouting.java     | 101 ++++++++++++++
 .../file/source/reader/MarkdownReadStrategy.java   |  52 ++-----
 .../source/split/FileSourceSplitEnumerator.java    |  50 ++++++-
 .../MultipleTableFileSourceSplitEnumerator.java    |  66 ++++++++-
 .../split/FileSourceSplitEnumeratorTest.java       | 149 ++++++++++++++++-----
 ...MultipleTableFileSourceSplitEnumeratorTest.java |  66 +++++++++
 26 files changed, 463 insertions(+), 94 deletions(-)

diff --git a/docs/en/connectors/source/CosFile.md 
b/docs/en/connectors/source/CosFile.md
index 380b5fc9f2..e1759f2e06 100644
--- a/docs/en/connectors/source/CosFile.md
+++ b/docs/en/connectors/source/CosFile.md
@@ -205,6 +205,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/FtpFile.md 
b/docs/en/connectors/source/FtpFile.md
index 1d11794f72..59cbe2d19d 100644
--- a/docs/en/connectors/source/FtpFile.md
+++ b/docs/en/connectors/source/FtpFile.md
@@ -307,6 +307,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/HdfsFile.md 
b/docs/en/connectors/source/HdfsFile.md
index 42bad28e6d..d25f4b0ccf 100644
--- a/docs/en/connectors/source/HdfsFile.md
+++ b/docs/en/connectors/source/HdfsFile.md
@@ -132,6 +132,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/LocalFile.md 
b/docs/en/connectors/source/LocalFile.md
index b270975623..63d7fec912 100644
--- a/docs/en/connectors/source/LocalFile.md
+++ b/docs/en/connectors/source/LocalFile.md
@@ -219,6 +219,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/ObsFile.md 
b/docs/en/connectors/source/ObsFile.md
index 70df651de0..bf9b51882a 100644
--- a/docs/en/connectors/source/ObsFile.md
+++ b/docs/en/connectors/source/ObsFile.md
@@ -243,6 +243,8 @@ schema {
 > - `chunk_index`: One-based chunk order in the parsed document
 > - `content_hash`: SHA-256 hash of the emitted `text` value
 >
+> When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+>
 > The option defaults to `false`, so the original Markdown schema is unchanged 
 > unless you enable it.
 >
 > Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/OssFile.md 
b/docs/en/connectors/source/OssFile.md
index c85f9859b4..c65574c611 100644
--- a/docs/en/connectors/source/OssFile.md
+++ b/docs/en/connectors/source/OssFile.md
@@ -246,6 +246,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/OssJindoFile.md 
b/docs/en/connectors/source/OssJindoFile.md
index 6764d7a523..f9179ea0bf 100644
--- a/docs/en/connectors/source/OssJindoFile.md
+++ b/docs/en/connectors/source/OssJindoFile.md
@@ -207,6 +207,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/S3File.md 
b/docs/en/connectors/source/S3File.md
index e2f0d1ac43..cab46613ba 100644
--- a/docs/en/connectors/source/S3File.md
+++ b/docs/en/connectors/source/S3File.md
@@ -256,6 +256,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/en/connectors/source/SftpFile.md 
b/docs/en/connectors/source/SftpFile.md
index 97596ff268..6b6c0ce739 100644
--- a/docs/en/connectors/source/SftpFile.md
+++ b/docs/en/connectors/source/SftpFile.md
@@ -275,6 +275,8 @@ When `markdown_rag_metadata_enabled` is set to `true`, 
SeaTunnel appends the fol
 - `chunk_index`: One-based chunk order in the parsed document
 - `content_hash`: SHA-256 hash of the emitted `text` value
 
+When this option is enabled for bounded Markdown file sources, the source 
enumerator assigns each whole-file split by the same `document_id` hash so all 
rows derived from one document stay in the same source route bucket. The 
default round-robin split assignment is unchanged when the option is disabled.
+
 The option defaults to `false`, so the original Markdown schema is unchanged 
unless you enable it.
 
 Note: Markdown format only supports reading, not writing.
diff --git a/docs/zh/connectors/source/CosFile.md 
b/docs/zh/connectors/source/CosFile.md
index 3f81b513ae..0cf1b3c85c 100644
--- a/docs/zh/connectors/source/CosFile.md
+++ b/docs/zh/connectors/source/CosFile.md
@@ -202,6 +202,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/FtpFile.md 
b/docs/zh/connectors/source/FtpFile.md
index 874f5a766c..9571e08a8f 100644
--- a/docs/zh/connectors/source/FtpFile.md
+++ b/docs/zh/connectors/source/FtpFile.md
@@ -286,6 +286,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/HdfsFile.md 
b/docs/zh/connectors/source/HdfsFile.md
index 180bb911c2..ec3d42b873 100644
--- a/docs/zh/connectors/source/HdfsFile.md
+++ b/docs/zh/connectors/source/HdfsFile.md
@@ -132,6 +132,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/LocalFile.md 
b/docs/zh/connectors/source/LocalFile.md
index 8cd76a60a3..ded499c9b8 100644
--- a/docs/zh/connectors/source/LocalFile.md
+++ b/docs/zh/connectors/source/LocalFile.md
@@ -219,6 +219,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/ObsFile.md 
b/docs/zh/connectors/source/ObsFile.md
index c914b3daa9..9faf6ecf99 100644
--- a/docs/zh/connectors/source/ObsFile.md
+++ b/docs/zh/connectors/source/ObsFile.md
@@ -109,6 +109,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/OssFile.md 
b/docs/zh/connectors/source/OssFile.md
index 57e389dfdc..edb354530a 100644
--- a/docs/zh/connectors/source/OssFile.md
+++ b/docs/zh/connectors/source/OssFile.md
@@ -273,6 +273,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/OssJindoFile.md 
b/docs/zh/connectors/source/OssJindoFile.md
index e873b535df..64dbc636c1 100644
--- a/docs/zh/connectors/source/OssJindoFile.md
+++ b/docs/zh/connectors/source/OssJindoFile.md
@@ -116,6 +116,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/S3File.md 
b/docs/zh/connectors/source/S3File.md
index 3367b8eac7..10bc105ec8 100644
--- a/docs/zh/connectors/source/S3File.md
+++ b/docs/zh/connectors/source/S3File.md
@@ -406,6 +406,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git a/docs/zh/connectors/source/SftpFile.md 
b/docs/zh/connectors/source/SftpFile.md
index 58af1adb4a..05edf11e0e 100644
--- a/docs/zh/connectors/source/SftpFile.md
+++ b/docs/zh/connectors/source/SftpFile.md
@@ -273,6 +273,8 @@ markdown 解析器提取各种元素,包括标题、段落、列表、代码
 - `chunk_index`:解析后文档中的一基 chunk 顺序
 - `content_hash`:已输出 `text` 值的 SHA-256 哈希
 
+启用该选项并读取有界 Markdown 文件时,source enumerator 会使用相同的 `document_id` 哈希分配整文件 
split,使同一文档派生的所有行留在同一个 source 路由 bucket 中。禁用该选项时,默认的轮询 split 分配行为保持不变。
+
 该选项默认值为 `false`,因此只有显式启用后才会改变原始 Markdown schema。
 
 注意:Markdown 格式仅支持读取,不支持写入。
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseFileSource.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseFileSource.java
index ecbacb8e98..63310aba43 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseFileSource.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseFileSource.java
@@ -54,6 +54,7 @@ public abstract class BaseFileSource
     private final CatalogTable catalogTable;
     private final List<String> filePaths;
     private final ReadStrategy readStrategy;
+    private final boolean documentRoutingEnabled;
 
     /** shouldn't use this construct method. just for testing */
     @VisibleForTesting
@@ -61,13 +62,17 @@ public abstract class BaseFileSource
         this.catalogTable = null;
         this.filePaths = null;
         this.readStrategy = null;
+        this.documentRoutingEnabled = false;
     }
 
     public BaseFileSource(ReadonlyConfig pluginConfig) {
         this.pluginConfig = pluginConfig;
         HadoopConf hadoopConf = initHadoopConf();
-        this.readStrategy =
-                
pluginConfig.get(FileBaseSourceOptions.FILE_FORMAT_TYPE).getReadStrategy();
+        FileFormat fileFormat = 
pluginConfig.get(FileBaseSourceOptions.FILE_FORMAT_TYPE);
+        this.readStrategy = fileFormat.getReadStrategy();
+        this.documentRoutingEnabled =
+                fileFormat == FileFormat.MARKDOWN
+                        && 
pluginConfig.get(FileBaseSourceOptions.MARKDOWN_RAG_METADATA_ENABLED);
         this.readStrategy.setPluginConfig(pluginConfig.toConfig());
         this.readStrategy.init(hadoopConf);
         String path = pluginConfig.get(FileBaseSourceOptions.FILE_PATH);
@@ -81,7 +86,6 @@ public abstract class BaseFileSource
 
         // support user-defined schema
         CatalogTable userDefinedCatalogTable;
-        FileFormat fileFormat = 
pluginConfig.get(FileBaseSourceOptions.FILE_FORMAT_TYPE);
         // only json text csv type support user-defined schema now
         if 
(pluginConfig.getOptional(ConnectorCommonOptions.SCHEMA).isPresent()) {
             switch (fileFormat) {
@@ -156,7 +160,7 @@ public abstract class BaseFileSource
     @Override
     public SourceSplitEnumerator<FileSourceSplit, FileSourceState> 
createEnumerator(
             SourceSplitEnumerator.Context<FileSourceSplit> enumeratorContext) 
throws Exception {
-        return new FileSourceSplitEnumerator(enumeratorContext, filePaths);
+        return new FileSourceSplitEnumerator(enumeratorContext, filePaths, 
documentRoutingEnabled);
     }
 
     @Override
@@ -164,6 +168,7 @@ public abstract class BaseFileSource
             SourceSplitEnumerator.Context<FileSourceSplit> enumeratorContext,
             FileSourceState checkpointState)
             throws Exception {
-        return new FileSourceSplitEnumerator(enumeratorContext, filePaths, 
checkpointState);
+        return new FileSourceSplitEnumerator(
+                enumeratorContext, filePaths, checkpointState, 
documentRoutingEnabled);
     }
 }
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseMultipleTableFileSource.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseMultipleTableFileSource.java
index 3350286d89..2c5bdbab11 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseMultipleTableFileSource.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/BaseMultipleTableFileSource.java
@@ -30,6 +30,7 @@ import 
org.apache.seatunnel.connectors.seatunnel.file.config.BaseFileSourceConfi
 import 
org.apache.seatunnel.connectors.seatunnel.file.config.BaseMultipleTableFileSourceConfig;
 import 
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions;
 import org.apache.seatunnel.connectors.seatunnel.file.config.FileDiscoveryMode;
+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.exception.FileConnectorException;
 import 
org.apache.seatunnel.connectors.seatunnel.file.source.reader.MultipleTableFileSourceReader;
@@ -43,8 +44,10 @@ import 
org.apache.seatunnel.connectors.seatunnel.file.source.split.MultipleTable
 import 
org.apache.seatunnel.connectors.seatunnel.file.source.state.FileSourceState;
 
 import java.util.HashMap;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
+import java.util.Set;
 import java.util.stream.Collectors;
 
 public abstract class BaseMultipleTableFileSource
@@ -114,7 +117,10 @@ public abstract class BaseMultipleTableFileSource
                     enumeratorContext, baseMultipleTableFileSourceConfig, 
fileSplitStrategy);
         }
         return new MultipleTableFileSourceSplitEnumerator(
-                enumeratorContext, baseMultipleTableFileSourceConfig, 
fileSplitStrategy);
+                enumeratorContext,
+                baseMultipleTableFileSourceConfig,
+                fileSplitStrategy,
+                resolveDocumentRoutingTableIds());
     }
 
     @Override
@@ -132,9 +138,23 @@ public abstract class BaseMultipleTableFileSource
                 enumeratorContext,
                 baseMultipleTableFileSourceConfig,
                 fileSplitStrategy,
+                resolveDocumentRoutingTableIds(),
                 checkpointState);
     }
 
+    private Set<String> resolveDocumentRoutingTableIds() {
+        Set<String> tableIds = new HashSet<>();
+        for (BaseFileSourceConfig config :
+                baseMultipleTableFileSourceConfig.getFileSourceConfigs()) {
+            if (config.getFileFormat() == FileFormat.MARKDOWN
+                    && config.getBaseFileSourceConfig()
+                            
.get(FileBaseSourceOptions.MARKDOWN_RAG_METADATA_ENABLED)) {
+                
tableIds.add(config.getCatalogTable().getTableId().toTablePath().toString());
+            }
+        }
+        return tableIds;
+    }
+
     private FileDiscoveryMode resolveDiscoveryMode() {
         List<BaseFileSourceConfig> configs =
                 baseMultipleTableFileSourceConfig.getFileSourceConfigs();
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/FileSourceDocumentRouting.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/FileSourceDocumentRouting.java
new file mode 100644
index 0000000000..94567a0404
--- /dev/null
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/FileSourceDocumentRouting.java
@@ -0,0 +1,101 @@
+/*
+ * 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.source;
+
+import java.net.URI;
+import java.nio.ByteBuffer;
+import java.nio.charset.StandardCharsets;
+import java.nio.file.Paths;
+import java.security.MessageDigest;
+import java.security.NoSuchAlgorithmException;
+
+/** Utilities for stable file-backed document identity and route bucket 
calculation. */
+public final class FileSourceDocumentRouting {
+
+    private static final char[] HEX_CHARS = "0123456789abcdef".toCharArray();
+
+    private FileSourceDocumentRouting() {}
+
+    /**
+     * Builds the stable file-backed document id used by markdown RAG metadata 
and split routing.
+     *
+     * @param sourceUri source URI or path
+     * @return stable document id
+     */
+    public static String buildDocumentId(String sourceUri) {
+        return "doc_" + sha256Hex(normalizeSourceUri(sourceUri));
+    }
+
+    /**
+     * Calculates the deterministic route bucket for a document id.
+     *
+     * @param documentId stable document id
+     * @param routeParallelism planned route parallelism
+     * @return bucket in the range [0, routeParallelism)
+     */
+    public static int routeBucket(String documentId, int routeParallelism) {
+        if (routeParallelism <= 0) {
+            throw new IllegalArgumentException("routeParallelism must be 
greater than zero");
+        }
+        byte[] digest = sha256(documentId == null ? "" : documentId);
+        return Math.floorMod(ByteBuffer.wrap(digest).getInt(), 
routeParallelism);
+    }
+
+    /**
+     * Normalizes local file URIs to the path form emitted by existing 
local-file reads.
+     *
+     * @param sourceUri source URI or path
+     * @return normalized source URI
+     */
+    public static String normalizeSourceUri(String sourceUri) {
+        if (sourceUri == null || !sourceUri.startsWith("file:")) {
+            return sourceUri;
+        }
+        try {
+            return Paths.get(URI.create(sourceUri)).toString();
+        } catch (IllegalArgumentException e) {
+            return sourceUri;
+        }
+    }
+
+    /**
+     * Returns the lower-case SHA-256 hexadecimal digest used by document 
metadata fields.
+     *
+     * @param value source value to hash
+     * @return lower-case SHA-256 hexadecimal digest
+     */
+    public static String sha256Hex(String value) {
+        byte[] bytes = sha256(value == null ? "" : value);
+        char[] chars = new char[bytes.length * 2];
+        for (int i = 0; i < bytes.length; i++) {
+            int unsigned = bytes[i] & 0xFF;
+            chars[i * 2] = HEX_CHARS[unsigned >>> 4];
+            chars[i * 2 + 1] = HEX_CHARS[unsigned & 0x0F];
+        }
+        return new String(chars);
+    }
+
+    private static byte[] sha256(String value) {
+        try {
+            MessageDigest digest = MessageDigest.getInstance("SHA-256");
+            return digest.digest(value.getBytes(StandardCharsets.UTF_8));
+        } catch (NoSuchAlgorithmException e) {
+            throw new IllegalStateException("SHA-256 is not available", e);
+        }
+    }
+}
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/reader/MarkdownReadStrategy.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/reader/MarkdownReadStrategy.java
index fcb1f5e95b..f1cef012f5 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/reader/MarkdownReadStrategy.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/reader/MarkdownReadStrategy.java
@@ -26,6 +26,7 @@ import org.apache.seatunnel.api.table.type.SeaTunnelRow;
 import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
 import 
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions;
 import 
org.apache.seatunnel.connectors.seatunnel.file.exception.FileConnectorException;
+import 
org.apache.seatunnel.connectors.seatunnel.file.source.FileSourceDocumentRouting;
 
 import org.apache.commons.io.IOUtils;
 
@@ -49,11 +50,7 @@ import lombok.extern.slf4j.Slf4j;
 
 import java.io.IOException;
 import java.io.InputStream;
-import java.net.URI;
 import java.nio.charset.StandardCharsets;
-import java.nio.file.Paths;
-import java.security.MessageDigest;
-import java.security.NoSuchAlgorithmException;
 import java.util.ArrayList;
 import java.util.HashMap;
 import java.util.IdentityHashMap;
@@ -65,7 +62,6 @@ public class MarkdownReadStrategy extends 
AbstractReadStrategy {
 
     private static final int DEFAULT_PAGE_NUMBER = 1;
     private static final int DEFAULT_POSITION = 1;
-    private static final char[] HEX_CHARS = "0123456789abcdef".toCharArray();
     private static final String[] DEFAULT_FIELD_NAMES = {
         "element_id",
         "element_type",
@@ -124,7 +120,7 @@ public class MarkdownReadStrategy extends 
AbstractReadStrategy {
         }
         Parser parser = Parser.builder().build();
         Node document = parser.parse(markdown);
-        String sourceUri = normalizeSourceUri(path);
+        String sourceUri = FileSourceDocumentRouting.normalizeSourceUri(path);
 
         Map<Node, NodeInfo> nodeInfoMap = new IdentityHashMap<>();
         Map<String, Integer> typeCounters = new HashMap<>();
@@ -137,7 +133,7 @@ public class MarkdownReadStrategy extends 
AbstractReadStrategy {
                 nodeInfoMap,
                 DEFAULT_PAGE_NUMBER,
                 sourceUri,
-                buildDocumentId(sourceUri));
+                FileSourceDocumentRouting.buildDocumentId(sourceUri));
 
         for (SeaTunnelRow row : rows) {
             output.collect(row);
@@ -337,10 +333,13 @@ public class MarkdownReadStrategy extends 
AbstractReadStrategy {
 
     private Object[] appendRagMetadata(
             Object[] fields, String sourceUri, String documentId, int 
chunkIndex, String text) {
-        String contentHash = sha256Hex(text == null ? "" : text);
+        String contentHash = FileSourceDocumentRouting.sha256Hex(text == null 
? "" : text);
         // Keep chunk ids stable across re-reads of the same logical document 
while still changing
         // when the chunk content changes.
-        String chunkId = "chunk_" + sha256Hex(documentId + ":" + chunkIndex + 
":" + contentHash);
+        String chunkId =
+                "chunk_"
+                        + FileSourceDocumentRouting.sha256Hex(
+                                documentId + ":" + chunkIndex + ":" + 
contentHash);
         Object[] enriched = new Object[fields.length + 
RAG_METADATA_FIELD_NAMES.length];
         System.arraycopy(fields, 0, enriched, 0, fields.length);
         enriched[fields.length] = sourceUri;
@@ -351,41 +350,6 @@ public class MarkdownReadStrategy extends 
AbstractReadStrategy {
         return enriched;
     }
 
-    private static String buildDocumentId(String sourceUri) {
-        // Document ids stay anchored to the normalized source location so 
every chunk from the same
-        // file shares one stable parent id.
-        return "doc_" + sha256Hex(sourceUri);
-    }
-
-    private static String normalizeSourceUri(String sourceUri) {
-        // Normalize local file URIs to the path form emitted by existing 
local-file reads so the
-        // metadata contract stays stable between "file:/..." and plain local 
paths.
-        if (!sourceUri.startsWith("file:")) {
-            return sourceUri;
-        }
-        try {
-            return Paths.get(URI.create(sourceUri)).toString();
-        } catch (IllegalArgumentException e) {
-            return sourceUri;
-        }
-    }
-
-    private static String sha256Hex(String value) {
-        try {
-            MessageDigest digest = MessageDigest.getInstance("SHA-256");
-            byte[] bytes = 
digest.digest(value.getBytes(StandardCharsets.UTF_8));
-            char[] chars = new char[bytes.length * 2];
-            for (int i = 0; i < bytes.length; i++) {
-                int unsigned = bytes[i] & 0xFF;
-                chars[i * 2] = HEX_CHARS[unsigned >>> 4];
-                chars[i * 2 + 1] = HEX_CHARS[unsigned & 0x0F];
-            }
-            return new String(chars);
-        } catch (NoSuchAlgorithmException e) {
-            throw new IllegalStateException("SHA-256 is not available", e);
-        }
-    }
-
     private static String[] concat(String[] left, String[] right) {
         String[] result = new String[left.length + right.length];
         System.arraycopy(left, 0, result, 0, left.length);
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumerator.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumerator.java
index dadb7066cf..112df48631 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumerator.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumerator.java
@@ -18,6 +18,7 @@
 package org.apache.seatunnel.connectors.seatunnel.file.source.split;
 
 import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+import 
org.apache.seatunnel.connectors.seatunnel.file.source.FileSourceDocumentRouting;
 import 
org.apache.seatunnel.connectors.seatunnel.file.source.state.FileSourceState;
 
 import org.slf4j.Logger;
@@ -43,13 +44,22 @@ public class FileSourceSplitEnumerator
             new TreeSet<>(Comparator.comparing(FileSourceSplit::splitId));
     private Set<FileSourceSplit> assignedSplit;
     private final List<String> filePaths;
+    private final boolean documentRoutingEnabled;
     private final Object lock = new Object();
     private final AtomicInteger assignCount = new AtomicInteger(0);
 
     public FileSourceSplitEnumerator(
             SourceSplitEnumerator.Context<FileSourceSplit> context, 
List<String> filePaths) {
+        this(context, filePaths, false);
+    }
+
+    public FileSourceSplitEnumerator(
+            SourceSplitEnumerator.Context<FileSourceSplit> context,
+            List<String> filePaths,
+            boolean documentRoutingEnabled) {
         this.context = context;
         this.filePaths = filePaths;
+        this.documentRoutingEnabled = documentRoutingEnabled;
         this.assignedSplit = new HashSet<>();
     }
 
@@ -57,7 +67,15 @@ public class FileSourceSplitEnumerator
             SourceSplitEnumerator.Context<FileSourceSplit> context,
             List<String> filePaths,
             FileSourceState sourceState) {
-        this(context, filePaths);
+        this(context, filePaths, sourceState, false);
+    }
+
+    public FileSourceSplitEnumerator(
+            SourceSplitEnumerator.Context<FileSourceSplit> context,
+            List<String> filePaths,
+            FileSourceState sourceState,
+            boolean documentRoutingEnabled) {
+        this(context, filePaths, documentRoutingEnabled);
         this.assignedSplit = sourceState.getAssignedSplit();
     }
 
@@ -97,16 +115,19 @@ public class FileSourceSplitEnumerator
 
     private void assignSplit(int taskId) {
         ArrayList<FileSourceSplit> currentTaskSplits = new ArrayList<>();
-        if (context.currentParallelism() == 1) {
+        if (!documentRoutingEnabled && context.currentParallelism() == 1) {
             // if parallelism == 1, we should assign all the splits to reader
             currentTaskSplits.addAll(allSplit);
         } else {
-            // if parallelism > 1, according to polling strategy to determine 
whether to
-            // allocate the current task
+            // if parallelism > 1, according to polling strategy or document 
routing to determine
+            // whether to allocate the current task
             assignCount.set(0);
             for (FileSourceSplit fileSourceSplit : allSplit) {
                 int splitOwner =
-                        getSplitOwner(assignCount.getAndIncrement(), 
context.currentParallelism());
+                        getSplitOwner(
+                                fileSourceSplit,
+                                assignCount.getAndIncrement(),
+                                context.currentParallelism());
                 if (splitOwner == taskId) {
                     currentTaskSplits.add(fileSourceSplit);
                 }
@@ -126,10 +147,27 @@ public class FileSourceSplitEnumerator
         context.signalNoMoreSplits(taskId);
     }
 
-    private static int getSplitOwner(int assignCount, int numReaders) {
+    private int getSplitOwner(FileSourceSplit split, int assignCount, int 
numReaders) {
+        if (documentRoutingEnabled) {
+            return getDocumentRouteOwner(split, numReaders);
+        }
+        return getRoundRobinSplitOwner(assignCount, numReaders);
+    }
+
+    private static int getRoundRobinSplitOwner(int assignCount, int 
numReaders) {
         return assignCount % numReaders;
     }
 
+    private static int getDocumentRouteOwner(FileSourceSplit split, int 
numReaders) {
+        if (split.getStart() != 0L || split.getLength() >= 0L) {
+            throw new IllegalStateException(
+                    "Document routing requires whole-file splits, but got 
split "
+                            + split.splitId());
+        }
+        String documentId = 
FileSourceDocumentRouting.buildDocumentId(split.getFilePath());
+        return FileSourceDocumentRouting.routeBucket(documentId, numReaders);
+    }
+
     @Override
     public int currentUnassignedSplitSize() {
         return allSplit.size() - assignedSplit.size();
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumerator.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumerator.java
index 92acf52c80..188714cf26 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumerator.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/main/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumerator.java
@@ -20,6 +20,7 @@ package 
org.apache.seatunnel.connectors.seatunnel.file.source.split;
 import org.apache.seatunnel.api.source.SourceSplitEnumerator;
 import 
org.apache.seatunnel.connectors.seatunnel.file.config.BaseFileSourceConfig;
 import 
org.apache.seatunnel.connectors.seatunnel.file.config.BaseMultipleTableFileSourceConfig;
+import 
org.apache.seatunnel.connectors.seatunnel.file.source.FileSourceDocumentRouting;
 import 
org.apache.seatunnel.connectors.seatunnel.file.source.state.FileSourceState;
 
 import org.apache.commons.collections4.CollectionUtils;
@@ -29,6 +30,7 @@ import lombok.extern.slf4j.Slf4j;
 import java.io.Closeable;
 import java.io.IOException;
 import java.util.ArrayList;
+import java.util.Collections;
 import java.util.Comparator;
 import java.util.HashMap;
 import java.util.HashSet;
@@ -52,32 +54,64 @@ public class MultipleTableFileSourceSplitEnumerator
     private final AtomicInteger assignCount = new AtomicInteger(0);
     private final Object lock = new Object();
     private final FileSplitStrategy fileSplitStrategy;
+    private final Set<String> documentRoutingTableIds;
 
     public MultipleTableFileSourceSplitEnumerator(
             Context<FileSourceSplit> context,
             BaseMultipleTableFileSourceConfig multipleTableFileSourceConfig,
             FileSplitStrategy fileSplitStrategy) {
+        this(context, multipleTableFileSourceConfig, fileSplitStrategy, 
Collections.emptySet());
+    }
+
+    public MultipleTableFileSourceSplitEnumerator(
+            Context<FileSourceSplit> context,
+            BaseMultipleTableFileSourceConfig multipleTableFileSourceConfig,
+            FileSplitStrategy fileSplitStrategy,
+            Set<String> documentRoutingTableIds) {
         this.context = context;
         this.fileSourceConfigs = 
multipleTableFileSourceConfig.getFileSourceConfigs();
         this.assignedSplit = new HashSet<>();
         this.allSplit = new 
TreeSet<>(Comparator.comparing(FileSourceSplit::splitId));
         this.fileSplitStrategy = fileSplitStrategy;
+        this.documentRoutingTableIds =
+                new HashSet<>(
+                        documentRoutingTableIds == null
+                                ? Collections.emptySet()
+                                : documentRoutingTableIds);
     }
 
     public MultipleTableFileSourceSplitEnumerator(
             Context<FileSourceSplit> context,
             BaseMultipleTableFileSourceConfig multipleTableFileSourceConfig,
             FileSourceState fileSourceState) {
-        this(context, multipleTableFileSourceConfig, new 
DefaultFileSplitStrategy());
-        this.assignedSplit.addAll(fileSourceState.getAssignedSplit());
+        this(
+                context,
+                multipleTableFileSourceConfig,
+                new DefaultFileSplitStrategy(),
+                Collections.emptySet(),
+                fileSourceState);
+    }
+
+    public MultipleTableFileSourceSplitEnumerator(
+            Context<FileSourceSplit> context,
+            BaseMultipleTableFileSourceConfig multipleTableFileSourceConfig,
+            FileSplitStrategy fileSplitStrategy,
+            FileSourceState fileSourceState) {
+        this(
+                context,
+                multipleTableFileSourceConfig,
+                fileSplitStrategy,
+                Collections.emptySet(),
+                fileSourceState);
     }
 
     public MultipleTableFileSourceSplitEnumerator(
             Context<FileSourceSplit> context,
             BaseMultipleTableFileSourceConfig multipleTableFileSourceConfig,
             FileSplitStrategy fileSplitStrategy,
+            Set<String> documentRoutingTableIds,
             FileSourceState fileSourceState) {
-        this(context, multipleTableFileSourceConfig, fileSplitStrategy);
+        this(context, multipleTableFileSourceConfig, fileSplitStrategy, 
documentRoutingTableIds);
         this.assignedSplit.addAll(fileSourceState.getAssignedSplit());
     }
 
@@ -145,7 +179,7 @@ public class MultipleTableFileSourceSplitEnumerator
 
     private void assignSplit(int taskId) {
         List<FileSourceSplit> currentTaskSplits = new ArrayList<>();
-        if (context.currentParallelism() == 1) {
+        if (documentRoutingTableIds.isEmpty() && context.currentParallelism() 
== 1) {
             // if parallelism == 1, we should assign all the splits to reader
             currentTaskSplits.addAll(allSplit);
         } else {
@@ -154,7 +188,10 @@ public class MultipleTableFileSourceSplitEnumerator
             assignCount.set(0);
             for (FileSourceSplit fileSourceSplit : allSplit) {
                 int splitOwner =
-                        getSplitOwner(assignCount.getAndIncrement(), 
context.currentParallelism());
+                        getSplitOwner(
+                                fileSourceSplit,
+                                assignCount.getAndIncrement(),
+                                context.currentParallelism());
                 if (splitOwner == taskId) {
                     currentTaskSplits.add(fileSourceSplit);
                 }
@@ -189,10 +226,27 @@ public class MultipleTableFileSourceSplitEnumerator
                 + " more)";
     }
 
-    private static int getSplitOwner(int assignCount, int numReaders) {
+    private int getSplitOwner(FileSourceSplit split, int assignCount, int 
numReaders) {
+        if (documentRoutingTableIds.contains(split.getTableId())) {
+            return getDocumentRouteOwner(split, numReaders);
+        }
+        return getRoundRobinSplitOwner(assignCount, numReaders);
+    }
+
+    private static int getRoundRobinSplitOwner(int assignCount, int 
numReaders) {
         return assignCount % numReaders;
     }
 
+    private static int getDocumentRouteOwner(FileSourceSplit split, int 
numReaders) {
+        if (split.getStart() != 0L || split.getLength() >= 0L) {
+            throw new IllegalStateException(
+                    "Document routing requires whole-file splits, but got 
split "
+                            + split.splitId());
+        }
+        String documentId = 
FileSourceDocumentRouting.buildDocumentId(split.getFilePath());
+        return FileSourceDocumentRouting.routeBucket(documentId, numReaders);
+    }
+
     @Override
     public void run() throws Exception {
         for (int i = 0; i < context.currentParallelism(); i++) {
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumeratorTest.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumeratorTest.java
index 0c81f69c9b..4c49562da6 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumeratorTest.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/FileSourceSplitEnumeratorTest.java
@@ -21,6 +21,8 @@ import org.apache.seatunnel.api.common.metrics.MetricsContext;
 import org.apache.seatunnel.api.event.EventListener;
 import org.apache.seatunnel.api.source.SourceEvent;
 import org.apache.seatunnel.api.source.SourceSplitEnumerator;
+import 
org.apache.seatunnel.connectors.seatunnel.file.source.FileSourceDocumentRouting;
+import 
org.apache.seatunnel.connectors.seatunnel.file.source.state.FileSourceState;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
@@ -50,38 +52,7 @@ public class FileSourceSplitEnumeratorTest {
         Map<Integer, List<FileSourceSplit>> assignSplitMap = new HashMap<>();
 
         SourceSplitEnumerator.Context<FileSourceSplit> context =
-                new SourceSplitEnumerator.Context<FileSourceSplit>() {
-                    @Override
-                    public int currentParallelism() {
-                        return parallelism;
-                    }
-
-                    @Override
-                    public Set<Integer> registeredReaders() {
-                        return null;
-                    }
-
-                    @Override
-                    public void assignSplit(int subtaskId, 
List<FileSourceSplit> splits) {
-                        assignSplitMap.put(subtaskId, splits);
-                    }
-
-                    @Override
-                    public void signalNoMoreSplits(int subtask) {}
-
-                    @Override
-                    public void sendEventToSourceReader(int subtaskId, 
SourceEvent event) {}
-
-                    @Override
-                    public MetricsContext getMetricsContext() {
-                        return null;
-                    }
-
-                    @Override
-                    public EventListener getEventListener() {
-                        return null;
-                    }
-                };
+                new CapturingEnumeratorContext(parallelism, assignSplitMap);
 
         FileSourceSplitEnumerator fileSourceSplitEnumerator =
                 new FileSourceSplitEnumerator(context, filePaths);
@@ -105,4 +76,118 @@ public class FileSourceSplitEnumeratorTest {
                     "The number of files assigned to adjacent subtasks is more 
than 1.");
         }
     }
+
+    @Test
+    void assignSplitByDocumentRouteWhenDocumentRoutingEnabled() {
+        List<String> filePaths = new ArrayList<>();
+        filePaths.add("file:/tmp/knowledge/a.md");
+        filePaths.add("file:/tmp/knowledge/b.md");
+        filePaths.add("file:/tmp/knowledge/c.md");
+        filePaths.add("file:/tmp/knowledge/d.md");
+
+        int parallelism = 4;
+        Map<Integer, List<FileSourceSplit>> assignSplitMap = new HashMap<>();
+        SourceSplitEnumerator.Context<FileSourceSplit> context =
+                new CapturingEnumeratorContext(parallelism, assignSplitMap);
+
+        FileSourceSplitEnumerator fileSourceSplitEnumerator =
+                new FileSourceSplitEnumerator(context, filePaths, true);
+        fileSourceSplitEnumerator.open();
+
+        fileSourceSplitEnumerator.run();
+
+        for (String filePath : filePaths) {
+            String documentId = 
FileSourceDocumentRouting.buildDocumentId(filePath);
+            int expectedOwner = 
FileSourceDocumentRouting.routeBucket(documentId, parallelism);
+            Assertions.assertTrue(
+                    assignSplitMap.get(expectedOwner).stream()
+                            .anyMatch(split -> 
split.getFilePath().equals(filePath)),
+                    "File should be assigned to the reader that owns its 
document route bucket.");
+        }
+    }
+
+    @Test
+    void restoredEnumeratorKeepsDocumentRouteBucket() {
+        List<String> filePaths = new ArrayList<>();
+        filePaths.add("file:/tmp/knowledge/alpha.md");
+        filePaths.add("file:/tmp/knowledge/beta.md");
+        filePaths.add("file:/tmp/knowledge/gamma.md");
+
+        int parallelism = 4;
+        Map<Integer, List<FileSourceSplit>> firstAssignSplitMap = new 
HashMap<>();
+        SourceSplitEnumerator.Context<FileSourceSplit> firstContext =
+                new CapturingEnumeratorContext(parallelism, 
firstAssignSplitMap);
+        FileSourceSplitEnumerator firstEnumerator =
+                new FileSourceSplitEnumerator(firstContext, filePaths, true);
+        firstEnumerator.open();
+        firstEnumerator.run();
+
+        FileSourceState checkpointState =
+                new FileSourceState(
+                        firstAssignSplitMap.values().stream()
+                                .flatMap(List::stream)
+                                .collect(Collectors.toSet()));
+        Map<Integer, List<FileSourceSplit>> restoredAssignSplitMap = new 
HashMap<>();
+        SourceSplitEnumerator.Context<FileSourceSplit> restoredContext =
+                new CapturingEnumeratorContext(parallelism, 
restoredAssignSplitMap);
+        FileSourceSplitEnumerator restoredEnumerator =
+                new FileSourceSplitEnumerator(restoredContext, filePaths, 
checkpointState, true);
+        restoredEnumerator.open();
+        restoredEnumerator.run();
+
+        Assertions.assertEquals(
+                splitOwnersByPath(firstAssignSplitMap), 
splitOwnersByPath(restoredAssignSplitMap));
+    }
+
+    private static Map<String, Integer> splitOwnersByPath(
+            Map<Integer, List<FileSourceSplit>> assignSplitMap) {
+        Map<String, Integer> owners = new HashMap<>();
+        assignSplitMap.forEach(
+                (owner, splits) -> splits.forEach(split -> 
owners.put(split.getFilePath(), owner)));
+        return owners;
+    }
+
+    private static class CapturingEnumeratorContext
+            implements SourceSplitEnumerator.Context<FileSourceSplit> {
+
+        private final int parallelism;
+        private final Map<Integer, List<FileSourceSplit>> assignSplitMap;
+
+        private CapturingEnumeratorContext(
+                int parallelism, Map<Integer, List<FileSourceSplit>> 
assignSplitMap) {
+            this.parallelism = parallelism;
+            this.assignSplitMap = assignSplitMap;
+        }
+
+        @Override
+        public int currentParallelism() {
+            return parallelism;
+        }
+
+        @Override
+        public Set<Integer> registeredReaders() {
+            return null;
+        }
+
+        @Override
+        public void assignSplit(int subtaskId, List<FileSourceSplit> splits) {
+            assignSplitMap.put(subtaskId, splits);
+        }
+
+        @Override
+        public void signalNoMoreSplits(int subtask) {}
+
+        @Override
+        public void sendEventToSourceReader(int subtaskId, SourceEvent event) 
{}
+
+        @Override
+        public MetricsContext getMetricsContext() {
+            return null;
+        }
+
+        @Override
+        public EventListener getEventListener() {
+            return null;
+        }
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumeratorTest.java
 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumeratorTest.java
index 6028fb48a8..83729daf01 100644
--- 
a/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumeratorTest.java
+++ 
b/seatunnel-connectors-v2/connector-file/connector-file-base/src/test/java/org/apache/seatunnel/connectors/seatunnel/file/source/split/MultipleTableFileSourceSplitEnumeratorTest.java
@@ -30,6 +30,7 @@ import 
org.apache.seatunnel.connectors.seatunnel.file.config.BaseMultipleTableFi
 import 
org.apache.seatunnel.connectors.seatunnel.file.config.FileBaseSourceOptions;
 import org.apache.seatunnel.connectors.seatunnel.file.config.HadoopConf;
 import 
org.apache.seatunnel.connectors.seatunnel.file.exception.FileConnectorException;
+import 
org.apache.seatunnel.connectors.seatunnel.file.source.FileSourceDocumentRouting;
 
 import org.apache.hadoop.conf.Configuration;
 import org.apache.hadoop.fs.FSDataInputStream;
@@ -118,6 +119,71 @@ class MultipleTableFileSourceSplitEnumeratorTest {
         Assertions.assertEquals(0, enumerator.currentUnassignedSplitSize());
     }
 
+    @Test
+    void assignSplitByDocumentRouteForEnabledTable() throws Exception {
+        int parallelism = 4;
+        List<String> filePaths =
+                Arrays.asList(
+                        "file:/tmp/knowledge/table/doc-a.md",
+                        "file:/tmp/knowledge/table/doc-b.md",
+                        "file:/tmp/knowledge/table/doc-c.md");
+
+        BaseFileSourceConfig baseFileSourceConfig = 
Mockito.mock(BaseFileSourceConfig.class);
+        
Mockito.when(baseFileSourceConfig.getFilePathsForSplitEnumerator()).thenReturn(filePaths);
+
+        CatalogTable catalogTable =
+                CatalogTable.of(
+                        TableIdentifier.of("catalog", "test", "hive_table1"),
+                        null,
+                        Maps.newHashMap(),
+                        Lists.newArrayList(),
+                        null);
+        String tableId = catalogTable.getTableId().toTablePath().toString();
+        
Mockito.when(baseFileSourceConfig.getCatalogTable()).thenReturn(catalogTable);
+
+        BaseMultipleTableFileSourceConfig baseMultipleTableFileSourceConfig =
+                Mockito.mock(BaseMultipleTableFileSourceConfig.class);
+        Mockito.when(baseMultipleTableFileSourceConfig.getFileSourceConfigs())
+                .thenReturn(Collections.singletonList(baseFileSourceConfig));
+
+        SourceSplitEnumerator.Context<FileSourceSplit> context =
+                Mockito.mock(SourceSplitEnumerator.Context.class);
+        Mockito.when(context.currentParallelism()).thenReturn(parallelism);
+        MultipleTableFileSourceSplitEnumerator enumerator =
+                new MultipleTableFileSourceSplitEnumerator(
+                        context,
+                        baseMultipleTableFileSourceConfig,
+                        new DefaultFileSplitStrategy(),
+                        Collections.singleton(tableId));
+
+        enumerator.open();
+        enumerator.run();
+
+        ArgumentCaptor<Integer> subtaskId = 
ArgumentCaptor.forClass(Integer.class);
+        ArgumentCaptor<List> split = ArgumentCaptor.forClass(List.class);
+        Mockito.verify(context, Mockito.times(parallelism))
+                .assignSplit(subtaskId.capture(), split.capture());
+
+        Map<Integer, List<FileSourceSplit>> assignedSplits = new HashMap<>();
+        for (int i = 0; i < subtaskId.getAllValues().size(); i++) {
+            assignedSplits.put(subtaskId.getAllValues().get(i), 
split.getAllValues().get(i));
+        }
+
+        for (String filePath : filePaths) {
+            String documentId = 
FileSourceDocumentRouting.buildDocumentId(filePath);
+            int expectedOwner = 
FileSourceDocumentRouting.routeBucket(documentId, parallelism);
+            Assertions.assertTrue(
+                    assignedSplits.get(expectedOwner).stream()
+                            .anyMatch(
+                                    fileSourceSplit ->
+                                            
fileSourceSplit.getFilePath().equals(filePath)
+                                                    && fileSourceSplit
+                                                            .getTableId()
+                                                            .equals(tableId)),
+                    "File should be assigned to the reader that owns its 
document route bucket.");
+        }
+    }
+
     @Test
     void deferredFileDiscoveryDoesNotListFilesDuringConfigCreation() throws 
IOException {
         CountingFileSystem.reset();

Reply via email to