尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

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

使用 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-cdcFluss 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定义为flusspipeline 在创建 DataSink 时通过该标识完成工厂匹配见 FlussDataSinkFactory.java。关于 Schema 变更行为的整体说明可参考仓库中的 Schema 变更同步文档完整的 MySQL 到 Fluss 入门示例可参考 Postgres 到 Fluss 快速入门。Pipeline 连接器选项以下是 Fluss Pipeline 连接器支持的完整选项继承自原文档表格并整理为 Markdown 格式OptionRequiredDefaultTypeDescriptiontyperequired(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.servers、bucket.key、bucket.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.protocolproperties.table.*前缀的属性则以table.xxx的形式作为建表属性传递。工厂在解析时会跳过validateExceptproperties.client.*与properties.table.*两类前缀其余选项按严格校验。分桶配置的解析细节bucket.key与bucket.num均为字符串由FlussConfigUtils负责解析成结构化配置见 FlussConfigUtils.javabucket.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通过FlussSinkFlussEventSerializationSchema构建 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 typeFluss typeNoteTINYINTTINYINTSMALLINTSMALLINTINTINTBIGINTBIGINTFLOATFLOATDOUBLEDOUBLEDECIMAL(p, s)DECIMAL(p, s)BOOLEANBOOLEANDATEDATETIMETIMETIMESTAMPTIMESTAMPTIMESTAMP_LTZTIMESTAMP_LTZCHAR(n)CHAR(n)VARCHAR(n)STRINGBINARY(n)BINARY(n)VARBINARY(N)BYTESARRAYARRAY元素类型递归映射。MAPMAP键和值类型递归映射。ROWROW字段类型递归映射。上述映射由FlussConversions.CdcTypeToFlussType类型转换器实现与表格逐项对应见 FlussConversions.java。值得注意的几个细节VARCHAR(n)映射为STRINGFluss 不提供 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),仅供参考
返回列表