跳到主要内容
版本:Next

Fluss

Fluss Source 连接器

支持引擎

Spark
Flink
SeaTunnel Zeta

主要特性

描述

Fluss Source 用于在批处理或流处理作业中,从已有的 Fluss 表读取数据。

连接器通过 Fluss 的 log scanner 读取数据,每个表 bucket 对应一个分片(split),因此读取并行度与 bucket 数量一致。每条记录的变更类型(INSERTUPDATE_BEFOREUPDATE_AFTERDELETE)会映射为对应的 SeaTunnel RowKind,因此日志表以仅追加(append-only)的 INSERT 形式读取,主键表则以其变更日志(changelog)的形式读取。

主键表

对于主键表,连接器只从最早可用的 log offset 开始读取表的 changelog不会先读取 KV 快照。因此它能捕获持续发生的变更(插入/更新/删除)并带上正确的 RowKind,但不保证对已存在数据的完整初始加载:任何因日志保留(retention)或压缩(compaction)而从 changelog 中被清除的记录都会缺失。目前尚不支持主键表的“快照 + 增量”完整同步。如果需要完整的当前状态,请优先使用日志表(append-only),log scanner 总能将其完整读取。

有界性由作业模式决定:

  • BATCH 模式下,Source 是有界的:每个 bucket 读取到作业启动时捕获的最新 log offset 后,分片结束。
  • STREAMING 模式下,Source 是无界的:会持续读取新的 log 记录。每个 bucket 的读取位置会保存在 checkpoint 状态中,作业可从中断处恢复。

使用 start_mode 选择每个 bucket 的起始读取位点:

  • earliest(默认):从最早可用的 offset 开始读取整个 log。
  • latest:只读取作业启动之后新追加的记录。

start_mode=latest 仅对流处理作业有意义。

运行作业前,Fluss database 和 table 必须已经存在。Source 不会自动创建 Fluss database 或 table。表结构会自动从 Fluss 集群读取,因此不需要配置 schema 选项。

限制

  • 仅支持单表。 每个 source 只读取一张表,通过 database + table 配置。不支持在一个 source 中读取多张表。
  • 不支持指定任意起始位点。 起始位置只能通过 start_modeearliestlatest)选择;不支持从指定的 log offset 开始读取。
  • 不支持分区表。 将 source 指向分区 Fluss 表会在作业启动时直接失败(fail-fast)并给出错误。请使用非分区表。
  • 主键表仅按 changelog 读取。 连接器读取表的 changelog,而非 KV 快照,因此不做完整初始加载(“快照 + 增量”)。详见上方主键表 caution。

依赖

<dependency>
<groupId>com.alibaba.fluss</groupId>
<artifactId>fluss-client</artifactId>
<version>0.7.0</version>
</dependency>

Source 选项

名称类型是否必填默认值描述
bootstrap.serversstring-Fluss coordinator 地址,例如 fluss-coordinator:9123
databasestring-要读取的 Fluss database。
tablestring-要读取的 Fluss table。
client.configmap-传递给 Fluss 连接的额外 Fluss 客户端选项。
start_modestringearliest每个 bucket 的起始读取位点:earliest(整个 log)或 latest(仅作业启动后新追加的记录)。latestBATCH 模式下会被拒绝。
poll.timeout.mslong10000单次 Fluss log scanner poll 的最大阻塞时间,单位毫秒。
common-options--Source 通用选项,详见 Source 通用选项

client.config

使用 client.config 传递额外的 Fluss 客户端配置。

client.config = {
request.timeout = "30s"
}

支持的配置项请参考 Fluss 客户端文档。

数据类型映射

Fluss 数据类型SeaTunnel 数据类型
BOOLEANBOOLEAN
TINYINTTINYINT
SMALLINTSMALLINT
INTINT
BIGINTBIGINT
FLOATFLOAT
DOUBLEDOUBLE
DECIMALDECIMAL
CHARSTRING
STRINGSTRING
BINARYBYTES
BYTESBYTES
DATEDATE
TIMETIME
TIMESTAMPTIMESTAMP
TIMESTAMP_LTZTIMESTAMP_TZ

任务示例

批处理读取

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

source {
Fluss {
bootstrap.servers = "fluss-coordinator:9123"
database = "fluss_db"
table = "fluss_table"
plugin_output = "fluss_source"
}
}

sink {
Console {
plugin_input = "fluss_source"
}
}

流处理读取

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

source {
Fluss {
bootstrap.servers = "fluss-coordinator:9123"
database = "fluss_db"
table = "fluss_table"
start_mode = "latest"
}
}

sink {
Console {
}
}

Changelog