
Flink CDC 3.5.0 从理论到实践 —— 第 3 章 环境准备与版本依赖课程定位本系列教程以MySQL 为唯一数据源Sink 覆盖Doris / Paimon / Kafka三大目标从原理到生产落地全链路实战。版本基线Flink CDC 3.5.0 Flink 1.20.x MySQL 8.0/8.4 Doris 4.1 Paimon 1.4.2 Kafka 3.x章节导读3.1 软件版本矩阵与兼容性3.2 MySQL 端配置3.3 Flink 集群准备3.4 依赖获取与部署清单3.5 课堂环境快速拉起Docker Compose3.6 环境连通性验证3.7 本章小结3.1 软件版本矩阵与兼容性环境准备的第一步是锁定版本矩阵。CDC 链路组件多、依赖关系复杂版本不匹配是新手最容易踩坑的地方务必先对齐再动手。3.1.1 本教程版本基线组件版本角色说明JDK8 / 11 / 17运行时Flink 1.20 与 CDC 3.5.0 均支持推荐 JDK 17Flink1.20.x计算引擎建议 1.20.1 及以上Flink CDC3.5.0采集与管道本教程核心Pipeline YAML 方式MySQL8.0.x / 8.4数据源5.7 也支持但建议 8.0Doris4.1OLAP SinkFE BEStream Load 写入Paimon1.4.2数据湖 Sink推荐 1.1.x 及以上Kafka3.x消息 Sink2.8 即可推荐 3.7Hadoop / HDFS3.3.x可选存储Paimon 生产用 HDFS开发可用本地盘3.1.2 连接器兼容性矩阵Flink CDC 3.5.0 与各 Sink 连接器的兼容关系关键版本约束连接器制品兼容 Flink兼容目标库获取方式MySQL Pipelineflink-cdc-pipeline-connector-mysql-3.5.0.jar1.20MySQL 5.6~8.4随 Flink CDC dist 内置Doris Pipelineflink-cdc-pipeline-connector-doris-3.5.0.jar1.20Doris 1.0含 4.1随 Flink CDC dist 内置Kafka Pipelineflink-cdc-pipeline-connector-kafka-3.5.0.jar1.20Kafka 2.8随 Flink CDC dist 内置Paimon Pipelineflink-cdc-pipeline-connector-paimon-3.5.0.jar1.20Paimon 1.4.2随 dist 内置 额外 paimon-flink jarDoris Flinkflink-doris-connector-1.2025.1.01.15~1.20Doris 1.0含 4.1官网下载 / Maven注意Flink CDC 3.x 各版本对 Flink 小版本有严格要求如 3.5.0 面向 Flink 1.20混用会导致NoSuchMethodError等运行时错误flink-doris-connector与 Flink CDC 的 Doris Pipeline 是两个不同的制品前者用于 Flink SQL 读写 Doris后者用于整库同步管道不要混淆。3.1.3 部署拓扑规划本教程使用的组件拓扑与端口约定┌─────────────┐ 3306 ┌──────────────────────┐ │ MySQL │◄────────│ Flink 1.20 (Standalone) │ │ (数据源) │ Binlog │ JobManager : 8081 │ └─────────────┘ │ TaskManager x N │ └───────┬──────┬──────┬─────┘ │ │ │ ┌──────────▼┐ ┌───▼──────┐ ┌▼───────────┐ │ Doris 4.1│ │ Paimon │ │ Kafka │ │ FE : 8030 │ │ (HDFS/ │ │ Broker:9092│ │ 查询: 9030 │ │ 本地FS) │ │ │ │ BE : 8040 │ │ │ │ │ └───────────┘ └──────────┘ └────────────┘组件端口用途MySQL3306业务连接 Binlog 读取Flink WebUI8081作业管理与监控Doris FE HTTP8030Stream Load 提交端口fenodesDoris FE MySQL9030SQL 查询 / 建表验证Doris BE HTTP8040BE 直写benodesKafka9092消息读写3.2 MySQL 端配置MySQL 是整条链路的源头Binlog 配置与账号权限必须先就位否则 Flink CDC 作业启动即报错。3.2.1 前置要求清单检查项要求验证命令Binlog 开启log_bin ONSHOW VARIABLES LIKE log_bin;Binlog 格式binlog_format ROWSHOW VARIABLES LIKE binlog_format;行镜像binlog_row_image FULLSHOW VARIABLES LIKE binlog_row_image;日志压缩关闭8.0.20 默认关闭SHOW VARIABLES LIKE binlog_transaction_compression;server-id全局唯一正整数SHOW VARIABLES LIKE server_id;Binlog 保留expire_logs_days/binlog_expire_logs_seconds足够长避免作业停止期间 Binlog 被清理导致无法续读时区与下游约定一致SHOW VARIABLES LIKE time_zone;3.2.2 my.cnf 配置示例[mysqld] # Binlog 配置CDC 必需 server-id 1 log_bin mysql-bin binlog_format ROW binlog_row_image FULL # Binlog 保留 7 天按业务调整全量初始化耗时长的场景建议更长 binlog_expire_logs_seconds 604800 # GTID推荐开启便于按 GTID 恢复位点 gtid_mode ON enforce_gtid_consistency ON # 连接超时全量大表初始化时避免连接被回收 interactive_timeout 28800 wait_timeout 28800修改后重启 MySQL 生效systemctl restart mysqld验证配置SHOWVARIABLESLIKElog_bin;-- ONSHOWVARIABLESLIKEbinlog_format;-- ROWSHOWVARIABLESLIKEbinlog_row_image;-- FULLSHOWVARIABLESLIKEgtid_mode;-- ONSHOWVARIABLESLIKEbinlog_expire_logs_seconds;-- 6048003.2.3 创建 CDC 专用账号不要使用 root 做 CDC 采集按最小权限原则创建专用账号-- 1. 创建用户生产建议限制来源 IPCREATEUSERcdc_user%IDENTIFIEDBYCdc.2026;-- 2. 授予最小权限GRANTSELECT,SHOWDATABASES,REPLICATIONSLAVE,REPLICATIONCLIENTON*.*TOcdc_user%;-- 3. 刷新权限FLUSHPRIVILEGES;-- 4. 验证SHOWGRANTSFORcdc_user%;权限说明权限用途SELECT全量快照阶段读取表数据SHOW DATABASES整库同步时枚举库表REPLICATION SLAVE伪装 Slave 读取 BinlogREPLICATION CLIENT执行SHOW MASTER STATUS获取位点注意增量快照模式下无需RELOAD权限旧版本 Debezium 快照需要但若关闭增量快照则需要补充。3.2.4 准备业务测试数据创建一个业务库用于后续章节实战CREATEDATABASEappDEFAULTCHARACTERSETutf8mb4;USEapp;CREATETABLEorders(idBIGINTNOTNULLAUTO_INCREMENT,customer_idBIGINTNOTNULL,amountDECIMAL(10,2)NOTNULL,statusVARCHAR(20)NOTNULLDEFAULTcreated,create_timeDATETIMENOTNULLDEFAULTCURRENT_TIMESTAMP,update_timeDATETIMENOTNULLDEFAULTCURRENT_TIMESTAMPONUPDATECURRENT_TIMESTAMP,PRIMARYKEY(id))ENGINEInnoDB;INSERTINTOorders(customer_id,amount,status)VALUES(1001,99.90,paid),(1002,199.00,created),(1003,59.50,paid);要点CDC 表必须有主键否则无法做增量快照并发切分无主键表需指定chunk.key-column见第 2 章。3.3 Flink 集群准备Flink CDC 3.5.0 是运行在 Flink 之上的管道框架先准备好 Flink 集群。本节以Standalone 模式为主线最简单、适合学习并给出 YARN / K8s 的差异说明。3.3.1 下载与解压# 下载 Flink 1.20.xScala 版本后缀不影响1.20 起已无 Scala 区分wgethttps://archive.apache.org/dist/flink/flink-1.20.1/flink-1.20.1-bin-scala_2.12.tgz# 解压到部署目录tar-xzfflink-1.20.1-bin-scala_2.12.tgz-C/opt/apps/ln-s/opt/apps/flink-1.20.1 /opt/apps/flinkexportFLINK_HOME/opt/apps/flink3.3.2 config.yaml 核心配置Flink 1.20 使用新的config.yaml扁平化格式替代旧的flink-conf.yaml# 基础 jobmanager.rpc.address:localhostjobmanager.bind-host:0.0.0.0rest.bind-address:0.0.0.0rest.port:8081# 内存按机器规格调整jobmanager.memory.process.size:2048mtaskmanager.memory.process.size:4096mtaskmanager.numberOfTaskSlots:4# CheckpointExactly-Once 前提见第 2 章execution.checkpointing.interval:60sexecution.checkpointing.mode:EXACTLY_ONCEexecution.checkpointing.timeout:10minexecution.checkpointing.min-pause:30sexecution.checkpointing.externalized-checkpoint-retention:RETAIN_ON_CANCELLATIONstate.checkpoints.dir:file:///opt/apps/flink/checkpointsstate.savepoints.dir:file:///opt/apps/flink/savepoints# 重启策略 restart-strategy.type:fixed-delayrestart-strategy.fixed-delay.attempts:3restart-strategy.fixed-delay.delay:10s# 并行度默认值 parallelism.default:2说明开发环境 Checkpoint 目录用本地盘file://即可生产建议 HDFS / S3hdfs:///s3://RETAIN_ON_CANCELLATION保证作业取消后 Checkpoint 保留便于从断点恢复若继续使用旧版flink-conf.yamlFlink 1.20 仍兼容但推荐迁移到config.yaml。3.3.3 启动与验证cd/opt/apps/flink# 启动 Standalone 集群./bin/start-cluster.sh# 验证进程应有 StandaloneSessionClusterEntrypoint 和 TaskManagerRunnerjps|grep-EStandaloneSession|TaskManager# 验证 WebUI浏览器访问 http://host:8081curl-shttp://localhost:8081/overview返回类似 JSON 即集群正常{taskmanagers:1,slots-total:4,slots-available:4,jobs-running:0,jobs-finished:0,jobs-cancelled:0,jobs-failed:0,flink-version:1.20.1}3.3.4 其他部署模式速览模式启动方式适用场景额外准备Standalonestart-cluster.sh学习、开发、小规模生产无YARN Sessionyarn-session.sh -d共享集群资源、多作业HADOOP_CLASSPATH、hadoop classpath配置YARN Per-Job / Applicationflink run -t yarn-per-job生产隔离部署同上K8s OperatorFlinkDeployment CRD云原生生产cert-manager、Operator 安装YARN 模式需要先导出 Hadoop 环境变量exportHADOOP_CLASSPATH$(hadoop classpath)# 提交示例./bin/flink run-tyarn-per-job-ccom.xxx.MainApp app.jar本教程主线使用 Standalone 模式讲解YARN / K8s 的生产化部署在第 10 章运维篇展开。3.4 依赖获取与部署清单Flink CDC 的依赖分两组Flink CDC 发行包Pipeline 作业运行时与Flink lib 目录的连接器 jarSQL API 作业使用。3.4.1 Flink CDC 发行包# 下载 Flink CDC 3.5.0 二进制发行包wgethttps://archive.apache.org/dist/flink/flink-cdc-3.5.0/flink-cdc-3.5.0-bin.tar.gz# 解压tar-xzfflink-cdc-3.5.0-bin.tar.gz-C/opt/apps/# 目录结构/opt/apps/flink-cdc-3.5.0/ ├── bin/ │ └── flink-cdc.sh# Pipeline 作业提交入口├── lib/ │ ├── flink-cdc-dist-3.5.0.jar │ ├── flink-cdc-pipeline-connector-mysql-3.5.0.jar │ ├── flink-cdc-pipeline-connector-doris-3.5.0.jar │ ├── flink-cdc-pipeline-connector-kafka-3.5.0.jar │ ├── flink-cdc-pipeline-connector-paimon-3.5.0.jar │ └── flink-cdc-runtime-3.5.0.jar └── conf/# Pipeline YAML 存放处自建要点Pipeline 作业通过bin/flink-cdc.sh提交它会把自身lib/下的连接器一起加载MySQL / Doris / Kafka / Paimon 四个 Pipeline 连接器均已内置开箱即用。3.4.2 Flink lib 目录依赖若使用Flink SQL API方式单表 CDC需要将对应 jar 放入 Flink 的lib/用途制品Maven 坐标 / 下载MySQL CDC Sourceflink-sql-connector-mysql-cdc-3.5.0.jarorg.apache.flink:flink-sql-connector-mysql-cdc:3.5.0Doris 读写flink-doris-connector-1.20-25.1.0.jarorg.apache.doris:flink-doris-connector-1.20:25.1.0Paimonpaimon-flink-1.20-1.4.x.jarorg.apache.paimon:paimon-flink-1.20:1.4.xKafka 连接器flink-sql-connector-kafka-3.4.0-1.20.jarorg.apache.flink:flink-sql-connector-kafka:3.4.0-1.20MySQL JDBC 驱动mysql-connector-j-8.0.33.jarcom.mysql:mysql-connector-j:8.0.33cd$FLINK_HOME/lib# 示例版本号以实际下载为准wgethttps://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-mysql-cdc/3.5.0/flink-sql-connector-mysql-cdc-3.5.0.jarwgethttps://repo1.maven.org/maven2/org/apache/paimon/paimon-flink-1.20/1.4.2/paimon-flink-1.20-1.4.2.jarwgethttps://repo1.maven.org/maven2/org/apache/flink/flink-sql-connector-kafka/3.4.0-1.20/flink-sql-connector-kafka-3.4.0-1.20.jarwgethttps://repo1.maven.org/maven2/com/mysql/mysql-connector-j/8.0.33/mysql-connector-j-8.0.33.jar# 重启集群使 jar 生效./bin/stop-cluster.sh./bin/start-cluster.sh注意lib/下不要放同一制品的多个版本避免类冲突Paimon jar 同时放入flink-cdc-3.5.0/lib/Pipeline 方式写 Paimon 时也需要它Doris Pipeline 连接器已内置 MySQL 协议依赖无需额外驱动SQL API 方式写 Doris 时按需添加 JDBC 驱动。3.4.3 部署清单汇总/opt/apps/ ├── flink/ # Flink 1.20.1 │ ├── bin/{start-cluster.sh, sql-client.sh, flink} │ ├── conf/config.yaml │ └── lib/ # SQL API 依赖见 3.4.2 ├── flink-cdc-3.5.0/ # Flink CDC 发行包 │ ├── bin/flink-cdc.sh │ ├── lib/ # Pipeline 连接器内置 │ └── conf/ # Pipeline YAML ├── mysql-8.0.x / 或外部 RDS # 数据源 ├── doris-4.1/ (FE BE) # OLAP Sink └── kafka-3.7/ # 消息 Sink3.5 课堂环境快速拉起Docker Compose为了让大家零依赖、10 分钟内搭好实验环境本节提供一键 Docker Compose 套件MySQL Doris Kafka。Paimon 在开发模式直接使用本地文件系统无需额外部署。3.5.1 docker-compose.ymlversion:3.8services:# MySQL 数据源Binlog 已开启mysql:image:mysql:8.0container_name:cdc-mysqlports:-3306:3306environment:MYSQL_ROOT_PASSWORD:Root.2026MYSQL_DATABASE:appcommand:--server-id1--log-binmysql-bin--binlog-formatROW--binlog-row-imageFULL--gtid-modeON--enforce-gtid-consistencyON--default-authentication-pluginmysql_native_passwordvolumes:-mysql-data:/var/lib/mysqlhealthcheck:test:[CMD,mysqladmin,ping,-uroot,-pRoot.2026]interval:10sretries:5# Doris 4.1 FE doris-fe:image:apache/doris:fe-4.1.0-x86_64container_name:cdc-doris-feports:-8030:8030# Stream Load-9030:9030# MySQL 协议查询environment:FE_SERVERS:fe1:cdc-doris-fe:9010FE_ID:1volumes:-doris-fe-data:/opt/apache-doris/fe/doris-meta# Doris 4.1 BE doris-be:image:apache/doris:be-4.1.0-x86_64container_name:cdc-doris-beports:-8040:8040# BE HTTPenvironment:FE_SERVERS:fe1:cdc-doris-fe:9010BE_ADDR:cdc-doris-be:9050volumes:-doris-be-data:/opt/apache-doris/be/storagedepends_on:-doris-fe# KafkaKRaft 模式免 Zookeeperkafka:image:bitnami/kafka:3.7container_name:cdc-kafkaports:-9092:9092environment:KAFKA_CFG_NODE_ID:1KAFKA_CFG_PROCESS_ROLES:broker,controllerKAFKA_CFG_CONTROLLER_QUORUM_VOTERS:1kafka:9093KAFKA_CFG_LISTENERS:PLAINTEXT://:9092,CONTROLLER://:9093KAFKA_CFG_ADVERTISED_LISTENERS:PLAINTEXT://kafka:9092KAFKA_CFG_CONTROLLER_LISTENER_NAMES:CONTROLLERKAFKA_CFG_AUTO_CREATE_TOPICS_ENABLE:truevolumes:-kafka-data:/bitnami/kafkavolumes:mysql-data:doris-fe-data:doris-be-data:kafka-data:注意Doris 镜像 tag 以 Docker Hub 官方仓库 实际发布的 4.1 tag 为准Doris BE 容器内存建议 ≥ 4G宿主机内存不足时调低BE的mem_limit生产环境请勿使用容器内单节点 Doris按集群规划部署。3.5.2 一键启动与初始化# 1. 启动全部服务dockercompose up-d# 2. 等待健康检查通过约 1-2 分钟dockercomposeps# 3. 初始化 MySQL 业务数据dockerexec-icdc-mysql mysql-uroot-pRoot.2026EOF CREATE USER cdc_user% IDENTIFIED BY Cdc.2026; GRANT SELECT, SHOW DATABASES, REPLICATION SLAVE, REPLICATION CLIENT ON *.* TO cdc_user%; FLUSH PRIVILEGES; CREATE DATABASE app DEFAULT CHARACTER SET utf8mb4; USE app; CREATE TABLE orders ( id BIGINT NOT NULL AUTO_INCREMENT, customer_id BIGINT NOT NULL, amount DECIMAL(10,2) NOT NULL, status VARCHAR(20) NOT NULL DEFAULT created, create_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP, update_time DATETIME NOT NULL DEFAULT CURRENT_TIMESTAMP ON UPDATE CURRENT_TIMESTAMP, PRIMARY KEY (id) ) ENGINE InnoDB; INSERT INTO orders (customer_id, amount, status) VALUES (1001, 99.90, paid), (1002, 199.00, created), (1003, 59.50, paid); EOF# 4. 注册 Doris BE 到 FEdockerexeccdc-doris-be mysql-hcdc-doris-fe-P9030-uroot\-eALTER SYSTEM ADD BACKEND cdc-doris-be:9050;3.5.3 各组件快速验证# MySQL验证 Binlog 与权限dockerexec-icdc-mysql mysql-ucdc_user-pCdc.2026\-eSHOW VARIABLES LIKE log_bin; SHOW GRANTS;# Doris验证 FE/BE 状态dockerexeccdc-doris-be mysql-hcdc-doris-fe-P9030-uroot\-eSHOW FRONTENDS\G SHOW BACKENDS\G# Kafka验证 Brokerdockerexeccdc-kafka kafka-topics.sh --bootstrap-server kafka:9092--list3.6 环境连通性验证在进入实战章节前用下面的Checklist逐项确认环境就绪。任何一项不通先解决再继续否则后续章节会连环报错。3.6.1 连通性验证清单#验证项命令 / 方式预期结果1MySQL BinlogSHOW VARIABLES LIKE log_binON2MySQL ROW 格式SHOW VARIABLES LIKE binlog_formatROW3CDC 账号权限SHOW GRANTS FOR cdc_user%含 4 项复制权限4Flink 集群浏览器访问http://host:8081WebUI 正常Slot 可用5Flink Checkpointconfig.yaml三项配置interval / dir / retention 已配置6Doris FESHOW FRONTENDSAlive true7Doris BESHOW BACKENDSAlive true且TotalCapacity正常8Doris Stream Load 端口curl http://fe:8030/api/health返回 OK / 2009Kafka Brokerkafka-topics.sh --bootstrap-server host:9092 --list无报错10网络互通在 Flink 节点telnet mysql 3306等端口可达11CDC 发行包ls $CDC_HOME/lib4 个 Pipeline 连接器 jar 存在12Flink libls $FLINK_HOME/lib所需连接器 jar 存在且无重复版本3.6.2 常见环境问题速查问题现象原因解决Access denied; you need REPLICATION SLAVE privilege账号缺少复制权限补授REPLICATION SLAVE/CLIENTThe MySQL server is not using binlogBinlog 未开启或配置后未重启检查log_bin重启 MySQLFlink WebUI 打不开端口冲突或防火墙检查 8081 端口与rest.bind-addressCould not find any factory for identifier dorisDoris 连接器 jar 未放入对应 lib按 3.4 节补齐 jar 并重启Server runtime does not match/NoSuchMethodErrorFlink 与连接器版本不匹配严格对齐 3.1.2 兼容矩阵Doris BE Alive false内存不足或未注册检查 BE 日志重新ADD BACKENDKafka 连接超时监听地址为容器内网检查ADVERTISED_LISTENERS配置3.7 本章小结本章完成了整条实战链路的环境搭建版本矩阵锁定 Flink CDC 3.5.0 Flink 1.20 MySQL 8.0/8.4 Doris 4.1 Paimon 1.4.2 Kafka 3.x 的基线组合牢记连接器兼容性是第一踩坑点。MySQL 配置BinlogROW FULL、GTID、专用账号四项复制权限业务表必须有主键——这是 CDC 数据源的三板斧。Flink 集群Standalone 为主线config.yaml核心配置Checkpoint 重启策略并预告 YARN / K8s 生产化方式。依赖部署Flink CDC 发行包内置四大 Pipeline 连接器flink-cdc.sh提交即用SQL API 方式则按需向 Flinklib/添加连接器 jar。快速环境Docker Compose 一键拉起 MySQL Doris 4.1 KafkaPaimon 开发态用本地文件系统。连通性验证12 项 Checklist 高频问题速查环境通则后续顺。下一章预告第 4 章《MySQL CDC Source 全方位解析》将以本章环境为基础逐参数讲解 MySQL Pipeline Source 的配置、启动模式与快照调优并用 Print Sink 完成第一次端到端验证。参考资料Flink CDC 3.5.0 官方文档https://nightlies.apache.org/flink/flink-cdc-docs-release-3.5/Flink 1.20 下载页https://flink.apache.org/downloads/Flink 1.20 配置文档https://nightlies.apache.org/flink/flink-docs-release-1.20/docs/deployment/config/Doris 官网4.1https://doris.apache.org/Doris Docker 部署https://doris.apache.org/docs/install/construct-docker/Paimon 下载页https://paimon.apache.org/docs/master/project/download/Kafka 下载页https://kafka.apache.org/downloads本文是《Flink CDC 3.5.0 从理论到实践》系列教程的第 3 章后续章节将持续更新欢迎关注收藏。