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