news 2026/9/17 17:42:47

使用 Flink CDC Fluss Pipeline 连接器:将 MySQL 实时数据写入 Fluss(含自动建表、分桶策略与 Schema 变更同步)

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
使用 Flink CDC Fluss Pipeline 连接器:将 MySQL 实时数据写入 Fluss(含自动建表、分桶策略与 Schema 变更同步)

使用 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: LENIENT

sink段中的type: fluss是连接器的唯一标识。在源码中,FlussDataSinkFactoryIDENTIFIER定义为"fluss",pipeline 在创建 DataSink 时通过该标识完成工厂匹配(见 FlussDataSinkFactory.java)。

关于 Schema 变更行为的整体说明,可参考仓库中的 Schema 变更同步文档;完整的 MySQL 到 Fluss 入门示例可参考 Postgres 到 Fluss 快速入门。

Pipeline 连接器选项

以下是 Fluss Pipeline 连接器支持的完整选项(继承自原文档表格,并整理为 Markdown 格式):

OptionRequiredDefaultTypeDescription
typerequired(none)String指定要使用的连接器,这里需要设置成'fluss'
nameoptional(none)StringSink 的名称。
bootstrap.serversrequired(none)String用于建立与 Fluss 集群初始连接的主机/端口对列表。
bucket.keyoptional(none)String指定每个 Fluss 表的数据分布策略。表之间用;分隔,分桶键之间用,分隔。格式:database1.table1:key1,key2;database1.table2:key3。数据将根据分桶键的哈希值分配到各个桶中(分桶键必须是主键的子集,且不含主键表的分区键)。若表有主键但未指定分桶键,则分桶键默认为主键(不含分区键);若表无主键且未指定分桶键,则数据将随机分配到各个桶中。
bucket.numoptional(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.serversbucket.keybucket.num三个选项,并定义了TABLE_PROPERTIES_PREFIX = "properties.table."CLIENT_PROPERTIES_PREFIX = "properties.client."两个属性前缀(见 FlussDataSinkOptions.java)。

从工厂实现可以确认以下几点(见 FlussDataSinkFactory.java):

  • 必填项只有bootstrap.serversrequiredOptions()中仅包含BOOTSTRAP_SERVERS,其余均为可选项。
  • 前缀透传机制properties.client.*前缀的属性会去掉properties.前缀后写入 Fluss 客户端Configuration(例如properties.client.security.protocol会被转换成 Fluss 客户端的security.protocol);properties.table.*前缀的属性则以table.xxx的形式作为建表属性传递。
  • 工厂在解析时会跳过(validateExcept)properties.client.*properties.table.*两类前缀,其余选项按严格校验。

分桶配置的解析细节

bucket.keybucket.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配置(注意该配置不区分大小写,示例中同时出现了LENIENTlenient两种写法)。支持以下 Schema 变更事件:

  • 新增列— 新列会追加到 Fluss 表中(以LAST位置追加,源码见 FlussMetaDataApplier.java)。
  • 删除列— 在 lenient 模式下不会真正删除列,而是忽略该删除操作,后续写入时将该列的值设为 null。
  • 重命名列— 在 lenient 模式下,此操作会被转换为「新增列 + 将旧列类型修改为可空」的序列。
  • 修改列类型— 不支持。

要启用 Schema 变更同步,请在 pipeline 中配置schema.change.behavior: lenient。如果想要忽略所有 Schema 变更,使用schema.change.behavior: IGNORE

从源码看,FlussMetaDataApplier实现了MetadataApplier接口,其applySchemaChange目前支持CreateTableEventDropTableEventAddColumnEvent三类事件,其他 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.version0.9.0-incubating)构建,打包时会将org.apache.fluss:*shade 进连接器 JAR(见 pom.xml 与 shade 配置)。

数据类型映射

Flink CDC 类型到 Fluss 类型的映射关系如下(继承自原文档表格):

Flink CDC typeFluss typeNote
TINYINTTINYINT
SMALLINTSMALLINT
INTINT
BIGINTBIGINT
FLOATFLOAT
DOUBLEDOUBLE
DECIMAL(p, s)DECIMAL(p, s)
BOOLEANBOOLEAN
DATEDATE
TIMETIME
TIMESTAMPTIMESTAMP
TIMESTAMP_LTZTIMESTAMP_LTZ
CHAR(n)CHAR(n)
VARCHAR(n)STRING
BINARY(n)BINARY(n)
VARBINARY(N)BYTES
ARRAYARRAY元素类型递归映射。
MAPMAP键和值类型递归映射。
ROWROW字段类型递归映射。

上述映射由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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/17 17:41:42

GARCH模型实战:波动率建模、ARCH检验与VaR预测

简介:这份PPT课件面向金融学、计量经济学方向的学生与研究人员,系统讲解GARCH类模型在金融时间序列波动性分析中的原理与应用。内容从Engle提出的ARCH过程切入,梳理条件异方差性的来源与ARCH(q)的建模思路,进而展开GARCH(1,1)及高…

作者头像 李华
网站建设 2026/9/17 17:41:33

机器人基础模型OM-1解析:范式变革、技术原理与落地挑战

先说个背景。这几年“基础模型”这个词在AI圈已经被说烂了,从大语言模型到多模态模型,现在终于烧到了机器人领域。RewardAI这次发布的OM-1,名字听起来低调,但“机器人基础模型”这个定位本身就值得仔细拆一下。它不是某个机械臂的…

作者头像 李华
网站建设 2026/9/17 17:40:26

遥感影像配准:SIFT与Canny协同优化实战指南

简介:本资源是一篇聚焦多源遥感影像配准关键技术的学术研究文档,面向遥感图像处理、计算机视觉及地理信息科学领域的高校师生、科研人员与工程技术人员,旨在解决不同传感器获取的遥感影像因几何变形与辐射差异导致的配准难题。文档系统阐述了…

作者头像 李华
网站建设 2026/9/17 17:35:14

WeChatMsg教程:导出微信聊天记录为HTML/Word/CSV并生成年度报告

WeChatMsg教程:导出微信聊天记录为HTML/Word/CSV并生成年度报告 【免费下载链接】WeChatMsg 提取微信聊天记录,将其导出成HTML、Word、CSV文档永久保存,对聊天记录进行分析生成年度聊天报告 项目地址: https://gitcode.com/GitHub_Trending…

作者头像 李华