跳到主要内容
版本:3.0.0

Cassandra

Cassandra 源连接器

引擎支持​

Spark
Flink
SeaTunnel Zeta

描述​

以批处理方式从 Apache Cassandra 读取数据。

Cassandra source 支持两种读取方式:

  • 使用 cql 读取单张表。
  • 使用 tables_configs 读取多张表,每个条目里配置一个 cql。

连接器会根据 CQL 返回结果里的列名和数据类型生成下游数据结构,所以 CQL 应该返回下游真正需要的列。

支持的数据源信息​

数据源支持版本依赖
Cassandra通用下载

关键特性​

数据类型映射​

Cassandra 数据类型SeaTunnel 数据类型
asciiSTRING
varchar/textSTRING
varintSTRING
uuid/timeuuidSTRING
inetSTRING
tinyintBYTE
smallintSHORT
intINT
bigint/counterLONG
floatFLOAT
double/decimalDOUBLE
booleanBOOLEAN
timeTIME
dateDATE
timestampTIMESTAMP
blobARRAY\<BYTE>
listARRAY
setARRAY
mapMAP

Source 选项​

名称类型是否必填默认值描述
hostString是-Cassandra 集群地址,格式是 host:port,多个地址用逗号分隔。
keyspaceString是-Cassandra 会话使用的 keyspace。
cqlString否 *-读取单张表时使用的 CQL。
tables_configsList\<Map>否 *-多表读取配置,每个条目必须包含一个 cql。
usernameString否-Cassandra 用户名,需要和 password 一起配置。
passwordString否-Cassandra 密码,需要和 username 一起配置。
datacenterString否datacenter1Cassandra Java Driver 使用的本地数据中心名称。
consistency_levelString否LOCAL_ONE读取一致性级别,例如 LOCAL_ONE、ONE、QUORUM、LOCAL_QUORUM。
common-options否-Source 插件通用参数,例如 plugin_output。

* cql 与 tables_configs 二选一,必须提供其中之一。

host [string]​

Cassandra 的集群地址, 格式为 host:port , 允许指定多个 hosts . 例如 "cassandra1:9042,cassandra2:9042".

keyspace [string]​

Cassandra 的键空间.

cql [String]​

查询 CQL,用于读取单张表的数据。它和 tables_configs 互斥。

连接器会使用 CQL 返回结果里的元数据来生成输出结构。通常建议写成能返回真实表字段的查询,例如 select * from source_table 或 select id, name from source_table。

tables_configs [List\<Map>]​

多表读取配置,每个条目必须包含 cql 字段。它和根层级的 cql 互斥。

不要在 tables_configs 中重复配置同一张源表,连接器启动时会检查重复表名。

示例条目:

{
cql = "SELECT id, name FROM keyspace.table1"
}

username [string]​

Cassandra 用户的用户名.

password [string]​

Cassandra 用户的密码.

datacenter [String]​

Cassandra 数据中心, 默认为 datacenter1.

consistency_level [String]​

Cassandra 的读取一致性级别, 默认为 LOCAL_ONE.

common-options​

Source 插件通用参数,详情请参考 Source 常用选项。

注意事项​

  • username 和 password 是一组配置。集群开启认证时两个都要配;未开启认证时两个都可以不配。
  • datacenter 必须和 Cassandra 集群的本地数据中心名称一致。默认值是 datacenter1,这也是常见 Testcontainers 环境里的默认值。
  • cql 和 tables_configs 互斥。读取一个结果表时用 cql,需要让一个 source 读取多张 Cassandra 表时用 tables_configs。
  • 这是批处理 source。它读取当前查询结果后就会结束。
  • 一个 CQL 查询会作为一个 source split 读取。调大任务并行度不会自动把单张 Cassandra 表拆成多个扫描任务。
  • 连接器底层使用 Cassandra Java Driver,本文档列出的连接选项是连接器实际读取的全部设置;其他 DataStax Driver 选项沿用其内置默认值。

任务示例​

单表读取​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
Cassandra {
host = "localhost:9042"
username = "cassandra"
password = "cassandra"
datacenter = "datacenter1"
keyspace = "test"
cql = "SELECT * FROM test.source_table"
plugin_output = "source_table"
}
}

sink {
Console {}
}

多表读取​

env {
parallelism = 1
job.mode = "BATCH"
}

source {
Cassandra {
host = "localhost:9042"
username = "cassandra"
password = "cassandra"
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 = "localhost:9042"
username = "cassandra"
password = "cassandra"
datacenter = "datacenter1"
keyspace = "test"
table = "mt_sink_table"
}
}

提高读取一致性级别​

当读取结果必须满足配置的副本因子时,使用 consistency_level = "QUORUM",并配合 datacenter 让 Driver 连接到正确的本地协调节点:

source {
Cassandra {
host = "cassandra1:9042,cassandra2:9042"
username = "cassandra"
password = "cassandra"
datacenter = "datacenter1"
keyspace = "test"
consistency_level = "QUORUM"
cql = "SELECT id, name, score FROM test.accounts"
}
}

变更日志​

Change Log
ChangeCommitVersion
[Improve][Connector-V2] Migrate Cassandra validation to declarative OptionRule (#11964)https://github.com/apache/seatunnel/commit/cc548bb3f3.0.0
[Fix][Connector-V2] Fix Cassandra null timestamp conversion (#11068)https://github.com/apache/seatunnel/commit/9fdfc38143.0.0
[Feature][Connector-V2][Cassandra] Add multi-table source support via… (#10896)https://github.com/apache/seatunnel/commit/5d2ab3f313.0.0