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 850748c013 [Feature][Connector-V2][Neo4j] Support multi-table source 
reads (#11869)
850748c013 is described below

commit 850748c0133ddbbe7385af03446bfe2d3f2888d1
Author: Goutam Adwant <[email protected]>
AuthorDate: Sat Aug 22 06:54:06 2026 -0700

    [Feature][Connector-V2][Neo4j] Support multi-table source reads (#11869)
    
    Signed-off-by: goutamadwant <[email protected]>
---
 docs/en/connectors/source/Neo4j.md                 |  60 +++++++-
 docs/zh/connectors/source/Neo4j.md                 |  60 +++++++-
 .../seatunnel/neo4j/source/Neo4jSource.java        |  28 +++-
 .../seatunnel/neo4j/source/Neo4jSourceFactory.java | 169 +++++++++++++++++++--
 .../seatunnel/neo4j/source/Neo4jSourceReader.java  |  93 ++++++++----
 .../neo4j/source/Neo4jSourceTableConfig.java}      |  25 +--
 .../Neo4jSourceReaderTest.java                     | 118 ++++++++++++++
 .../seatunnel/neo4j/Neo4jFactoryTest.java          | 168 ++++++++++++++++++++
 .../connector-neo4j-e2e/pom.xml                    |   6 +
 .../seatunnel/e2e/connector/neo4j/Neo4jIT.java     |  13 ++
 .../resources/neo4j/neo4j_multi_table_source.conf  | 115 ++++++++++++++
 11 files changed, 792 insertions(+), 63 deletions(-)

diff --git a/docs/en/connectors/source/Neo4j.md 
b/docs/en/connectors/source/Neo4j.md
index c729dbc427..4b98715b7c 100644
--- a/docs/en/connectors/source/Neo4j.md
+++ b/docs/en/connectors/source/Neo4j.md
@@ -23,6 +23,7 @@ the returned fields to a SeaTunnel schema.
 - [ ] [stream](../../introduction/concepts/connector-v2-features.md)
 - [ ] [exactly-once](../../introduction/concepts/connector-v2-features.md)
 - [x] [column projection](../../introduction/concepts/connector-v2-features.md)
+- [x] [support multiple table 
read](../../introduction/concepts/connector-v2-features.md)
 - [ ] [parallelism](../../introduction/concepts/connector-v2-features.md)
 - [ ] [support user-defined 
split](../../introduction/concepts/connector-v2-features.md)
 
@@ -52,15 +53,21 @@ the returned fields to a SeaTunnel schema.
 | bearer_token               | String | No       | -       | Bearer token used 
for Neo4j authentication.                                                       
                                                |
 | kerberos_ticket            | String | No       | -       | Kerberos ticket 
used for Neo4j authentication.                                                  
                                                  |
 | database                   | String | Yes      | -       | Neo4j database 
name.                                                                           
                                                   |
-| query                      | String | Yes      | -       | Cypher query used 
to read data. The fields returned by this query must match `schema.fields`.     
                                                |
-| schema                     | Object | Yes      | -       | SeaTunnel schema 
of the query result. Configure it under `schema.fields`.                        
                                                 |
+| query                      | String | Yes *    | -       | Cypher query used 
for a single-table read. The returned fields must match `schema.fields`.        
                                                 |
+| schema                     | Object | Yes *    | -       | SeaTunnel schema 
for a single-table query result. Configure it under `schema.fields`.            
                                                  |
+| tables_configs             | List   | Yes *    | -       | Multi-table read 
configuration. Each item must contain its own `query` and `schema`, including a 
unique `schema.table`.                            |
 | max_transaction_retry_time | Long   | No       | 30      | Maximum 
transaction retry time, in seconds.                                             
                                                          |
 | max_connection_timeout     | Long   | No       | 30      | Maximum time to 
wait for a TCP connection to be established, in seconds.                        
                                                  |
 
+> * Configure either the root-level `query` and `schema`, or `tables_configs`.
+
 ## Notes
 
 - Use exactly one authentication method: username/password, bearer token, or 
Kerberos ticket.
 - `query` controls which fields are returned. `schema.fields` must list the 
returned field names and their SeaTunnel types.
+- In multi-table mode, keep connection and authentication options at the root 
level. Each `tables_configs` item defines one `query` and one `schema`.
+- Every multi-table `schema` must set a unique `table`. Rows use this value as 
their table ID for downstream routing.
+- Multi-table queries run in declaration order through one Neo4j driver and 
session. The source remains bounded and uses one reader.
 - Returned field names can contain dots, such as `t.string`, when the Cypher 
query returns properties from a node.
 - `MAP` fields must use `STRING` keys, for example `MAP<STRING, INT>`.
 - Neo4j integer and floating-point values are converted according to the 
SeaTunnel type declared in `schema.fields`. Use `BIGINT`/`DOUBLE` when the 
value may exceed the range of `INT`/`FLOAT`.
@@ -109,6 +116,55 @@ sink {
 }
 ```
 
+### Multi-table read
+
+```hocon
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+}
+
+source {
+  Neo4j {
+    uri = "neo4j://localhost:7687"
+    username = "neo4j"
+    password = "password"
+    database = "neo4j"
+
+    tables_configs = [
+      {
+        query = "MATCH (p:Person) RETURN p.name AS name"
+        schema {
+          table = "people"
+          fields {
+            name = STRING
+          }
+        }
+      },
+      {
+        query = "MATCH (c:Company) RETURN c.name AS name"
+        schema {
+          table = "companies"
+          fields {
+            name = STRING
+          }
+        }
+      }
+    ]
+  }
+}
+
+sink {
+  Console {
+    plugin_input = "people"
+  }
+
+  Console {
+    plugin_input = "companies"
+  }
+}
+```
+
 ## Changelog
 
 <ChangeLog />
diff --git a/docs/zh/connectors/source/Neo4j.md 
b/docs/zh/connectors/source/Neo4j.md
index 92ef02ccc5..1ebcef8d83 100644
--- a/docs/zh/connectors/source/Neo4j.md
+++ b/docs/zh/connectors/source/Neo4j.md
@@ -23,6 +23,7 @@ Neo4j 源连接器通过执行 Cypher 查询从 Neo4j 读取数据,并把查
 - [ ] [流处理](../../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)
 - [ ] [并行度](../../introduction/concepts/connector-v2-features.md)
 - [ ] [支持用户自定义切分](../../introduction/concepts/connector-v2-features.md)
 
@@ -52,15 +53,21 @@ Neo4j 源连接器通过执行 Cypher 查询从 Neo4j 读取数据,并把查
 | bearer_token               | String | 否    | -   | 用于 Neo4j 认证的 bearer 
token。                                                                       |
 | kerberos_ticket            | String | 否    | -   | 用于 Neo4j 认证的 Kerberos 
ticket。                                                                    |
 | database                   | String | 是    | -   | Neo4j 数据库名。               
                                                                        |
-| query                      | String | 是    | -   | 读取数据使用的 Cypher 
查询语句。查询返回字段必须和 `schema.fields` 对应。                                           |
-| schema                     | Object | 是    | -   | 查询结果对应的 SeaTunnel 表结构,在 
`schema.fields` 中配置字段名和类型。                                           |
+| query                      | String | 是 *  | -   | 单表读取使用的 Cypher 
查询语句,返回字段必须和 `schema.fields` 对应。                                               |
+| schema                     | Object | 是 *  | -   | 单表查询结果对应的 SeaTunnel 表结构,在 
`schema.fields` 中配置字段名和类型。                                        |
+| tables_configs             | List   | 是 *  | -   | 多表读取配置。每个配置项必须包含自己的 
`query` 和 `schema`,并设置唯一的 `schema.table`。                           |
 | max_transaction_retry_time | Long   | 否    | 30  | 最大事务重试时间,单位为秒。            
                                                                       |
 | max_connection_timeout     | Long   | 否    | 30  | 建立 TCP 连接的最大等待时间,单位为秒。    
                                                                      |
 
+> * 配置根级别的 `query` 和 `schema`,或者配置 `tables_configs`,二者选择其一。
+
 ## 注意事项
 
 - 认证方式只选一种:用户名密码、bearer token 或 Kerberos ticket。
 - `query` 决定返回哪些字段,`schema.fields` 必须写清这些返回字段和对应类型。
+- 多表模式下,连接和认证选项放在根级别;每个 `tables_configs` 配置项定义一个 `query` 和一个 `schema`。
+- 每个多表 `schema` 必须设置唯一的 `table`,该值会作为数据行的表 ID,用于下游路由。
+- 多表查询按配置顺序执行,并复用同一个 Neo4j driver 和 session。该 source 仍为有界单 reader source。
 - 查询返回字段名可以包含点号,例如从节点属性返回的 `t.string`。
 - `MAP` 字段的 key 必须是 `STRING`,例如 `MAP<STRING, INT>`。
 - Neo4j 的整数和浮点数会按 `schema.fields` 中声明的 SeaTunnel 类型转换;如果数值可能超过 `INT` 或 `FLOAT` 
范围,建议使用 `BIGINT` 或 `DOUBLE`。
@@ -109,6 +116,55 @@ sink {
 }
 ```
 
+### 多表读取
+
+```hocon
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+}
+
+source {
+  Neo4j {
+    uri = "neo4j://localhost:7687"
+    username = "neo4j"
+    password = "password"
+    database = "neo4j"
+
+    tables_configs = [
+      {
+        query = "MATCH (p:Person) RETURN p.name AS name"
+        schema {
+          table = "people"
+          fields {
+            name = STRING
+          }
+        }
+      },
+      {
+        query = "MATCH (c:Company) RETURN c.name AS name"
+        schema {
+          table = "companies"
+          fields {
+            name = STRING
+          }
+        }
+      }
+    ]
+  }
+}
+
+sink {
+  Console {
+    plugin_input = "people"
+  }
+
+  Console {
+    plugin_input = "companies"
+  }
+}
+```
+
 ## 变更日志
 
 <ChangeLog />
diff --git 
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSource.java
 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSource.java
index 65fd0ecf4c..76729f02d6 100644
--- 
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSource.java
+++ 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSource.java
@@ -21,12 +21,12 @@ import org.apache.seatunnel.api.source.Boundedness;
 import org.apache.seatunnel.api.source.SupportColumnProjection;
 import org.apache.seatunnel.api.table.catalog.CatalogTable;
 import org.apache.seatunnel.api.table.type.SeaTunnelRow;
-import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
 import 
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
 import 
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitSource;
 import 
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
 import 
org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jSourceQueryInfo;
 
+import java.util.ArrayList;
 import java.util.Collections;
 import java.util.List;
 
@@ -35,14 +35,28 @@ import static 
org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jSource
 public class Neo4jSource extends AbstractSingleSplitSource<SeaTunnelRow>
         implements SupportColumnProjection {
 
-    private final CatalogTable catalogTable;
+    private final List<CatalogTable> catalogTables;
     private final Neo4jSourceQueryInfo neo4jSourceQueryInfo;
-    private final SeaTunnelRowType rowType;
+    private final List<Neo4jSourceTableConfig> tableConfigs;
 
     public Neo4jSource(CatalogTable catalogTable, Neo4jSourceQueryInfo 
neo4jSourceQueryInfo) {
-        this.catalogTable = catalogTable;
+        this.catalogTables = Collections.singletonList(catalogTable);
         this.neo4jSourceQueryInfo = neo4jSourceQueryInfo;
-        this.rowType = catalogTable.getSeaTunnelRowType();
+        this.tableConfigs =
+                Collections.singletonList(
+                        new Neo4jSourceTableConfig(
+                                neo4jSourceQueryInfo.getQuery(),
+                                catalogTable.getSeaTunnelRowType(),
+                                null));
+    }
+
+    Neo4jSource(
+            List<CatalogTable> catalogTables,
+            Neo4jSourceQueryInfo neo4jSourceQueryInfo,
+            List<Neo4jSourceTableConfig> tableConfigs) {
+        this.catalogTables = Collections.unmodifiableList(new 
ArrayList<>(catalogTables));
+        this.neo4jSourceQueryInfo = neo4jSourceQueryInfo;
+        this.tableConfigs = Collections.unmodifiableList(new 
ArrayList<>(tableConfigs));
     }
 
     @Override
@@ -57,12 +71,12 @@ public class Neo4jSource extends 
AbstractSingleSplitSource<SeaTunnelRow>
 
     @Override
     public List<CatalogTable> getProducedCatalogTables() {
-        return Collections.singletonList(catalogTable);
+        return catalogTables;
     }
 
     @Override
     public AbstractSingleSplitReader<SeaTunnelRow> createReader(
             SingleSplitReaderContext readerContext) throws Exception {
-        return new Neo4jSourceReader(readerContext, neo4jSourceQueryInfo, 
rowType);
+        return new Neo4jSourceReader(readerContext, neo4jSourceQueryInfo, 
tableConfigs);
     }
 }
diff --git 
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
index d7e924d76e..eae2e8fc79 100644
--- 
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
+++ 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceFactory.java
@@ -17,10 +17,18 @@
 
 package org.apache.seatunnel.connectors.seatunnel.neo4j.source;
 
+import org.apache.seatunnel.shade.com.typesafe.config.ConfigValueFactory;
+
+import org.apache.seatunnel.api.common.SeaTunnelAPIErrorCode;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConditionExtension;
+import org.apache.seatunnel.api.configuration.util.Conditions;
 import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
 import org.apache.seatunnel.api.options.ConnectorCommonOptions;
 import org.apache.seatunnel.api.source.SeaTunnelSource;
 import org.apache.seatunnel.api.source.SourceSplit;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
 import org.apache.seatunnel.api.table.catalog.CatalogTableUtil;
 import org.apache.seatunnel.api.table.connector.TableSource;
 import org.apache.seatunnel.api.table.factory.Factory;
@@ -28,10 +36,16 @@ import 
org.apache.seatunnel.api.table.factory.TableSourceFactory;
 import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
 import 
org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jSourceOptions;
 import 
org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jSourceQueryInfo;
+import 
org.apache.seatunnel.connectors.seatunnel.neo4j.exception.Neo4jConnectorException;
 
 import com.google.auto.service.AutoService;
 
 import java.io.Serializable;
+import java.util.ArrayList;
+import java.util.HashSet;
+import java.util.List;
+import java.util.Map;
+import java.util.Set;
 
 @AutoService(Factory.class)
 public class Neo4jSourceFactory implements TableSourceFactory {
@@ -43,18 +57,26 @@ public class Neo4jSourceFactory implements 
TableSourceFactory {
     @Override
     public OptionRule optionRule() {
         return OptionRule.builder()
-                .required(
-                        Neo4jSourceOptions.KEY_NEO4J_URI,
-                        Neo4jSourceOptions.KEY_DATABASE,
+                .required(Neo4jSourceOptions.KEY_NEO4J_URI, 
Neo4jSourceOptions.KEY_DATABASE)
+                .exclusive(Neo4jSourceOptions.KEY_QUERY, 
ConnectorCommonOptions.TABLE_CONFIGS)
+                .optional(
                         Neo4jSourceOptions.KEY_QUERY,
-                        ConnectorCommonOptions.SCHEMA)
+                        Conditions.notBlank(Neo4jSourceOptions.KEY_QUERY),
+                        Conditions.extension(
+                                Neo4jSourceOptions.KEY_QUERY, new 
SingleTableConfigValidator()))
+                .optional(
+                        ConnectorCommonOptions.TABLE_CONFIGS,
+                        
Conditions.notEmpty(ConnectorCommonOptions.TABLE_CONFIGS),
+                        Conditions.extension(
+                                ConnectorCommonOptions.TABLE_CONFIGS, new 
TableConfigsValidator()))
                 .optional(
                         Neo4jSourceOptions.KEY_USERNAME,
                         Neo4jSourceOptions.KEY_PASSWORD,
                         Neo4jSourceOptions.KEY_BEARER_TOKEN,
                         Neo4jSourceOptions.KEY_KERBEROS_TICKET,
                         Neo4jSourceOptions.KEY_MAX_CONNECTION_TIMEOUT,
-                        Neo4jSourceOptions.KEY_MAX_TRANSACTION_RETRY_TIME)
+                        Neo4jSourceOptions.KEY_MAX_TRANSACTION_RETRY_TIME,
+                        ConnectorCommonOptions.SCHEMA)
                 .build();
     }
 
@@ -66,12 +88,135 @@ public class Neo4jSourceFactory implements 
TableSourceFactory {
     @Override
     public <T, SplitT extends SourceSplit, StateT extends Serializable>
             TableSource<T, SplitT, StateT> 
createSource(TableSourceFactoryContext context) {
-        Neo4jSourceQueryInfo neo4jSourceQueryInfo =
-                new Neo4jSourceQueryInfo(context.getOptions().toConfig());
-        return () ->
-                (SeaTunnelSource<T, SplitT, StateT>)
-                        new Neo4jSource(
-                                
CatalogTableUtil.buildWithConfig(context.getOptions()),
-                                neo4jSourceQueryInfo);
+        return () -> (SeaTunnelSource<T, SplitT, StateT>) 
createNeo4jSource(context.getOptions());
+    }
+
+    private Neo4jSource createNeo4jSource(ReadonlyConfig config) {
+        if 
(!config.getOptional(ConnectorCommonOptions.TABLE_CONFIGS).isPresent()) {
+            return new Neo4jSource(
+                    CatalogTableUtil.buildWithConfig(config),
+                    new Neo4jSourceQueryInfo(config.toConfig()));
+        }
+
+        List<Map<String, Object>> entries = 
config.get(ConnectorCommonOptions.TABLE_CONFIGS);
+        if (entries.isEmpty()) {
+            throw configError("'tables_configs' must not be empty");
+        }
+
+        List<CatalogTable> catalogTables = new ArrayList<>(entries.size());
+        List<Neo4jSourceTableConfig> tableConfigs = new 
ArrayList<>(entries.size());
+        Set<String> tableIds = new HashSet<>();
+
+        for (int i = 0; i < entries.size(); i++) {
+            ReadonlyConfig tableConfig = 
ReadonlyConfig.fromMap(entries.get(i));
+            String query = 
tableConfig.getOptional(Neo4jSourceOptions.KEY_QUERY).orElse(null);
+            if (query == null || query.trim().isEmpty()) {
+                throw configError(
+                        String.format(
+                                "tables_configs[%d]: 'query' must be 
configured and non-blank", i));
+            }
+
+            CatalogTable catalogTable;
+            try {
+                catalogTable = CatalogTableUtil.buildWithConfig(tableConfig);
+            } catch (RuntimeException e) {
+                throw new Neo4jConnectorException(
+                        SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED,
+                        String.format("tables_configs[%d]: invalid 'schema' 
configuration", i),
+                        e);
+            }
+
+            String tableId = 
catalogTable.getTableId().toTablePath().toString();
+            if (!tableIds.add(tableId)) {
+                throw configError(
+                        String.format(
+                                "Duplicate schema.table '%s' found in 
tables_configs", tableId));
+            }
+
+            catalogTables.add(catalogTable);
+            tableConfigs.add(
+                    new Neo4jSourceTableConfig(query, 
catalogTable.getSeaTunnelRowType(), tableId));
+        }
+
+        Neo4jSourceQueryInfo connectionInfo =
+                new Neo4jSourceQueryInfo(
+                        config.toConfig()
+                                .withValue(
+                                        Neo4jSourceOptions.KEY_QUERY.key(),
+                                        ConfigValueFactory.fromAnyRef(
+                                                
tableConfigs.get(0).getQuery())));
+        return new Neo4jSource(catalogTables, connectionInfo, tableConfigs);
+    }
+
+    private static Neo4jConnectorException configError(String message) {
+        return new 
Neo4jConnectorException(SeaTunnelAPIErrorCode.CONFIG_VALIDATION_FAILED, 
message);
+    }
+
+    static class SingleTableConfigValidator implements 
ConditionExtension<String> {
+
+        @Override
+        public String description() {
+            return "'schema' must be configured when using a root-level 
'query'";
+        }
+
+        @Override
+        public boolean evaluate(ReadonlyConfig config, String query)
+                throws OptionValidationException {
+            Map<String, Object> schema =
+                    
config.getOptional(ConnectorCommonOptions.SCHEMA).orElse(null);
+            if (schema == null || schema.isEmpty()) {
+                throw new OptionValidationException(
+                        "'schema' must be configured when using a root-level 
'query'");
+            }
+            return true;
+        }
+    }
+
+    static class TableConfigsValidator implements 
ConditionExtension<List<Map<String, Object>>> {
+
+        @Override
+        public String description() {
+            return "each 'tables_configs' entry must contain a non-blank 
'query' and a schema with a unique 'table'";
+        }
+
+        @Override
+        public boolean evaluate(ReadonlyConfig config, List<Map<String, 
Object>> entries)
+                throws OptionValidationException {
+            if (config.getOptional(ConnectorCommonOptions.SCHEMA).isPresent()) 
{
+                throw new OptionValidationException(
+                        "root-level 'schema' cannot be used with 
'tables_configs'");
+            }
+
+            Set<String> tableIds = new HashSet<>();
+            for (int i = 0; i < entries.size(); i++) {
+                Map<String, Object> entry = entries.get(i);
+                Object query = entry.get(Neo4jSourceOptions.KEY_QUERY.key());
+                if (!(query instanceof String) || ((String) 
query).trim().isEmpty()) {
+                    throw new OptionValidationException(
+                            "tables_configs[%d]: 'query' must be configured 
and non-blank", i);
+                }
+
+                Object schemaValue = 
entry.get(ConnectorCommonOptions.SCHEMA.key());
+                if (!(schemaValue instanceof Map) || ((Map<?, ?>) 
schemaValue).isEmpty()) {
+                    throw new OptionValidationException(
+                            "tables_configs[%d]: 'schema' must be configured 
and non-empty", i);
+                }
+
+                Object tableValue =
+                        ((Map<?, ?>) 
schemaValue).get(ConnectorCommonOptions.TABLE.key());
+                if (!(tableValue instanceof String) || ((String) 
tableValue).trim().isEmpty()) {
+                    throw new OptionValidationException(
+                            "tables_configs[%d]: 'schema.table' must be 
configured and non-blank",
+                            i);
+                }
+
+                String tableId = ((String) tableValue).trim();
+                if (!tableIds.add(tableId)) {
+                    throw new OptionValidationException(
+                            "tables_configs[%d]: duplicate 'schema.table' 
value '%s'", i, tableId);
+                }
+            }
+            return true;
+        }
     }
 }
diff --git 
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceReader.java
 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceReader.java
index 283e74f98b..a3ab0c0942 100644
--- 
a/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceReader.java
+++ 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceReader.java
@@ -32,6 +32,7 @@ import 
org.apache.seatunnel.connectors.seatunnel.neo4j.exception.Neo4jConnectorE
 
 import org.neo4j.driver.Driver;
 import org.neo4j.driver.Query;
+import org.neo4j.driver.Record;
 import org.neo4j.driver.Result;
 import org.neo4j.driver.Session;
 import org.neo4j.driver.SessionConfig;
@@ -40,14 +41,16 @@ import org.neo4j.driver.exceptions.value.LossyCoercion;
 
 import java.io.IOException;
 import java.lang.reflect.Array;
+import java.util.ArrayList;
+import java.util.Collections;
 import java.util.List;
 import java.util.Objects;
 
 public class Neo4jSourceReader extends AbstractSingleSplitReader<SeaTunnelRow> 
{
 
     private final SingleSplitReaderContext context;
-    private final Neo4jSourceQueryInfo neo4jSourceQueryInfo;
-    private final SeaTunnelRowType rowType;
+    private final String database;
+    private final List<Neo4jSourceTableConfig> tableConfigs;
     private final Driver driver;
     private Session session;
 
@@ -55,18 +58,27 @@ public class Neo4jSourceReader extends 
AbstractSingleSplitReader<SeaTunnelRow> {
             SingleSplitReaderContext context,
             Neo4jSourceQueryInfo neo4jSourceQueryInfo,
             SeaTunnelRowType rowType) {
+        this(
+                context,
+                neo4jSourceQueryInfo,
+                Collections.singletonList(
+                        new Neo4jSourceTableConfig(
+                                neo4jSourceQueryInfo.getQuery(), rowType, 
null)));
+    }
+
+    Neo4jSourceReader(
+            SingleSplitReaderContext context,
+            Neo4jSourceQueryInfo neo4jSourceQueryInfo,
+            List<Neo4jSourceTableConfig> tableConfigs) {
         this.context = context;
-        this.neo4jSourceQueryInfo = neo4jSourceQueryInfo;
+        this.database = neo4jSourceQueryInfo.getDriverBuilder().getDatabase();
+        this.tableConfigs = Collections.unmodifiableList(new 
ArrayList<>(tableConfigs));
         this.driver = neo4jSourceQueryInfo.getDriverBuilder().build();
-        this.rowType = rowType;
     }
 
     @Override
     public void open() throws Exception {
-        this.session =
-                driver.session(
-                        SessionConfig.forDatabase(
-                                
neo4jSourceQueryInfo.getDriverBuilder().getDatabase()));
+        this.session = driver.session(SessionConfig.forDatabase(database));
     }
 
     @Override
@@ -112,27 +124,50 @@ public class Neo4jSourceReader extends 
AbstractSingleSplitReader<SeaTunnelRow> {
 
     @Override
     public void internalPollNext(Collector<SeaTunnelRow> output) throws 
Exception {
-        final Query query = new Query(neo4jSourceQueryInfo.getQuery());
-        session.readTransaction(
-                tx -> {
-                    final Result result = tx.run(query);
-                    result.stream()
-                            .forEach(
-                                    row -> {
-                                        final Object[] fields =
-                                                new 
Object[rowType.getTotalFields()];
-                                        for (int i = 0; i < 
rowType.getTotalFields(); i++) {
-                                            final String fieldName = 
rowType.getFieldName(i);
-                                            final SeaTunnelDataType<?> 
fieldType =
-                                                    rowType.getFieldType(i);
-                                            final Value value = 
row.get(fieldName);
-                                            fields[i] = convertType(fieldType, 
value);
-                                        }
-                                        output.collect(new 
SeaTunnelRow(fields));
-                                    });
-                    return null;
-                });
-        this.context.signalNoMoreElement();
+        try {
+            for (Neo4jSourceTableConfig tableConfig : tableConfigs) {
+                readTable(output, tableConfig);
+            }
+        } finally {
+            this.context.signalNoMoreElement();
+        }
+    }
+
+    private void readTable(Collector<SeaTunnelRow> output, 
Neo4jSourceTableConfig tableConfig) {
+        final Query query = new Query(tableConfig.getQuery());
+        try {
+            session.readTransaction(
+                    tx -> {
+                        final Result result = tx.run(query);
+                        result.stream()
+                                .forEach(row -> 
output.collect(convertRecord(row, tableConfig)));
+                        return null;
+                    });
+        } catch (RuntimeException exception) {
+            if (tableConfig.getTableId() == null) {
+                throw exception;
+            }
+            throw new Neo4jConnectorException(
+                    CommonErrorCodeDeprecated.READER_OPERATION_FAILED,
+                    "Failed to read Neo4j table '" + tableConfig.getTableId() 
+ "'.",
+                    exception);
+        }
+    }
+
+    static SeaTunnelRow convertRecord(Record record, Neo4jSourceTableConfig 
tableConfig) {
+        SeaTunnelRowType rowType = tableConfig.getRowType();
+        Object[] fields = new Object[rowType.getTotalFields()];
+        for (int i = 0; i < rowType.getTotalFields(); i++) {
+            String fieldName = rowType.getFieldName(i);
+            SeaTunnelDataType<?> fieldType = rowType.getFieldType(i);
+            Value value = record.get(fieldName);
+            fields[i] = convertType(fieldType, value);
+        }
+        SeaTunnelRow seaTunnelRow = new SeaTunnelRow(fields);
+        if (tableConfig.getTableId() != null) {
+            seaTunnelRow.setTableId(tableConfig.getTableId());
+        }
+        return seaTunnelRow;
     }
 
     /**
diff --git 
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceTableConfig.java
similarity index 61%
copy from 
seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
copy to 
seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceTableConfig.java
index a77813d702..f52e7ece4e 100644
--- 
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
+++ 
b/seatunnel-connectors-v2/connector-neo4j/src/main/java/org/apache/seatunnel/connectors/seatunnel/neo4j/source/Neo4jSourceTableConfig.java
@@ -15,19 +15,22 @@
  * limitations under the License.
  */
 
-package org.apache.seatunnel.connectors.seatunnel.neo4j;
+package org.apache.seatunnel.connectors.seatunnel.neo4j.source;
 
-import org.apache.seatunnel.connectors.seatunnel.neo4j.sink.Neo4jSinkFactory;
-import 
org.apache.seatunnel.connectors.seatunnel.neo4j.source.Neo4jSourceFactory;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
 
-import org.junit.jupiter.api.Assertions;
-import org.junit.jupiter.api.Test;
+import lombok.AllArgsConstructor;
+import lombok.Getter;
 
-class Neo4jFactoryTest {
+import java.io.Serializable;
 
-    @Test
-    void optionRule() {
-        Assertions.assertNotNull((new Neo4jSourceFactory()).optionRule());
-        Assertions.assertNotNull((new Neo4jSinkFactory()).optionRule());
-    }
+@Getter
+@AllArgsConstructor
+final class Neo4jSourceTableConfig implements Serializable {
+
+    private static final long serialVersionUID = 1L;
+
+    private final String query;
+    private final SeaTunnelRowType rowType;
+    private final String tableId;
 }
diff --git 
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org.apache.seatunnel.connectors.seatunnel.neo4j.source/Neo4jSourceReaderTest.java
 
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org.apache.seatunnel.connectors.seatunnel.neo4j.source/Neo4jSourceReaderTest.java
index e839f62ad3..aee6f7f0bc 100644
--- 
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org.apache.seatunnel.connectors.seatunnel.neo4j.source/Neo4jSourceReaderTest.java
+++ 
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org.apache.seatunnel.connectors.seatunnel.neo4j.source/Neo4jSourceReaderTest.java
@@ -17,14 +17,28 @@
 
 package org.apache.seatunnel.connectors.seatunnel.neo4j.source;
 
+import org.apache.seatunnel.api.source.Collector;
 import org.apache.seatunnel.api.table.type.BasicType;
 import org.apache.seatunnel.api.table.type.LocalTimeType;
 import org.apache.seatunnel.api.table.type.MapType;
 import org.apache.seatunnel.api.table.type.PrimitiveByteArrayType;
+import org.apache.seatunnel.api.table.type.SeaTunnelDataType;
+import org.apache.seatunnel.api.table.type.SeaTunnelRow;
+import org.apache.seatunnel.api.table.type.SeaTunnelRowType;
+import 
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
+import org.apache.seatunnel.connectors.seatunnel.neo4j.config.DriverBuilder;
+import 
org.apache.seatunnel.connectors.seatunnel.neo4j.config.Neo4jSourceQueryInfo;
 import 
org.apache.seatunnel.connectors.seatunnel.neo4j.exception.Neo4jConnectorException;
 
 import org.junit.jupiter.api.Test;
+import org.neo4j.driver.Driver;
+import org.neo4j.driver.Session;
+import org.neo4j.driver.SessionConfig;
+import org.neo4j.driver.TransactionWork;
+import org.neo4j.driver.Value;
+import org.neo4j.driver.exceptions.ServiceUnavailableException;
 import org.neo4j.driver.exceptions.value.LossyCoercion;
+import org.neo4j.driver.internal.InternalRecord;
 import org.neo4j.driver.internal.value.BooleanValue;
 import org.neo4j.driver.internal.value.BytesValue;
 import org.neo4j.driver.internal.value.DateValue;
@@ -46,9 +60,96 @@ import static 
org.apache.seatunnel.api.table.type.ArrayType.STRING_ARRAY_TYPE;
 import static org.junit.jupiter.api.Assertions.assertArrayEquals;
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.when;
 
 class Neo4jSourceReaderTest {
+
+    @Test
+    void mapsRowsUsingTableSpecificSchemaAndTableId() {
+        SeaTunnelRowType peopleRowType =
+                new SeaTunnelRowType(
+                        new String[] {"name"}, new SeaTunnelDataType<?>[] 
{BasicType.STRING_TYPE});
+        Neo4jSourceTableConfig peopleConfig =
+                new Neo4jSourceTableConfig("people query", peopleRowType, 
"people");
+        InternalRecord peopleRecord =
+                new InternalRecord(
+                        Collections.singletonList("name"), new Value[] {new 
StringValue("Alice")});
+
+        SeaTunnelRow peopleRow = Neo4jSourceReader.convertRecord(peopleRecord, 
peopleConfig);
+
+        assertEquals("Alice", peopleRow.getField(0));
+        assertEquals("people", peopleRow.getTableId());
+
+        SeaTunnelRowType companiesRowType =
+                new SeaTunnelRowType(
+                        new String[] {"id"}, new SeaTunnelDataType<?>[] 
{BasicType.INT_TYPE});
+        Neo4jSourceTableConfig companiesConfig =
+                new Neo4jSourceTableConfig("companies query", 
companiesRowType, "companies");
+        InternalRecord companiesRecord =
+                new InternalRecord(
+                        Collections.singletonList("id"), new Value[] {new 
IntegerValue(7)});
+
+        SeaTunnelRow companiesRow =
+                Neo4jSourceReader.convertRecord(companiesRecord, 
companiesConfig);
+
+        assertEquals(7, companiesRow.getField(0));
+        assertEquals("companies", companiesRow.getTableId());
+
+        Neo4jSourceTableConfig singleTableConfig =
+                new Neo4jSourceTableConfig("single query", peopleRowType, 
null);
+        SeaTunnelRow singleTableRow =
+                Neo4jSourceReader.convertRecord(peopleRecord, 
singleTableConfig);
+        assertEquals("", singleTableRow.getTableId());
+    }
+
+    @Test
+    void includesTableIdAndSignalsCompletionWhenMultiTableReadFails() throws 
Exception {
+        SingleSplitReaderContext context = 
mock(SingleSplitReaderContext.class);
+        Session session = mock(Session.class);
+        ServiceUnavailableException failure = new 
ServiceUnavailableException("connection refused");
+        
when(session.readTransaction(any(TransactionWork.class))).thenThrow(failure);
+        Neo4jSourceTableConfig tableConfig =
+                new Neo4jSourceTableConfig("MATCH (n) RETURN n", rowType(), 
"people");
+        Neo4jSourceReader reader = reader(context, session, tableConfig);
+        Collector<SeaTunnelRow> collector = mock(Collector.class);
+
+        reader.open();
+        Neo4jConnectorException thrown =
+                assertThrows(
+                        Neo4jConnectorException.class, () -> 
reader.internalPollNext(collector));
+
+        assertTrue(thrown.getMessage().contains("people"));
+        assertSame(failure, thrown.getCause());
+        verify(context).signalNoMoreElement();
+    }
+
+    @Test
+    void keepsOriginalFailureForSingleTableRead() throws Exception {
+        SingleSplitReaderContext context = 
mock(SingleSplitReaderContext.class);
+        Session session = mock(Session.class);
+        ServiceUnavailableException failure = new 
ServiceUnavailableException("connection refused");
+        
when(session.readTransaction(any(TransactionWork.class))).thenThrow(failure);
+        Neo4jSourceTableConfig tableConfig =
+                new Neo4jSourceTableConfig("MATCH (n) RETURN n", rowType(), 
null);
+        Neo4jSourceReader reader = reader(context, session, tableConfig);
+        Collector<SeaTunnelRow> collector = mock(Collector.class);
+
+        reader.open();
+        ServiceUnavailableException thrown =
+                assertThrows(
+                        ServiceUnavailableException.class,
+                        () -> reader.internalPollNext(collector));
+
+        assertSame(failure, thrown);
+        verify(context).signalNoMoreElement();
+    }
+
     @Test
     void convertType() {
         assertEquals(
@@ -110,4 +211,21 @@ class Neo4jSourceReaderTest {
                                 new MapType<>(BasicType.INT_TYPE, 
BasicType.BOOLEAN_TYPE),
                                 new MapValue(Collections.singletonMap("1", 
BooleanValue.FALSE))));
     }
+
+    private Neo4jSourceReader reader(
+            SingleSplitReaderContext context, Session session, 
Neo4jSourceTableConfig tableConfig) {
+        Driver driver = mock(Driver.class);
+        DriverBuilder driverBuilder = mock(DriverBuilder.class);
+        Neo4jSourceQueryInfo queryInfo = mock(Neo4jSourceQueryInfo.class);
+        when(driverBuilder.build()).thenReturn(driver);
+        when(driverBuilder.getDatabase()).thenReturn("neo4j");
+        when(driver.session(any(SessionConfig.class))).thenReturn(session);
+        when(queryInfo.getDriverBuilder()).thenReturn(driverBuilder);
+        return new Neo4jSourceReader(context, queryInfo, 
Collections.singletonList(tableConfig));
+    }
+
+    private SeaTunnelRowType rowType() {
+        return new SeaTunnelRowType(
+                new String[] {"name"}, new SeaTunnelDataType<?>[] 
{BasicType.STRING_TYPE});
+    }
 }
diff --git 
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
 
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
index a77813d702..4c6220a97f 100644
--- 
a/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
+++ 
b/seatunnel-connectors-v2/connector-neo4j/src/test/java/org/apache/seatunnel/connectors/seatunnel/neo4j/Neo4jFactoryTest.java
@@ -17,12 +17,24 @@
 
 package org.apache.seatunnel.connectors.seatunnel.neo4j;
 
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
+import org.apache.seatunnel.api.table.factory.TableSourceFactoryContext;
 import org.apache.seatunnel.connectors.seatunnel.neo4j.sink.Neo4jSinkFactory;
+import org.apache.seatunnel.connectors.seatunnel.neo4j.source.Neo4jSource;
 import 
org.apache.seatunnel.connectors.seatunnel.neo4j.source.Neo4jSourceFactory;
 
 import org.junit.jupiter.api.Assertions;
 import org.junit.jupiter.api.Test;
 
+import java.util.Arrays;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
 class Neo4jFactoryTest {
 
     @Test
@@ -30,4 +42,160 @@ class Neo4jFactoryTest {
         Assertions.assertNotNull((new Neo4jSourceFactory()).optionRule());
         Assertions.assertNotNull((new Neo4jSinkFactory()).optionRule());
     }
+
+    @Test
+    void sourceOptionRuleAcceptsSingleTableConfig() {
+        Map<String, Object> config = sourceConnectionConfig();
+        config.put("query", "MATCH (p:Person) RETURN p.name");
+        config.put("schema", schema("people"));
+
+        Assertions.assertDoesNotThrow(() -> validateSource(config));
+    }
+
+    @Test
+    void sourceOptionRuleAcceptsTablesConfigs() {
+        Map<String, Object> config = sourceConnectionConfig();
+        config.put(
+                "tables_configs",
+                Arrays.asList(
+                        tableConfig("people", "MATCH (p:Person) RETURN 
p.name"),
+                        tableConfig("companies", "MATCH (c:Company) RETURN 
c.name")));
+
+        Assertions.assertDoesNotThrow(() -> validateSource(config));
+    }
+
+    @Test
+    void sourceOptionRuleRejectsSingleAndMultiTableConfigTogether() {
+        Map<String, Object> config = sourceConnectionConfig();
+        config.put("query", "MATCH (p:Person) RETURN p.name");
+        config.put("schema", schema("people"));
+        config.put(
+                "tables_configs",
+                Collections.singletonList(
+                        tableConfig("companies", "MATCH (c:Company) RETURN 
c.name")));
+
+        Assertions.assertThrows(OptionValidationException.class, () -> 
validateSource(config));
+    }
+
+    @Test
+    void sourceOptionRuleRejectsMissingTableConfiguration() {
+        Assertions.assertThrows(
+                OptionValidationException.class, () -> 
validateSource(sourceConnectionConfig()));
+    }
+
+    @Test
+    void sourceOptionRuleRejectsRootSchemaWithTablesConfigs() {
+        Map<String, Object> config = sourceConnectionConfig();
+        config.put("schema", schema("people"));
+        config.put(
+                "tables_configs",
+                Collections.singletonList(
+                        tableConfig("companies", "MATCH (c:Company) RETURN 
c.name")));
+
+        Assertions.assertThrows(OptionValidationException.class, () -> 
validateSource(config));
+    }
+
+    @Test
+    void sourceOptionRuleRejectsInvalidTablesConfigs() {
+        Map<String, Object> empty = sourceConnectionConfig();
+        empty.put("tables_configs", Collections.emptyList());
+
+        Map<String, Object> missingQuery = sourceConnectionConfig();
+        Map<String, Object> tableWithoutQuery = new HashMap<>();
+        tableWithoutQuery.put("schema", schema("people"));
+        missingQuery.put("tables_configs", 
Collections.singletonList(tableWithoutQuery));
+
+        Map<String, Object> missingSchema = sourceConnectionConfig();
+        Map<String, Object> tableWithoutSchema = new HashMap<>();
+        tableWithoutSchema.put("query", "MATCH (p:Person) RETURN p.name");
+        missingSchema.put("tables_configs", 
Collections.singletonList(tableWithoutSchema));
+
+        Map<String, Object> blankTable = sourceConnectionConfig();
+        blankTable.put(
+                "tables_configs",
+                Collections.singletonList(tableConfig(" ", "MATCH (p:Person) 
RETURN p.name")));
+
+        Map<String, Object> duplicateTable = sourceConnectionConfig();
+        duplicateTable.put(
+                "tables_configs",
+                Arrays.asList(
+                        tableConfig("people", "MATCH (p:Person) RETURN 
p.name"),
+                        tableConfig("people", "MATCH (p:Person) RETURN 
p.name")));
+
+        Assertions.assertAll(
+                () ->
+                        Assertions.assertThrows(
+                                OptionValidationException.class, () -> 
validateSource(empty)),
+                () ->
+                        Assertions.assertThrows(
+                                OptionValidationException.class,
+                                () -> validateSource(missingQuery)),
+                () ->
+                        Assertions.assertThrows(
+                                OptionValidationException.class,
+                                () -> validateSource(missingSchema)),
+                () ->
+                        Assertions.assertThrows(
+                                OptionValidationException.class, () -> 
validateSource(blankTable)),
+                () ->
+                        Assertions.assertThrows(
+                                OptionValidationException.class,
+                                () -> validateSource(duplicateTable)));
+    }
+
+    @Test
+    void sourceFactoryProducesCatalogTableForEachQuery() {
+        Map<String, Object> config = sourceConnectionConfig();
+        config.put(
+                "tables_configs",
+                Arrays.asList(
+                        tableConfig("people", "MATCH (p:Person) RETURN 
p.name"),
+                        tableConfig("companies", "MATCH (c:Company) RETURN 
c.name")));
+        validateSource(config);
+
+        Object createdSource =
+                new Neo4jSourceFactory()
+                        .createSource(
+                                new TableSourceFactoryContext(
+                                        ReadonlyConfig.fromMap(config),
+                                        getClass().getClassLoader()))
+                        .createSource();
+        Neo4jSource source = (Neo4jSource) createdSource;
+
+        List<CatalogTable> tables = source.getProducedCatalogTables();
+        Assertions.assertEquals(2, tables.size());
+        Assertions.assertEquals("people", 
tables.get(0).getTableId().toTablePath().toString());
+        Assertions.assertEquals("companies", 
tables.get(1).getTableId().toTablePath().toString());
+    }
+
+    private static Map<String, Object> sourceConnectionConfig() {
+        Map<String, Object> config = new HashMap<>();
+        config.put("uri", "neo4j://localhost:7687");
+        config.put("database", "neo4j");
+        config.put("username", "neo4j");
+        config.put("password", "password");
+        return config;
+    }
+
+    private static Map<String, Object> tableConfig(String table, String query) 
{
+        Map<String, Object> tableConfig = new HashMap<>();
+        tableConfig.put("query", query);
+        tableConfig.put("schema", schema(table));
+        return tableConfig;
+    }
+
+    private static Map<String, Object> schema(String table) {
+        Map<String, Object> fields = new HashMap<>();
+        fields.put("name", "STRING");
+
+        Map<String, Object> schema = new HashMap<>();
+        schema.put("table", table);
+        schema.put("fields", fields);
+        return schema;
+    }
+
+    private static void validateSource(Map<String, Object> config) {
+        ConfigValidator.of(ReadonlyConfig.fromMap(config))
+                .validate(new Neo4jSourceFactory().optionRule());
+    }
 }
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/pom.xml 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/pom.xml
index 240e8a87d1..f86e4450ec 100644
--- a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/pom.xml
+++ b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/pom.xml
@@ -33,6 +33,12 @@
             <version>${project.version}</version>
             <scope>test</scope>
         </dependency>
+        <dependency>
+            <groupId>org.apache.seatunnel</groupId>
+            <artifactId>connector-assert</artifactId>
+            <version>${project.version}</version>
+            <scope>test</scope>
+        </dependency>
         <dependency>
             <groupId>org.apache.seatunnel</groupId>
             <artifactId>connector-console</artifactId>
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
index 350ed68f52..ad21d44672 100644
--- 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/java/org/apache/seatunnel/e2e/connector/neo4j/Neo4jIT.java
@@ -184,6 +184,19 @@ public class Neo4jIT extends TestSuiteBase implements 
TestResource {
         assertEquals(FAKE_ROW_NUM, cnt);
     }
 
+    @TestTemplate
+    public void testMultiTableSource(TestContainer container)
+            throws IOException, InterruptedException {
+        neo4jSession.run("MATCH (n) WHERE n:MultiPerson OR n:MultiCompany 
DELETE n");
+        neo4jSession.run("CREATE (:MultiPerson {name:'Alice'})");
+        neo4jSession.run("CREATE (:MultiCompany {name:'Acme'})");
+
+        Container.ExecResult execResult =
+                container.executeJob("/neo4j/neo4j_multi_table_source.conf");
+
+        Assertions.assertEquals(0, execResult.getExitCode());
+    }
+
     @AfterAll
     @Override
     public void tearDown() {
diff --git 
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/resources/neo4j/neo4j_multi_table_source.conf
 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/resources/neo4j/neo4j_multi_table_source.conf
new file mode 100644
index 0000000000..3620e005ff
--- /dev/null
+++ 
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-neo4j-e2e/src/test/resources/neo4j/neo4j_multi_table_source.conf
@@ -0,0 +1,115 @@
+#
+# Licensed to the Apache Software Foundation (ASF) under one or more
+# contributor license agreements.  See the NOTICE file distributed with
+# this work for additional information regarding copyright ownership.
+# The ASF licenses this file to You under the Apache License, Version 2.0
+# (the "License"); you may not use this file except in compliance with
+# the License.  You may obtain a copy of the License at
+#
+#    http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing, software
+# distributed under the License is distributed on an "AS IS" BASIS,
+# WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+# See the License for the specific language governing permissions and
+# limitations under the License.
+#
+
+env {
+  parallelism = 1
+  job.mode = "BATCH"
+
+  spark.app.name = "SeaTunnel"
+  spark.executor.instances = 1
+  spark.executor.cores = 1
+  spark.executor.memory = "1g"
+  spark.master = local
+}
+
+source {
+  Neo4j {
+    uri = "neo4j://neo4j-host:7687"
+    username = "neo4j"
+    password = "Test@12343"
+    database = "neo4j"
+
+    tables_configs = [
+      {
+        query = "MATCH (n:MultiPerson) RETURN n.name AS name"
+        schema {
+          table = "people"
+          fields {
+            name = STRING
+          }
+        }
+      },
+      {
+        query = "MATCH (n:MultiCompany) RETURN n.name AS name"
+        schema {
+          table = "companies"
+          fields {
+            name = STRING
+          }
+        }
+      }
+    ]
+  }
+}
+
+sink {
+  Assert {
+    rules = {
+      table-names = ["people", "companies"]
+      tables_configs = [
+        {
+          table_path = "people"
+          row_rules = [
+            {
+              rule_type = MAX_ROW
+              rule_value = 1
+            },
+            {
+              rule_type = MIN_ROW
+              rule_value = 1
+            }
+          ]
+          field_rules = [
+            {
+              field_name = name
+              field_type = string
+              field_value = [
+                {
+                  equals_to = "Alice"
+                }
+              ]
+            }
+          ]
+        },
+        {
+          table_path = "companies"
+          row_rules = [
+            {
+              rule_type = MAX_ROW
+              rule_value = 1
+            },
+            {
+              rule_type = MIN_ROW
+              rule_value = 1
+            }
+          ]
+          field_rules = [
+            {
+              field_name = name
+              field_type = string
+              field_value = [
+                {
+                  equals_to = "Acme"
+                }
+              ]
+            }
+          ]
+        }
+      ]
+    }
+  }
+}

Reply via email to