This is an automated email from the ASF dual-hosted git repository.
dybyte 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 5d2ab3f317 [Feature][Connector-V2][Cassandra] Add multi-table source
support via… (#10896)
5d2ab3f317 is described below
commit 5d2ab3f317157be9f2780c5472eb32fb3ba219e4
Author: NixonWahome <[email protected]>
AuthorDate: Mon May 25 13:08:35 2026 +0300
[Feature][Connector-V2][Cassandra] Add multi-table source support via…
(#10896)
---
docs/en/connectors/source/Cassandra.md | 83 ++++++++---
docs/zh/connectors/source/Cassandra.md | 81 ++++++++---
.../cassandra/config/CassandraBaseOptions.java | 3 +-
.../cassandra/config/CassandraTableConfig.java} | 30 ++--
.../cassandra/source/CassandraSource.java | 151 ++++++++++++++-------
.../cassandra/source/CassandraSourceFactory.java | 4 +-
.../cassandra/source/CassandraSourceReader.java | 28 +++-
.../seatunnel/cassandra/CassandraFactoryTest.java | 62 +++++++++
.../seatunnel/cassandra/CassandraIT.java | 113 ++++++++++++++-
.../resources/cassandra_multitable_source.conf | 50 +++++++
.../src/test/resources/init/cassandra_init.conf | 32 +++++
11 files changed, 524 insertions(+), 113 deletions(-)
diff --git a/docs/en/connectors/source/Cassandra.md
b/docs/en/connectors/source/Cassandra.md
index 67216841ea..b45f502c0f 100644
--- a/docs/en/connectors/source/Cassandra.md
+++ b/docs/en/connectors/source/Cassandra.md
@@ -19,15 +19,18 @@ Read data from Apache Cassandra.
## Options
-| name | type | required | default value |
-|-------------------|--------|----------|---------------|
-| host | String | Yes | - |
-| keyspace | String | Yes | - |
-| cql | String | Yes | - |
-| username | String | No | - |
-| password | String | No | - |
-| datacenter | String | No | datacenter1 |
-| consistency_level | String | No | LOCAL_ONE |
+| name | type | required | default value |
+|-------------------|--------------------|----------|---------------|
+| host | String | Yes | - |
+| keyspace | String | Yes | - |
+| cql | String | No * | - |
+| tables_configs | List\<Map\> | No * | - |
+| username | String | No | - |
+| password | String | No | - |
+| datacenter | String | No | datacenter1 |
+| consistency_level | String | No | LOCAL_ONE |
+
+> \* Exactly one of `cql` or `tables_configs` must be provided.
### host [string]
@@ -40,7 +43,21 @@ The `Cassandra` keyspace.
### cql [String]
-The query cql used to search data though Cassandra session.
+The query CQL used to read data from Cassandra. Use this for single-table
reads.
+Mutually exclusive with `tables_configs`.
+
+### tables_configs [List\<Map\>]
+
+Multi-table read configuration. Each entry must contain a `cql` field with the
query for that table.
+Mutually exclusive with `cql`.
+
+Example entry:
+
+```
+{
+ cql = "SELECT id, name FROM keyspace.table1"
+}
+```
### username [string]
@@ -56,24 +73,48 @@ The `Cassandra` datacenter, default is `datacenter1`.
### consistency_level [String]
-The `Cassandra` write consistency level, default is `LOCAL_ONE`.
+The `Cassandra` read consistency level, default is `LOCAL_ONE`.
## Examples
+### Single-table mode
+
+```hocon
+source {
+ Cassandra {
+ host = "localhost:9042"
+ username = "cassandra"
+ password = "cassandra"
+ datacenter = "datacenter1"
+ keyspace = "test"
+ cql = "SELECT * FROM test.source_table"
+ plugin_output = "source_table"
+ }
+}
+```
+
+### Multi-table mode
+
```hocon
source {
- Cassandra {
- host = "localhost:9042"
- username = "cassandra"
- password = "cassandra"
- datacenter = "datacenter1"
- keyspace = "test"
- cql = "select * from source_table"
- plugin_output = "source_table"
- }
+ Cassandra {
+ host = "localhost:9042"
+ username = "cassandra"
+ password = "cassandra"
+ datacenter = "datacenter1"
+ keyspace = "test"
+ tables_configs = [
+ {
+ cql = "SELECT id, name FROM test.table1"
+ },
+ {
+ cql = "SELECT id, value FROM test.table2"
+ }
+ ]
+ }
}
```
## Changelog
-<ChangeLog />
\ No newline at end of file
+<ChangeLog />
diff --git a/docs/zh/connectors/source/Cassandra.md
b/docs/zh/connectors/source/Cassandra.md
index d1dc69918d..17a3d36a92 100644
--- a/docs/zh/connectors/source/Cassandra.md
+++ b/docs/zh/connectors/source/Cassandra.md
@@ -19,15 +19,18 @@ import ChangeLog from '../changelog/connector-cassandra.md';
## 选项
-| 名称 | 类型 | 必需 | 默认值 |
-|-------------------|--------|----|---------------|
-| host | String | 是 | - |
-| keyspace | String | 是 | - |
-| cql | String | 是 | - |
-| username | String | 否 | - |
-| password | String | 否 | - |
-| datacenter | String | 否 | datacenter1 |
-| consistency_level | String | 否 | LOCAL_ONE |
+| 名称 | 类型 | 必需 | 默认值 |
+|-------------------|-------------|------|---------------|
+| host | String | 是 | - |
+| keyspace | String | 是 | - |
+| cql | String | 否 * | - |
+| tables_configs | List\<Map\> | 否 * | - |
+| username | String | 否 | - |
+| password | String | 否 | - |
+| datacenter | String | 否 | datacenter1 |
+| consistency_level | String | 否 | LOCAL_ONE |
+
+> \* `cql` 与 `tables_configs` 二选一,必须提供其中之一。
### host [string]
@@ -40,7 +43,19 @@ import ChangeLog from '../changelog/connector-cassandra.md';
### cql [String]
-查询cql,用于通过Cassandra会话搜索数据.
+查询 CQL,用于通过 Cassandra 会话读取单张表的数据。与 `tables_configs` 互斥。
+
+### tables_configs [List\<Map\>]
+
+多表读取配置,每个条目必须包含 `cql` 字段。与 `cql` 互斥。
+
+示例条目:
+
+```
+{
+ cql = "SELECT id, name FROM keyspace.table1"
+}
+```
### username [string]
@@ -56,24 +71,48 @@ import ChangeLog from '../changelog/connector-cassandra.md';
### consistency_level [String]
-`Cassandra` 的写入一致性级别, 默认为 `LOCAL_ONE`.
+`Cassandra` 的读取一致性级别, 默认为 `LOCAL_ONE`.
## 示例
+### 单表模式
+
+```hocon
+source {
+ Cassandra {
+ host = "localhost:9042"
+ username = "cassandra"
+ password = "cassandra"
+ datacenter = "datacenter1"
+ keyspace = "test"
+ cql = "SELECT * FROM test.source_table"
+ plugin_output = "source_table"
+ }
+}
+```
+
+### 多表模式
+
```hocon
source {
- Cassandra {
- host = "localhost:9042"
- username = "cassandra"
- password = "cassandra"
- datacenter = "datacenter1"
- keyspace = "test"
- cql = "select * from source_table"
- plugin_output = "source_table"
- }
+ Cassandra {
+ host = "localhost:9042"
+ username = "cassandra"
+ password = "cassandra"
+ datacenter = "datacenter1"
+ keyspace = "test"
+ tables_configs = [
+ {
+ cql = "SELECT id, name FROM test.table1"
+ },
+ {
+ cql = "SELECT id, value FROM test.table2"
+ }
+ ]
+ }
}
```
## 变更日志
-<ChangeLog />
\ No newline at end of file
+<ChangeLog />
diff --git
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/config/CassandraBaseOptions.java
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/config/CassandraBaseOptions.java
index 7addd8e876..372dd8abcb 100644
---
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/config/CassandraBaseOptions.java
+++
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/config/CassandraBaseOptions.java
@@ -19,8 +19,9 @@ package
org.apache.seatunnel.connectors.seatunnel.cassandra.config;
import org.apache.seatunnel.api.configuration.Option;
import org.apache.seatunnel.api.configuration.Options;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
-public class CassandraBaseOptions {
+public class CassandraBaseOptions extends ConnectorCommonOptions {
public static final Integer DEFAULT_BATCH_SIZE = 5000;
diff --git
a/seatunnel-connectors-v2/connector-cassandra/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraFactoryTest.java
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/config/CassandraTableConfig.java
similarity index 53%
copy from
seatunnel-connectors-v2/connector-cassandra/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraFactoryTest.java
copy to
seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/config/CassandraTableConfig.java
index 3a07a7edb7..796e140b8f 100644
---
a/seatunnel-connectors-v2/connector-cassandra/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraFactoryTest.java
+++
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/config/CassandraTableConfig.java
@@ -15,19 +15,29 @@
* limitations under the License.
*/
-package org.apache.seatunnel.connectors.seatunnel.cassandra;
+package org.apache.seatunnel.connectors.seatunnel.cassandra.config;
-import
org.apache.seatunnel.connectors.seatunnel.cassandra.sink.CassandraSinkFactory;
-import
org.apache.seatunnel.connectors.seatunnel.cassandra.source.CassandraSourceFactory;
+import org.apache.seatunnel.api.table.catalog.CatalogTable;
-import org.junit.jupiter.api.Assertions;
-import org.junit.jupiter.api.Test;
+import lombok.Getter;
-class CassandraFactoryTest {
+import java.io.Serializable;
- @Test
- void optionRule() {
- Assertions.assertNotNull((new CassandraSourceFactory()).optionRule());
- Assertions.assertNotNull((new CassandraSinkFactory()).optionRule());
+@Getter
+public class CassandraTableConfig implements Serializable {
+
+ private static final long serialVersionUID = 1L;
+
+ private final String cql;
+
+ private final CatalogTable catalogTable;
+
+ /** Pre-computed routing key used to set {@code SeaTunnelRow#tableId}. */
+ private final String tableId;
+
+ public CassandraTableConfig(String cql, CatalogTable catalogTable) {
+ this.cql = cql;
+ this.catalogTable = catalogTable;
+ this.tableId = catalogTable.getTableId().toTablePath().toString();
}
}
diff --git
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSource.java
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSource.java
index 46e19b5a37..d92f8b810a 100644
---
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSource.java
+++
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSource.java
@@ -18,6 +18,7 @@
package org.apache.seatunnel.connectors.seatunnel.cassandra.source;
import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
import org.apache.seatunnel.api.source.Boundedness;
import org.apache.seatunnel.api.source.SupportColumnProjection;
import org.apache.seatunnel.api.table.catalog.CatalogTable;
@@ -28,7 +29,8 @@ import org.apache.seatunnel.api.table.type.SeaTunnelRow;
import org.apache.seatunnel.common.exception.CommonErrorCodeDeprecated;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.client.CassandraClient;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.config.CassandraParameters;
-import
org.apache.seatunnel.connectors.seatunnel.cassandra.exception.CassandraConnectorErrorCode;
+import
org.apache.seatunnel.connectors.seatunnel.cassandra.config.CassandraSourceOptions;
+import
org.apache.seatunnel.connectors.seatunnel.cassandra.config.CassandraTableConfig;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.exception.CassandraConnectorException;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.util.TypeConvertUtil;
import
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
@@ -36,77 +38,122 @@ import
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSpl
import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
import com.datastax.oss.driver.api.core.CqlSession;
-import com.datastax.oss.driver.api.core.cql.Row;
+import com.datastax.oss.driver.api.core.cql.ColumnDefinitions;
+import java.util.ArrayList;
import java.util.Collections;
import java.util.List;
-
-import static
org.apache.seatunnel.connectors.seatunnel.cassandra.config.CassandraSourceOptions.CQL;
+import java.util.Map;
+import java.util.stream.Collectors;
public class CassandraSource extends AbstractSingleSplitSource<SeaTunnelRow>
implements SupportColumnProjection {
+ private static final String PLUGIN_NAME = "Cassandra";
+
private final CassandraParameters cassandraParameters;
- private final CatalogTable catalogTable;
+ private final List<CassandraTableConfig> tableConfigs;
- public CassandraSource(CassandraParameters cassandraParameters,
ReadonlyConfig pluginConfig) {
+ public CassandraSource(CassandraParameters cassandraParameters,
ReadonlyConfig config) {
this.cassandraParameters = cassandraParameters;
+ this.tableConfigs = buildTableConfigs(cassandraParameters, config);
+ }
- try (CqlSession currentSession =
+ private List<CassandraTableConfig> buildTableConfigs(
+ CassandraParameters params, ReadonlyConfig config) {
+ try (CqlSession session =
CassandraClient.getCqlSessionBuilder(
- cassandraParameters.getHost(),
- cassandraParameters.getKeyspace(),
- cassandraParameters.getUsername(),
- cassandraParameters.getPassword(),
- cassandraParameters.getDatacenter())
+ params.getHost(),
+ params.getKeyspace(),
+ params.getUsername(),
+ params.getPassword(),
+ params.getDatacenter())
.build()) {
- Row rs =
- currentSession
- .execute(
- CassandraClient.createSimpleStatement(
- pluginConfig.get(CQL),
-
cassandraParameters.getConsistencyLevel()))
- .one();
- if (rs == null) {
- throw new CassandraConnectorException(
- CassandraConnectorErrorCode.NO_DATA_IN_SOURCE_TABLE,
- "No data select from this cql: " +
pluginConfig.get(CQL));
- }
- int columnSize = rs.getColumnDefinitions().size();
- TableSchema.Builder schemaBuilder = TableSchema.builder();
- String tableName = "default";
- for (int i = 0; i < columnSize; i++) {
- PhysicalColumn physicalColumn =
- PhysicalColumn.of(
-
rs.getColumnDefinitions().get(i).getName().asInternal(),
-
TypeConvertUtil.convert(rs.getColumnDefinitions().get(i).getType()),
- null,
- null,
- true,
- null,
- null);
- schemaBuilder.column(physicalColumn);
- tableName =
rs.getColumnDefinitions().get(i).getTable().asInternal();
+
+ if
(config.getOptional(ConnectorCommonOptions.TABLE_CONFIGS).isPresent()) {
+ List<Map<String, Object>> tableConfigMaps =
+ config.get(ConnectorCommonOptions.TABLE_CONFIGS);
+ List<CassandraTableConfig> configs = new ArrayList<>();
+ for (Map<String, Object> tableConfigMap : tableConfigMaps) {
+ ReadonlyConfig tableConfig =
ReadonlyConfig.fromMap(tableConfigMap);
+ String cql = tableConfig.get(CassandraSourceOptions.CQL);
+ CassandraTableConfig built =
+ buildTableConfig(
+ cql,
+ session,
+ params.getKeyspace(),
+ params.getConsistencyLevel());
+ String tableId = built.getTableId();
+ boolean duplicate =
+ configs.stream().anyMatch(c ->
c.getTableId().equals(tableId));
+ if (duplicate) {
+ throw new CassandraConnectorException(
+ CommonErrorCodeDeprecated.ILLEGAL_ARGUMENT,
+ "Duplicate table found in tables_configs: " +
tableId);
+ }
+ configs.add(built);
+ }
+ return configs;
+ } else {
+ String cql = config.get(CassandraSourceOptions.CQL);
+ return Collections.singletonList(
+ buildTableConfig(
+ cql, session, params.getKeyspace(),
params.getConsistencyLevel()));
}
- catalogTable =
- CatalogTable.of(
- TableIdentifier.of(
- getPluginName(),
cassandraParameters.getKeyspace(), tableName),
- schemaBuilder.build(),
- Collections.emptyMap(),
- Collections.emptyList(),
- "");
+ } catch (CassandraConnectorException e) {
+ throw e;
} catch (Exception e) {
throw new CassandraConnectorException(
CommonErrorCodeDeprecated.TABLE_SCHEMA_GET_FAILED,
- "Get table schema from cassandra source data failed",
+ "Get table schema from Cassandra source failed",
e);
}
}
+ private CassandraTableConfig buildTableConfig(
+ String cql,
+ CqlSession session,
+ String keyspace,
+ com.datastax.oss.driver.api.core.ConsistencyLevel
consistencyLevel) {
+ ColumnDefinitions columnDefs =
+ session.execute(CassandraClient.createSimpleStatement(cql,
consistencyLevel))
+ .getColumnDefinitions();
+
+ if (columnDefs.size() == 0) {
+ throw new CassandraConnectorException(
+ CommonErrorCodeDeprecated.TABLE_SCHEMA_GET_FAILED,
+ "No columns returned by CQL: " + cql);
+ }
+
+ TableSchema.Builder schemaBuilder = TableSchema.builder();
+ String tableName = "default";
+ for (int i = 0; i < columnDefs.size(); i++) {
+ schemaBuilder.column(
+ PhysicalColumn.of(
+ columnDefs.get(i).getName().asInternal(),
+
TypeConvertUtil.convert(columnDefs.get(i).getType()),
+ null,
+ null,
+ true,
+ null,
+ null));
+ tableName = columnDefs.get(i).getTable().asInternal();
+ }
+
+ CatalogTable catalogTable =
+ CatalogTable.of(
+ TableIdentifier.of(PLUGIN_NAME, keyspace, tableName),
+ schemaBuilder.build(),
+ Collections.emptyMap(),
+ Collections.emptyList(),
+ "");
+
+ return new CassandraTableConfig(cql, catalogTable);
+ }
+
@Override
public String getPluginName() {
- return "Cassandra";
+ return PLUGIN_NAME;
}
@Override
@@ -116,12 +163,14 @@ public class CassandraSource extends
AbstractSingleSplitSource<SeaTunnelRow>
@Override
public List<CatalogTable> getProducedCatalogTables() {
- return Collections.singletonList(catalogTable);
+ return tableConfigs.stream()
+ .map(CassandraTableConfig::getCatalogTable)
+ .collect(Collectors.toList());
}
@Override
public AbstractSingleSplitReader<SeaTunnelRow> createReader(
SingleSplitReaderContext readerContext) throws Exception {
- return new CassandraSourceReader(cassandraParameters, readerContext);
+ return new CassandraSourceReader(cassandraParameters, tableConfigs,
readerContext);
}
}
diff --git
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceFactory.java
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceFactory.java
index 7160d19e1e..f71f2dc43f 100644
---
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceFactory.java
+++
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceFactory.java
@@ -18,6 +18,7 @@
package org.apache.seatunnel.connectors.seatunnel.cassandra.source;
import org.apache.seatunnel.api.configuration.util.OptionRule;
+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.connector.TableSource;
@@ -48,7 +49,8 @@ public class CassandraSourceFactory implements
TableSourceFactory {
@Override
public OptionRule optionRule() {
return OptionRule.builder()
- .required(HOST, KEYSPACE, CQL)
+ .required(HOST, KEYSPACE)
+ .exclusive(CQL, ConnectorCommonOptions.TABLE_CONFIGS)
.bundled(USERNAME, PASSWORD)
.optional(DATACENTER, CONSISTENCY_LEVEL)
.build();
diff --git
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceReader.java
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceReader.java
index 6e4917ca9f..f367501954 100644
---
a/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceReader.java
+++
b/seatunnel-connectors-v2/connector-cassandra/src/main/java/org/apache/seatunnel/connectors/seatunnel/cassandra/source/CassandraSourceReader.java
@@ -21,6 +21,7 @@ import org.apache.seatunnel.api.source.Collector;
import org.apache.seatunnel.api.table.type.SeaTunnelRow;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.client.CassandraClient;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.config.CassandraParameters;
+import
org.apache.seatunnel.connectors.seatunnel.cassandra.config.CassandraTableConfig;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.util.TypeConvertUtil;
import
org.apache.seatunnel.connectors.seatunnel.common.source.AbstractSingleSplitReader;
import
org.apache.seatunnel.connectors.seatunnel.common.source.SingleSplitReaderContext;
@@ -30,16 +31,21 @@ import com.datastax.oss.driver.api.core.cql.ResultSet;
import lombok.extern.slf4j.Slf4j;
import java.io.IOException;
+import java.util.List;
@Slf4j
public class CassandraSourceReader extends
AbstractSingleSplitReader<SeaTunnelRow> {
private final CassandraParameters cassandraParameters;
+ private final List<CassandraTableConfig> tableConfigs;
private final SingleSplitReaderContext readerContext;
private CqlSession session;
CassandraSourceReader(
- CassandraParameters cassandraParameters, SingleSplitReaderContext
readerContext) {
+ CassandraParameters cassandraParameters,
+ List<CassandraTableConfig> tableConfigs,
+ SingleSplitReaderContext readerContext) {
this.cassandraParameters = cassandraParameters;
+ this.tableConfigs = tableConfigs;
this.readerContext = readerContext;
}
@@ -65,12 +71,20 @@ public class CassandraSourceReader extends
AbstractSingleSplitReader<SeaTunnelRo
@Override
public void pollNext(Collector<SeaTunnelRow> output) throws Exception {
try {
- ResultSet resultSet =
- session.execute(
- CassandraClient.createSimpleStatement(
- cassandraParameters.getCql(),
-
cassandraParameters.getConsistencyLevel()));
- resultSet.forEach(row ->
output.collect(TypeConvertUtil.buildSeaTunnelRow(row)));
+ for (CassandraTableConfig tableConfig : tableConfigs) {
+ ResultSet resultSet =
+ session.execute(
+ CassandraClient.createSimpleStatement(
+ tableConfig.getCql(),
+
cassandraParameters.getConsistencyLevel()));
+ String tableId = tableConfig.getTableId();
+ resultSet.forEach(
+ row -> {
+ SeaTunnelRow seaTunnelRow =
TypeConvertUtil.buildSeaTunnelRow(row);
+ seaTunnelRow.setTableId(tableId);
+ output.collect(seaTunnelRow);
+ });
+ }
} finally {
this.readerContext.signalNoMoreElement();
}
diff --git
a/seatunnel-connectors-v2/connector-cassandra/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraFactoryTest.java
b/seatunnel-connectors-v2/connector-cassandra/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraFactoryTest.java
index 3a07a7edb7..9ca2acc2ef 100644
---
a/seatunnel-connectors-v2/connector-cassandra/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraFactoryTest.java
+++
b/seatunnel-connectors-v2/connector-cassandra/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraFactoryTest.java
@@ -17,12 +17,23 @@
package org.apache.seatunnel.connectors.seatunnel.cassandra;
+import org.apache.seatunnel.api.configuration.ReadonlyConfig;
+import org.apache.seatunnel.api.configuration.util.ConfigValidator;
+import org.apache.seatunnel.api.configuration.util.OptionRule;
+import org.apache.seatunnel.api.configuration.util.OptionValidationException;
+import org.apache.seatunnel.api.options.ConnectorCommonOptions;
+import
org.apache.seatunnel.connectors.seatunnel.cassandra.config.CassandraSourceOptions;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.sink.CassandraSinkFactory;
import
org.apache.seatunnel.connectors.seatunnel.cassandra.source.CassandraSourceFactory;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import java.util.Collections;
+import java.util.HashMap;
+import java.util.List;
+import java.util.Map;
+
class CassandraFactoryTest {
@Test
@@ -30,4 +41,55 @@ class CassandraFactoryTest {
Assertions.assertNotNull((new CassandraSourceFactory()).optionRule());
Assertions.assertNotNull((new CassandraSinkFactory()).optionRule());
}
+
+ @Test
+ void testSourceOptionRuleWithCqlOnly() {
+ OptionRule rule = new CassandraSourceFactory().optionRule();
+ Map<String, Object> cfg = baseConfig();
+ cfg.put(CassandraSourceOptions.CQL.key(), "select * from test.table1");
+ ConfigValidator.of(ReadonlyConfig.fromMap(cfg)).validate(rule);
+ }
+
+ @Test
+ void testSourceOptionRuleWithTablesConfigsOnly() {
+ OptionRule rule = new CassandraSourceFactory().optionRule();
+ Map<String, Object> cfg = baseConfig();
+ List<Map<String, Object>> tablesConfigs =
+ Collections.singletonList(
+ Collections.singletonMap(
+ CassandraSourceOptions.CQL.key(), "select *
from test.table1"));
+ cfg.put(ConnectorCommonOptions.TABLE_CONFIGS.key(), tablesConfigs);
+ ConfigValidator.of(ReadonlyConfig.fromMap(cfg)).validate(rule);
+ }
+
+ @Test
+ void testSourceOptionRuleWithBothCqlAndTablesConfigsThrows() {
+ OptionRule rule = new CassandraSourceFactory().optionRule();
+ Map<String, Object> cfg = baseConfig();
+ cfg.put(CassandraSourceOptions.CQL.key(), "select * from test.table1");
+ List<Map<String, Object>> tablesConfigs =
+ Collections.singletonList(
+ Collections.singletonMap(
+ CassandraSourceOptions.CQL.key(), "select *
from test.table2"));
+ cfg.put(ConnectorCommonOptions.TABLE_CONFIGS.key(), tablesConfigs);
+ Assertions.assertThrows(
+ OptionValidationException.class,
+ () ->
ConfigValidator.of(ReadonlyConfig.fromMap(cfg)).validate(rule));
+ }
+
+ @Test
+ void testSourceOptionRuleWithNeitherCqlNorTablesConfigsThrows() {
+ OptionRule rule = new CassandraSourceFactory().optionRule();
+ Map<String, Object> cfg = baseConfig();
+ Assertions.assertThrows(
+ OptionValidationException.class,
+ () ->
ConfigValidator.of(ReadonlyConfig.fromMap(cfg)).validate(rule));
+ }
+
+ private Map<String, Object> baseConfig() {
+ Map<String, Object> cfg = new HashMap<>();
+ cfg.put("host", "localhost:9042");
+ cfg.put("keyspace", "test");
+ return cfg;
+ }
}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraIT.java
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraIT.java
index e86f5b21b9..55c8228a32 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraIT.java
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/java/org/apache/seatunnel/connectors/seatunnel/cassandra/CassandraIT.java
@@ -92,12 +92,20 @@ public class CassandraIT extends TestSuiteBase implements
TestResource {
private static final Integer PORT = 9042;
private static final String INIT_CASSANDRA_PATH =
"/init/cassandra_init.conf";
private static final String CASSANDRA_JOB_CONFIG =
"/cassandra_to_cassandra.conf";
+ private static final String CASSANDRA_MULTI_TABLE_JOB_CONFIG =
+ "/cassandra_multitable_source.conf";
private static final String CASSANDRA_DRIVER_CONFIG = "/application.conf";
private static final String DATACENTER = "datacenter1";
private static final String KEYSPACE = "test";
private static final String SOURCE_TABLE = "source_table";
private static final String SINK_TABLE = "sink_table";
private static final String INSERT_CQL = "insert_cql";
+ private static final String MT_SOURCE_TABLE_A = "mt_source_a";
+ private static final String MT_SOURCE_TABLE_B = "mt_source_b";
+ private static final String MT_SINK_TABLE = "mt_sink_table";
+ private static final String MT_INSERT_A = "mt_insert_a";
+ private static final String MT_INSERT_B = "mt_insert_b";
+ private static final int MT_ROWS_PER_TABLE = 5;
private static final Pair<SeaTunnelRowType, List<SeaTunnelRow>>
TEST_DATASET =
generateTestDataSet();
private Config config;
@@ -114,6 +122,24 @@ public class CassandraIT extends TestSuiteBase implements
TestResource {
Assertions.assertNull(getRow());
}
+ @TestTemplate
+ public void testCassandraMultiTableSource(TestContainer container) throws
Exception {
+ Container.ExecResult execResult =
container.executeJob(CASSANDRA_MULTI_TABLE_JOB_CONFIG);
+ Assertions.assertEquals(0, execResult.getExitCode());
+ long sinkCount = getMTSinkRowCount();
+ Assertions.assertEquals(MT_ROWS_PER_TABLE * 2L, sinkCount);
+ Set<Long> sinkIds = getMTSinkIds();
+ for (long i = 0; i < MT_ROWS_PER_TABLE; i++) {
+ Assertions.assertTrue(
+ sinkIds.contains(i), "Expected id " + i + " from
mt_source_a in sink");
+ }
+ for (long i = MT_ROWS_PER_TABLE; i < MT_ROWS_PER_TABLE * 2L; i++) {
+ Assertions.assertTrue(
+ sinkIds.contains(i), "Expected id " + i + " from
mt_source_b in sink");
+ }
+ clearMTSinkTable();
+ }
+
@BeforeAll
@Override
public void startUp() throws Exception {
@@ -134,6 +160,8 @@ public class CassandraIT extends TestSuiteBase implements
TestResource {
.untilAsserted(this::initConnection);
this.initializeCassandraTable();
this.batchInsertData();
+ this.initializeMTTables();
+ this.batchInsertMTData();
}
private void initializeCassandraTable() {
@@ -388,12 +416,95 @@ public class CassandraIT extends TestSuiteBase implements
TestResource {
}
}
+ private void initializeMTTables() {
+ try {
+ session.execute(
+
SimpleStatement.builder(config.getString(MT_SOURCE_TABLE_A))
+ .setKeyspace(KEYSPACE)
+ .setTimeout(Duration.ofSeconds(10))
+ .build());
+ session.execute(
+
SimpleStatement.builder(config.getString(MT_SOURCE_TABLE_B))
+ .setKeyspace(KEYSPACE)
+ .setTimeout(Duration.ofSeconds(10))
+ .build());
+ session.execute(
+ SimpleStatement.builder(config.getString(MT_SINK_TABLE))
+ .setKeyspace(KEYSPACE)
+ .setTimeout(Duration.ofSeconds(10))
+ .build());
+ } catch (Exception e) {
+ throw new RuntimeException("Initializing MT tables failed!", e);
+ }
+ }
+
+ private void batchInsertMTData() {
+ try {
+ BoundStatement stmtA =
+ session.prepare(
+
SimpleStatement.builder(config.getString(MT_INSERT_A))
+ .setKeyspace(KEYSPACE)
+ .build())
+ .bind();
+ BatchStatement batchA =
BatchStatement.builder(BatchType.UNLOGGED).build();
+ for (int i = 0; i < MT_ROWS_PER_TABLE; i++) {
+ batchA = batchA.add(stmtA.setLong(0, i).setInt(1, i));
+ }
+ session.execute(batchA);
+
+ BoundStatement stmtB =
+ session.prepare(
+
SimpleStatement.builder(config.getString(MT_INSERT_B))
+ .setKeyspace(KEYSPACE)
+ .build())
+ .bind();
+ BatchStatement batchB =
BatchStatement.builder(BatchType.UNLOGGED).build();
+ for (int i = MT_ROWS_PER_TABLE; i < MT_ROWS_PER_TABLE * 2; i++) {
+ batchB = batchB.add(stmtB.setLong(0, i).setInt(1, i));
+ }
+ session.execute(batchB);
+ } catch (Exception e) {
+ throw new RuntimeException("Batch insert MT data failed!", e);
+ }
+ }
+
+ private long getMTSinkRowCount() {
+ ResultSet rs =
+ session.execute(
+ SimpleStatement.builder("select count(*) from " +
MT_SINK_TABLE)
+ .setKeyspace(KEYSPACE)
+ .build());
+ Row row = rs.one();
+ return row != null ? row.getLong(0) : 0;
+ }
+
+ private Set<Long> getMTSinkIds() {
+ ResultSet rs =
+ session.execute(
+ SimpleStatement.builder("select id from " +
MT_SINK_TABLE)
+ .setKeyspace(KEYSPACE)
+ .build());
+ return rs.all().stream().map(row ->
row.getLong(0)).collect(Collectors.toSet());
+ }
+
+ private void clearMTSinkTable() {
+ session.execute(
+ SimpleStatement.builder(String.format("truncate table %s",
MT_SINK_TABLE))
+ .setKeyspace(KEYSPACE)
+ .build());
+ }
+
private void initCassandraConfig() {
File file = ContainerUtil.getResourcesFile(INIT_CASSANDRA_PATH);
Config config = ConfigFactory.parseFile(file);
assert config.hasPath(SOURCE_TABLE)
&& config.hasPath(SINK_TABLE)
- && config.hasPath(INSERT_CQL);
+ && config.hasPath(INSERT_CQL)
+ && config.hasPath(MT_SOURCE_TABLE_A)
+ && config.hasPath(MT_SOURCE_TABLE_B)
+ && config.hasPath(MT_SINK_TABLE)
+ && config.hasPath(MT_INSERT_A)
+ && config.hasPath(MT_INSERT_B);
this.config = config;
}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/resources/cassandra_multitable_source.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/resources/cassandra_multitable_source.conf
new file mode 100644
index 0000000000..a5ef5e0d1b
--- /dev/null
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/resources/cassandra_multitable_source.conf
@@ -0,0 +1,50 @@
+#
+# 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"
+}
+
+source {
+ Cassandra {
+ host = "cassandra:9042"
+ username = ""
+ password = ""
+ datacenter = "datacenter1"
+ keyspace = "test"
+ tables_configs = [
+ {
+ cql = "select id, c_int from mt_source_a"
+ },
+ {
+ cql = "select id, c_int from mt_source_b"
+ }
+ ]
+ }
+}
+
+sink {
+ Cassandra {
+ host = "cassandra:9042"
+ username = ""
+ password = ""
+ datacenter = "datacenter1"
+ keyspace = "test"
+ table = "mt_sink_table"
+ }
+}
diff --git
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/resources/init/cassandra_init.conf
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/resources/init/cassandra_init.conf
index 62b4952440..219412e324 100644
---
a/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/resources/init/cassandra_init.conf
+++
b/seatunnel-e2e/seatunnel-connector-v2-e2e/connector-cassandra-e2e/src/test/resources/init/cassandra_init.conf
@@ -73,6 +73,38 @@ create table if not exists sink_table(
);
"""
+mt_source_a = """
+create table if not exists mt_source_a(
+ id bigint,
+ c_int int,
+ PRIMARY KEY (id)
+);
+"""
+
+mt_source_b = """
+create table if not exists mt_source_b(
+ id bigint,
+ c_int int,
+ PRIMARY KEY (id)
+);
+"""
+
+mt_sink_table = """
+create table if not exists mt_sink_table(
+ id bigint,
+ c_int int,
+ PRIMARY KEY (id)
+);
+"""
+
+mt_insert_a = """
+insert into mt_source_a (id, c_int) values (?, ?)
+"""
+
+mt_insert_b = """
+insert into mt_source_b (id, c_int) values (?, ?)
+"""
+
insert_cql = """
insert into source_table
(