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
 (


Reply via email to