ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Flink CDC 入门指南:基于 YAML 的流式数据集成管道、CLI 提交机制与源码级实现解析

Flink CDC 入门指南:基于 YAML 的流式数据集成管道、CLI 提交机制与源码级实现解析 Flink CDC 入门指南基于 YAML 的流式数据集成管道、CLI 提交机制与源码级实现解析【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc本文围绕 Flink CDC 官方入门文档展开系统讲解 Flink CDC 的定位与核心能力、JDK/Flink 版本要求、受支持的 Connector 生态以及如何使用 YAML 定义数据管道Data Pipeline并通过flink-cdc.sh将其编译、提交到 Flink 集群同时结合仓库中的 CLI、Composer 等模块源码剖析一条 YAML 管道定义如何被解析、翻译并最终成为一个可运行的 Flink 作业帮助读者从能提交作业进阶到理解提交链路。一、Flink CDC 是什么YAML 驱动的流式数据集成工具Flink CDC 是一个流式数据集成工具streaming data integration tool其核心设计目标是为用户提供更健壮的数据集成 API。与传统需要手写 Flink DataStream/SQL 代码不同Flink CDC 允许用户通过 YAML 文件优雅地描述 ETL 管道逻辑框架会自动生成定制的 Flink 算子并负责作业提交。按照 入门文档 的表述Flink CDC 优先优化的是作业提交过程即把写作业的成本降到最低并在此之上提供了若干增强能力端到端数据集成框架End-to-end data integration framework面向数据集成用户的简易作业构建 APISource / Sink 侧的多表Multi-table支持整库同步Synchronization of entire databasesSchema 演进Schema evolution能力从源码结构看这些能力对应清晰的模块划分见根 pom.xml 的modules列表flink-cdc-cli命令行入口、flink-cdc-composer把管道定义翻译成 Flink 作业的编译器、flink-cdc-runtime管道运行时如分片、事件序列化、flink-cdc-common公共事件模型与配置、flink-cdc-connectSource 与 Pipeline Connector 集合以及flink-cdc-flink1-compat/flink-cdc-flink2-compat两个 Flink 版本兼容层——后者解释了它为何能同时适配 Flink 1.20.x 与 Flink 2.2.x 两个大版本。二、环境要求Flink CDC 对运行环境有如下硬性要求引自入门文档并可用仓库构建配置交叉验证JDKJDK 11 或更高版本Flink CDC 自 3.6.0 版本起基于 JDK 11 构建。仓库根 pom.xml 中java.version、source.java.version、target.java.version均为11与文档一致。Apache FlinkFlink 1.20.x 或 Flink 2.2.x。根 pom.xml 中的依赖版本属性印证了这一点flink.1.x.version为1.20.3flink.2.x.version为2.2.0默认flink.version取 1.20.x。java -version # 运行 Flink CDC 前先用该命令确认 JDK 版本正确三、支持的 Connector 生态Flink CDC 提供丰富的 Connector 生态以对接各类外部系统。入门文档将 Connector 分为两类Source Connector面向 Flink Source API用于构建 CDC 源表与Pipeline Connector面向数据管道 API分 Pipeline Source / Pipeline Sink。当前仓库支持矩阵如下链接指向仓库内对应 Connector 文档Connector类型MySQLSource Connector / Pipeline Source ConnectorOracleSource Connector / Pipeline Source ConnectorPostgreSQLSource Connector / Pipeline Source ConnectorDb2Source ConnectorMongoDBSource ConnectorSQL ServerSource ConnectorTiDBSource ConnectorVitessSource ConnectorApache DorisPipeline Sink ConnectorElasticsearchPipeline Sink ConnectorFlussPipeline Sink ConnectorHudiPipeline Sink ConnectorIcebergPipeline Sink ConnectorKafkaPipeline Sink ConnectorMaxComputePipeline Sink ConnectorOceanBasePipeline Sink ConnectorPaimonPipeline Sink ConnectorStarRocksPipeline Sink Connector各 Connector 的下载地址与打包方式可参考 Flink Source Connectors 总览 与 Pipeline Connectors 总览 页面。从仓库结构看每类 Connector 都遵循核心模块 SQL 连接器模块的命名约定例如flink-connector-mysql-cdc与flink-sql-connector-mysql-cdc、flink-cdc-pipeline-connector-doris均位于flink-cdc-connect/目录下。四、用 YAML 定义一条数据管道Flink CDC 提供了一个更适合数据集成场景的 YAML 格式用户 API。以下 YAML 定义了一条完整管道摄入 MySQL 的实时变更并同步到 Apache Doris该示例完整继承自入门文档source: type: mysql hostname: localhost port: 3306 username: root password: 123456 tables: app_db.\.* server-id: 5400-5404 server-time-zone: UTC sink: type: doris fenodes: 127.0.0.1:8030 username: root password: table.create.properties.light_schema_change: true table.create.properties.replication_num: 1 pipeline: name: Sync MySQL Database to Doris parallelism: 24.1 三个必备部分与可选部分一条管道对应的是一条 Flink 算子链。依据 Data Pipeline 概念文档 与 PipelineDef 的实现YAML 中必备source数据源、sink数据目标、pipeline管道级配置可选route源表到目标表的路由规则、transform对数据变更事件施加投影/过滤/计算、user-defined-function自定义函数等。PipelineDef类以source、sink、routes、transforms、udfs、models、config七个字段承载这些定义其 Javadoc 明确说明一个定义在提交给计算引擎前会先被PipelineComposer翻译成PipelineExecution。构造函数中还能看到两处重要的隐式行为时区归一化local-time-zone若未显式设置则回退为系统默认时区若显式设置则会校验其必须是合法的 Time Zone Database ID如America/Los_Angeles、GMT08:00、UTC否则抛出IllegalArgumentException运行模式默认值execution.runtime-mode未指定时自动补全为STREAMING。4.2 pipeline 段配置项一览pipeline段整体必填不能为空但其中每个参数单独都是可选的。完整配置项如下继承自 Data Pipeline 文档参数含义可选/必填name管道名称作为作业名提交到 Flink 集群可选parallelism管道全局并行度默认 1可选local-time-zone当前会话时区 ID可选execution.runtime-mode运行模式STREAMING / BATCH默认 STREAMING可选route-moderoute 规则匹配模式ALL_MATCH默认应用所有匹配规则或FIRST_MATCH仅应用第一条匹配规则可选schema.change.behavior处理 Schema 变更 的方式exception/evolve/try_evolve/lenient默认 /ignore可选schema.operator.uidSchema 算子的唯一 ID用于算子间通信。已废弃请改用operator.uid.prefix可选schema-operator.rpc-timeoutSchemaOperator 等待下游 SchemaChangeEvent 应用完成的超时时间默认 3 分钟可选operator.uid.prefix所有管道算子 UID 的前缀。未设置时 UID 由 Flink 生成建议显式设置以保证状态升级、排障与 Flink UI 诊断时的 UID 稳定可识别可选sink.partitioning.strategy写入 Sink 时的分区策略默认SINK_DEFINED可选值SINK_DEFINED使用 Sink 定义的策略、PRIMARY_KEY按表 ID 主键分区、TABLE_ID仅按表 ID 分区可选transform.decimal.precision.modetransform 表达式求值中 DECIMAL 类型的最大精度模式UP_TO_19默认对齐 Calcite 默认类型系统或UP_TO_38可选在更复杂的场景中管道还可以叠加transform与route。例如 Data Pipeline 文档中的完整示例 演示了如何对adb.web_order01、adb.web_order02两张表做投影与过滤projection: \*, format(%S, product_name) as product_name再把app_db.orders等三张表分别路由到ods_db.ods_orders等目标表并通过user-defined-function注册addone、format两个 UDF——这正是多表同步 转换 路由组合的典型写法。五、从 YAML 到 Flink 作业源码级编译链路flink-cdc.sh提交动作的背后是一条解析 → 翻译 → 执行的调用链。结合仓库源码可以还原如下1. 解析层flink-cdc-cli 模块CliExecutor 是核心入口。其deployWithComposer方法展示了主干流程PipelineDefinitionParser pipelineDefinitionParser new YamlPipelineDefinitionParser(); PipelineDef pipelineDef pipelineDefinitionParser.parse(pipelineDefPath, globalPipelineConfig); PipelineExecution execution composer.compose(pipelineDef); return execution.execute();即YamlPipelineDefinitionParser把 YAML 文件解析为强类型PipelineDefPipelineComposer实际实现为FlinkPipelineComposer把定义编译成PipelineExecution最后execute()触发作业提交。CliExecutor还支持 Application 模式的main入口解析 YAML 后通过FlinkPipelineComposer.ofApplicationCluster(env)直接运行这是 K8s / YARN Application 模式作业的入口。2. 翻译层flink-cdc-composer 模块FlinkPipelineComposer 的compose方法完成三件事读取PIPELINE_PARALLELISM并设置到StreamExecutionEnvironment调用translate(...)翻译管道——从同目录的translator包结构看翻译被拆分为若干专职组件DataSourceTranslatorSource → Flink 源算子、SchemaOperatorTranslatorSchema 变更协调算子、TransformTranslatortransform 逻辑、PartitioningTranslator按主键/表 ID 分片、DataSinkTranslatorSink 算子配合OperatorUidGenerator生成算子 UID通过addFrameworkJars()把框架 JAR 追加到 Flink 环境最后返回携带管道名称即pipeline.name的FlinkPipelineExecution。3. 部署分支仍由 CliExecutor 决定CliExecutor.run()依据 Flink 配置中的部署目标选择不同路径见 CliExecutor.java#L66-L87部署目标执行路径kubernetes-applicationK8SApplicationDeploymentExecutorApplication 模式部署yarn-applicationYarnApplicationDeploymentExecutorApplication 模式部署localFlinkPipelineComposer.ofMiniCluster()在进程内 MiniCluster 运行remote/yarn-sessionFlinkPipelineComposer.ofRemoteCluster(flinkConfig, additionalJars)向已存在的 Session 集群提交并可携带额外 JAR这一结构印证了文档中standalone / Kubernetes / YARN 三种部署模式的说法Standalone 集群走remote目标提交到 Session 集群而 K8s 与 YARN 的 Application 模式则由各自的 DeploymentExecutor 在集群侧拉起作业。4. 分发打包flink-cdc-dist 模块flink-cdc-dist模块通过 assembly.xml 组装出flink-cdc-bin发行包其中conf/flink-cdc.yaml是 Flink CDC 的全局配置模板。flink-cdc.sh脚本即运行在发行包bin/目录下配合上述 CLI 选项完成提交。六、flink-cdc.sh 命令行参数详解CLI 的完整参数定义在 CliFrontendOptions.java 中整理如下选项说明--flink-home pathFlink 安装目录Flink home路径-t/--target target部署目标可选值local、remote、yarn-session、yarn-application、kubernetes-application--global-config pathFlink CDC 管道全局配置文件路径--jar jar...与管道一起提交的额外 JAR可多个如 UDF JAR--use-mini-cluster使用 Flink MiniCluster 运行管道-s/--from-savepoint path从指定 Savepoint 恢复作业例如hdfs:///flink/savepoint-1537-cm/--claim-mode modeSavepoint 认领模式claim认领所有权被取代后删除、no_claim默认、legacy旧行为可复用部分共享文件-n/--allow-nonRestored-state允许跳过无法恢复的 Savepoint 状态当删除了 Savepoint 时刻程序中存在的算子时需要开启-D keyval指定通用 Flink 配置项可多次使用-h/--help显示帮助信息典型用法示例提交管道到 Flink 集群并附加 Flink 配置与 UDF JARbin/flink-cdc.sh pipeline-definition.yaml \ --target remote \ --flink-home /opt/flink \ --jar ./my-udf.jar \ -D parallelism.default4七、动手实践与进阶学习路径7.1 快速上手示例Quickstart入门文档按 Flink 版本提供了四条端到端 Quickstart同一管道分别适配 1.20.x 与 2.2.x 两套环境示例1.20.x2.2.xMySQL 到 Apache DorisMySQL → DorisMySQL → DorisMySQL 到 StarRocksMySQL → StarRocksMySQL → StarRocksMySQL 到 KafkaMySQL → KafkaMySQL → KafkaPostgreSQL 到 FlussPostgreSQL → FlussPostgreSQL → Fluss此外仓库还内置了基于 Docker 的cdcup 快速实验环境tools/cdcup 目录含 Dockerfile 与 cdcup.sh 脚本用法见 cdcup 快速上手指南。执行./cdcup.sh后可用init初始化并交互式选择 Flink/CDC 版本与 Connector用up启动容器pipeline yaml提交 YAML 管道mysql打开 MySQL 控制台建表造数flink打印 Web UI 地址stop/down停止并清理环境。该工具要求本地具备可用的 Docker 与 Docker Compose V2 环境。7.2 核心概念要构建更复杂的管道建议先熟悉以下核心概念文档Data Pipeline管道定义、必备/可选部分与 pipeline 级配置Data SourceSource 定义与事件流模型Data SinkSink 定义与写入门控Table ID表标识的命名规范Transform投影、过滤与计算表达式Route源表到目标表的路由规则与匹配模式Schema EvolutionSchema 变更的多种处理模式Type Mappings源端与 Flink CDC 逻辑类型的映射7.3 部署到不同模式的 Flink 集群Standalone 集群KubernetesYARN7.4 开发者与贡献者如果想把 Flink CDC 对接到自定义外部系统或直接参与框架开发可以参考理解 Flink CDC API学习开发自定义 Flink CDC Connector 的 API 设计贡献指南如何向 Flink CDC 提交贡献许可协议Flink CDC 使用的许可证说明UdfDef 与 ModelDef 的解析入口如 UdfDef以及flink-cdc-pipeline-udf-examples模块中的 UDF 示例也为在管道中注册自定义函数与 AI Model提供了可直接对照的实现参考。八、小结Flink CDC 的入门路径可以概括为四步确认环境JDK 11、Flink 1.20.x/2.2.x→编写 YAMLsource sink pipeline 三段必备route/transform 可选→flink-cdc.sh 提交按目标选择 remote/K8s/YARN 部署路径→按需扩展UDF、Schema 演进策略、分片策略。从源码结构看YamlPipelineDefinitionParser → PipelineDef → FlinkPipelineComposer → 各 Translator → FlinkPipelineExecution的编译链路正是YAML 描述即作业这一设计承诺的具体实现读者可沿着本文给出的文件路径在仓库中逐层深入验证。【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表