使用 Flink CDC Fluss Pipeline 连接器:将 MySQL 实时数据写入 Fluss(含自动建表、分桶策略与 Schema 变更同步)
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
Fluss Pipeline 连接器是 Flink CDC 提供的 PipelineData Sink实现,它让 YAML 定义的整库同步任务能够把 MySQL 等源端的数据变更直接写入 Fluss(一个支持主键表与日志表的流式存储系统)。读完本文,你将掌握完整的 Fluss Sink YAML 配置、bucket.key/bucket.num的分桶与数据分布策略、LENIENT模式下 Schema 变更同步的具体语义,以及 Flink CDC 类型到 Fluss 类型的映射规则,并借助仓库源码理解自动建表、分桶哈希与元数据应用等底层实现。
连接器能做什么
Fluss Pipeline 连接器作为 pipeline 的 Data Sink 使用,主要提供三项能力:
- 自动创建不存在的表:当源端表在 Fluss 中不存在时,连接器会根据源表 Schema 自动建表(含自动建库),无需预先在 Fluss 中手工创建。
- 数据同步:通过 Fluss Java Client 将 CDC 事件(Insert/Update/Delete)写入 Fluss 的对应表。
- Schema 变更同步(lenient 模式):在
schema.change.behavior: lenient配置下,源端的部分 Schema 变更(如新增列)可以传递并应用到 Fluss 表。
如何创建 Pipeline
从 MySQL 读取数据并写入 Fluss 的 Pipeline 可定义如下(原文示例,补充注释说明):
source: type: mysql name: MySQL Source hostname: 127.0.0.1 port: 3306 username: admin password: pass tables: adb.\.*, bdb.user_table_[0-9]+, [app|web].order_\.* server-id: 5401-5404 sink: type: fluss name: Fluss Sink # Fluss 集群的 bootstrap 地址(必填) bootstrap.servers: localhost:9123 # 传给 Fluss client 的安全相关属性 properties.client.security.protocol: sasl properties.client.security.sasl.mechanism: PLAIN properties.client.security.sasl.username: developer properties.client.security.sasl.password: developer-pass pipeline: name: MySQL to Fluss Pipeline parallelism: 2 # LENIENT 模式:允许以宽松方式处理 Schema 变更 schema.change.behavior: LENIENTsink段中的type: fluss是连接器的唯一标识。在源码中,FlussDataSinkFactory将IDENTIFIER定义为"fluss",pipeline 在创建 DataSink 时通过该标识完成工厂匹配(见 FlussDataSinkFactory.java)。
关于 Schema 变更行为的整体说明,可参考仓库中的 Schema 变更同步文档;完整的 MySQL 到 Fluss 入门示例可参考 Postgres 到 Fluss 快速入门。
Pipeline 连接器选项
以下是 Fluss Pipeline 连接器支持的完整选项(继承自原文档表格,并整理为 Markdown 格式):
| Option | Required | Default | Type | Description |
|---|---|---|---|---|
| type | required | (none) | String | 指定要使用的连接器,这里需要设置成'fluss'。 |
| name | optional | (none) | String | Sink 的名称。 |
| bootstrap.servers | required | (none) | String | 用于建立与 Fluss 集群初始连接的主机/端口对列表。 |
| bucket.key | optional | (none) | String | 指定每个 Fluss 表的数据分布策略。表之间用;分隔,分桶键之间用,分隔。格式:database1.table1:key1,key2;database1.table2:key3。数据将根据分桶键的哈希值分配到各个桶中(分桶键必须是主键的子集,且不含主键表的分区键)。若表有主键但未指定分桶键,则分桶键默认为主键(不含分区键);若表无主键且未指定分桶键,则数据将随机分配到各个桶中。 |
| bucket.num | optional | (none) | String | 每个 Fluss 表的桶数量。表之间用;分隔。格式:database1.table1:4;database1.table2:8。 |
| properties.table.* | optional | (none) | String | 将 Fluss table 支持的参数传递给 pipeline(对应 Fluss 的存储相关选项,即 Fluss table options)。 |
| properties.client.* | optional | (none) | String | 将 Fluss client 支持的参数传递给 pipeline(对应 Fluss 的写入相关选项,即 Fluss client options)。 |
以上选项在源码中的定义与默认行为与文档完全一致:FlussDataSinkOptions声明了bootstrap.servers、bucket.key、bucket.num三个选项,并定义了TABLE_PROPERTIES_PREFIX = "properties.table."与CLIENT_PROPERTIES_PREFIX = "properties.client."两个属性前缀(见 FlussDataSinkOptions.java)。
从工厂实现可以确认以下几点(见 FlussDataSinkFactory.java):
- 必填项只有
bootstrap.servers:requiredOptions()中仅包含BOOTSTRAP_SERVERS,其余均为可选项。 - 前缀透传机制:
properties.client.*前缀的属性会去掉properties.前缀后写入 Fluss 客户端Configuration(例如properties.client.security.protocol会被转换成 Fluss 客户端的security.protocol);properties.table.*前缀的属性则以table.xxx的形式作为建表属性传递。 - 工厂在解析时会跳过(validateExcept)
properties.client.*与properties.table.*两类前缀,其余选项按严格校验。
分桶配置的解析细节
bucket.key与bucket.num均为字符串,由FlussConfigUtils负责解析成结构化配置(见 FlussConfigUtils.java):
bucket.key:先按;拆分为多个表配置,再对每项按:拆分为「表名 + 分桶键列表」,分桶键之间按,拆分。格式不合法(缺少:)会抛出IllegalArgumentException。bucket.num:按同样的;与:规则解析,桶数量必须是合法整数,否则抛出IllegalArgumentException。
实际使用示例:
sink: type: fluss bootstrap.servers: localhost:9123 # adb.orders 表按 order_id 分 4 个桶;bdb.users 表按 user_id,region 分 8 个桶 bucket.key: adb.orders:order_id;bdb.users:user_id,region bucket.num: adb.orders:4;bdb.users:8使用说明
支持 Fluss 主键表和日志表
连接器同时支持 Fluss 的主键表(primary key table)与日志表(log table)。源表带主键时,自动建表会保留主键约束;源表无主键时,自动创建为日志表。
关于自动建表
当 Fluss 中不存在目标表时,连接器会自动建库建表,自动建表遵循以下规则:
- 没有分区键:自动建的表不含分区键,源端 Schema 中的分区键信息会被忽略(见 FlussConversions.java 中
partitionedBy的使用方式)。 - 桶数量由
bucket.num选项控制:未配置时使用 Fluss 默认桶数量。 - 数据分布由
bucket.key选项控制:- 对于主键表,若未指定分桶键,则分桶键默认为主键(不含分区键)。在源码中,
toFlussTable会在未提供bucketKeys时自动计算「主键列 - 分区键」作为分桶键(见 FlussConversions.java)。 - 对于无主键的日志表,若未指定分桶键,则数据将随机分配到各个桶中。
- 对于主键表,若未指定分桶键,则分桶键默认为主键(不含分区键)。在源码中,
源码中applyCreateTable的完整建表链路为:通过ConnectionFactory建立 Fluss 连接 →createDatabase(不存在则创建)→ 若表不存在则createTable;若表已存在,则执行sanityCheck,校验 Fluss 现有表与 Flink CDC 推断出的表在主键列、分桶键、分区键三方面是否一致,不一致时抛出校验异常,防止意外的 Schema 演进(见 FlussMetaDataApplier.java 与 sanityCheck 实现)。
Schema 变更同步(lenient 模式)
连接器支持在lenient模式下进行 Schema 变更同步,通过schema.change.behavior: lenient配置(注意该配置不区分大小写,示例中同时出现了LENIENT与lenient两种写法)。支持以下 Schema 变更事件:
- 新增列— 新列会追加到 Fluss 表中(以
LAST位置追加,源码见 FlussMetaDataApplier.java)。 - 删除列— 在 lenient 模式下不会真正删除列,而是忽略该删除操作,后续写入时将该列的值设为 null。
- 重命名列— 在 lenient 模式下,此操作会被转换为「新增列 + 将旧列类型修改为可空」的序列。
- 修改列类型— 不支持。
要启用 Schema 变更同步,请在 pipeline 中配置schema.change.behavior: lenient。如果想要忽略所有 Schema 变更,使用schema.change.behavior: IGNORE。
从源码看,FlussMetaDataApplier实现了MetadataApplier接口,其applySchemaChange目前支持CreateTableEvent、DropTableEvent与AddColumnEvent三类事件,其他 Schema 变更事件会抛出异常(见 FlussMetaDataApplier.java)。这正好解释了文档中「删除列/重命名列由 lenient 模式在 Pipeline 层面转换为可接受事件」的语义:LENIENT模式会把上游的删除、重命名等操作降级为对现有表「追加列」等安全操作,而IGNORE模式则直接丢弃全部 Schema 变更。
关于数据同步
数据同步部分由 Fluss Java Client 完成:FlussDataSink通过FlussSink+FlussEventSerializationSchema构建 Flink Sink Provider,将 CDC 事件序列化后写入 Fluss(见 FlussDataSink.java)。
另外,连接器实现了FlussHashFunctionProvider,用于在 Flink CDC 内部将带主键的数据事件按「表 ID + 主键值」计算哈希并分发到对应子任务,从而保证同一主键的变更按序处理;对于无主键的日志表,为避免所有事件落到同一子任务,则加入随机数参与哈希(见 FlussHashFunctionProvider.java)。
依赖信息方面,该连接器模块基于fluss-client(当前仓库 pom 中fluss.version为0.9.0-incubating)构建,打包时会将org.apache.fluss:*shade 进连接器 JAR(见 pom.xml 与 shade 配置)。
数据类型映射
Flink CDC 类型到 Fluss 类型的映射关系如下(继承自原文档表格):
| Flink CDC type | Fluss type | Note |
|---|---|---|
| TINYINT | TINYINT | |
| SMALLINT | SMALLINT | |
| INT | INT | |
| BIGINT | BIGINT | |
| FLOAT | FLOAT | |
| DOUBLE | DOUBLE | |
| DECIMAL(p, s) | DECIMAL(p, s) | |
| BOOLEAN | BOOLEAN | |
| DATE | DATE | |
| TIME | TIME | |
| TIMESTAMP | TIMESTAMP | |
| TIMESTAMP_LTZ | TIMESTAMP_LTZ | |
| CHAR(n) | CHAR(n) | |
| VARCHAR(n) | STRING | |
| BINARY(n) | BINARY(n) | |
| VARBINARY(N) | BYTES | |
| ARRAY | ARRAY | 元素类型递归映射。 |
| MAP | MAP | 键和值类型递归映射。 |
| ROW | ROW | 字段类型递归映射。 |
上述映射由FlussConversions.CdcTypeToFlussType类型转换器实现,与表格逐项对应(见 FlussConversions.java)。值得注意的几个细节:
VARCHAR(n)映射为STRING:Fluss 不提供 varchar 类型,源码注释中明确说明,长度信息因此不会被保留。VARBINARY(n)映射为BYTES:同理,Fluss 不提供 varbinary 类型。TIMESTAMP_LTZ映射为 Fluss 的LocalZonedTimestampType;而ZONED_TIMESTAMP不支持,转换时会直接抛出UnsupportedOperationException。- ARRAY / MAP / ROW 均为递归映射,复合类型的子元素、键值类型与字段类型会逐一递归转换。
- 所有转换均保留字段的可空性(nullable)与精度/长度等元信息。
这些映射行为同时被单元测试覆盖,例如 FlussConversionsTest.java 中通过 18 个字段的 Schema 逐一断言了BOOLEAN → TINYINT → SMALLINT → INT → BIGINT → FLOAT → DOUBLE → DECIMAL → CHAR → STRING → BINARY → BYTES → DATE → TIME → TIMESTAMP → TIMESTAMP_LTZ等完整映射,并验证了分桶键默认值、表属性透传、注释、分区键与主键校验等行为。
总结
Fluss Pipeline 连接器是 Flink CDC Pipeline 体系中最具「流式存储」特色的 Sink 之一:它把自动建表、分桶数据分布、主键/日志双表类型支持和 lenient 模式 Schema 变更同步集成在一个 YAML 配置段中。结合 FlussDataSinkFactory、FlussMetaDataApplier 与 FlussConversions 等源码,你可以进一步按需定制分桶策略、安全配置与数据类型行为。相关的更多入门内容可继续阅读 Postgres 到 Fluss 快速入门、Data Sink 概念 与 Schema 变更同步。
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考