Flink CDC on Kubernetes:原生 Session 集群与 Kubernetes Operator 两种部署模式实战指南
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
Flink CDC 的 Kubernetes 部署指南围绕两条主线展开:一是借助 Flink 原生 Kubernetes 集成启动 Session 集群,再将 Flink CDC 任务提交其上;二是基于 Flink Kubernetes Operator,通过自定义 Docker 镜像加 FlinkDeployment 资源声明式地运行 CDC 管道。读完本文,你将掌握两种模式下的集群准备、镜像构建、配置挂载与任务提交全流程,并能结合 Flink CDC 源码理解 CLI 入口、--use-mini-cluster参数以及发行包结构背后的实现逻辑。
一、准备工作:Kubernetes 集群要求
部署 Flink CDC 到 Kubernetes 前,需要满足以下前提条件:
- Kubernetes 版本不低于 1.9;
- 具备
list、create、deletePod 与 Service 权限的 KubeConfig,可通过~/.kube/config配置。可使用kubectl auth can-i <list|create|edit|delete> pods验证权限; - 已启用 Kubernetes DNS;
defaultServiceAccount 拥有创建、删除 Pod 的 RBAC 权限。
Flink 的原生 Kubernetes 集成允许直接把 Flink 部署到运行中的 K8s 集群,并且 Flink 可以直接与 Kubernetes API 交互,根据资源需求动态分配和释放 TaskManager。此外,Apache Flink 还提供了 Kubernetes Operator,支持 standalone 与 native 两种部署模式,大幅简化 Flink 资源在 Kubernetes 上的部署、配置与生命周期管理。本文的 Operator 模式一节即基于该 Operator 展开。
版本适配方面,可参考 pipeline 连接器总览 中的版本映射表:Flink CDC 3.6.x 对应 Flink 1.20.x 与 2.2.x,3.2.x/3.1.x/3.0.x 对应 Flink 1.17~1.19。因此选择 Flink 发行版和 Operator 镜像基础版本时,务必与所用 Flink CDC 版本匹配。
二、Session 模式:原生 K8s 集成提交 Flink CDC 任务
2.1 准备 Flink 并设置 FLINK_HOME
Flink 运行在所有类 UNIX 环境(Linux、macOS、Cygwin/Windows)。从官方发布渠道下载与 Flink CDC 版本匹配的 Flink 二进制发行包并解压:
tar -xzf flink-*.tgz然后设置FLINK_HOME环境变量:
export FLINK_HOME=/path/flink-*这个变量在后面提交任务时会用到:从源码看,CLI 启动脚本 flink-cdc.sh 会优先从命令行参数--flink-home中解析 Flink 安装目录,若未提供则回退到环境变量FLINK_HOME,两者皆无时直接报错退出。
2.2 启动 Session 集群
进入 Flink 安装目录,执行随 Flink 提供的 bash 脚本:
cd /path/flink-* ./bin/kubernetes-session.sh -Dkubernetes.cluster-id=my-first-flink-cluster启动成功后的返回信息如下:
org.apache.flink.kubernetes.utils.KubernetesUtils [] - Kubernetes deployment requires a fixed port. Configuration blob.server.port will be set to 6124 org.apache.flink.kubernetes.utils.KubernetesUtils [] - Kubernetes deployment requires a fixed port. Configuration taskmanager.rpc.port will be set to 6122 org.apache.flink.kubernetes.KubernetesClusterDescriptor [] - Please note that Flink client operations(e.g. cancel, list, stop, savepoint, etc.) won't work from outside the Kubernetes cluster since 'kubernetes.rest-service.exposed.type' has been set to ClusterIP. org.apache.flink.kubernetes.KubernetesClusterDescriptor [] - Create flink session cluster my-first-flink-cluster successfully, JobManager Web Interface: http://my-first-flink-cluster-rest.default:8081提示:默认
kubernetes.rest-service.exposed.type为ClusterIP,在集群外部无法执行 cancel、list、savepoint 等客户端操作,Web UI 也需要通过相应方式暴露(可参考 Flink 官方文档中 “Accessing Flink's Web UI” 一节)。确保 REST endpoint 可被提交任务的节点访问到。
随后需要在flink-conf.yaml中补充两项配置,把 JobManager Web Interface 的实际端口与节点 IP 填入:
rest.bind-port: {{REST_PORT}} rest.address: {{NODE_IP}}其中{{REST_PORT}}与{{NODE_IP}}替换为上一步日志中输出的 JobManager Web Interface 对应的实际值。
2.3 部署 Flink CDC
从 Flink CDC 官方 release 页下载 Flink CDC 的 tar 包并解压:
tar -xzf flink-cdc-*.tar.gz解压后的flink-cdc目录包含bin、lib、log、conf四个子目录——这一结构由 flink-cdc-dist 模块的打包描述文件 定义:lib中放置 flink-cdc-dist uber jar,bin中是启动脚本,conf中是全局配置,log为空日志目录。
再从 release 页下载所需连接器 jar(如 MySQL、Doris 连接器及mysql-connector-java-8.0.27.jar),移动到lib目录。注意:release 下载链接仅覆盖稳定版本,SNAPSHOT 依赖需要自行基于对应分支构建。
conf目录下的 flink-cdc.yaml 是 Flink CDC 管道的全局配置文件,默认内容为:
# Parallelism of the pipeline parallelism: 4 # Behavior for handling schema change events from source schema.change.behavior: EVOLVE从 CliFrontend 源码可以看到全局配置的加载顺序:优先使用命令行--global-config指定的路径,其次回退到FLINK_CDC_HOME/conf/flink-cdc.yaml,两者都没有时使用空配置(并输出警告)。
2.4 提交 Flink CDC 任务
以下是一个同步整库的管道定义示例mysql-to-doris.yaml:
################################################################################ # Description: Sync MySQL all tables to 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: "" pipeline: name: Sync MySQL Database to Doris parallelism: 2请按需修改连接参数,更多参数说明可参考 MySQL pipeline connector 与 Apache Doris pipeline connector 文档。
最后通过 CLI 将任务提交到 Session 集群:
cd /path/flink-cdc-* ./bin/flink-cdc.sh mysql-to-doris.yaml提交成功后的返回信息如下:
Pipeline has been submitted to cluster. Job ID: ae30f4580f1918bebf16752d4963dc54 Job Description: Sync MySQL Database to Doris之后即可在 Flink Web UI 中找到名为Sync MySQL Database to Doris的运行中作业。
2.5 源码视角:CLI 入口与默认部署目标
flink-cdc.sh 的核心逻辑是:确定FLINK_HOME后 source Flink 的config.sh初始化类路径,把FLINK_HOME/lib、FLINK_CDC_HOME/lib以及 Hadoop 相关 jar 拼成完整 classpath,最终执行org.apache.flink.cdc.cli.CliFrontend。
CliFrontend 的createExecutor方法解析命令行并构造执行器;其中overrideFlinkConfiguration决定了部署目标的取值:
- 指定
--use-mini-cluster时,pipeline.deploy.target被设置为local,在进程内启动 MiniCluster 运行管道; - 未指定时默认取
--target参数值,缺省为remote,即通过rest.address/rest.bind-port指向的 REST endpoint 提交到远程集群——这正是 2.2 节需要配置这两个参数的原因。
CliFrontendOptions 中还定义了其他可用参数,可在提交时补充使用:
| 选项 | 说明 |
|---|---|
--flink-home | Flink 安装目录路径 |
--global-config | Flink CDC 管道全局配置文件路径 |
--jar | 随管道一起提交的 JAR(可多次指定) |
--target | 部署目标:local、remote、yarn-session、yarn-application、kubernetes-application |
--use-mini-cluster | 使用 Flink MiniCluster 在进程内运行管道 |
-s / --from-savepoint | 从指定 savepoint 恢复任务 |
-cm / --claim-mode | savepoint 恢复时的认领模式(claim/no_claim/legacy) |
-n / --allow-nonRestored-state | 允许跳过无法恢复的 savepoint 状态 |
-D key=value | 动态覆盖 Flink 配置项,可多次指定 |
三、Operator 模式:声明式部署 Flink CDC 管道
Operator 模式的前提是集群中已部署 Flink Kubernetes Operator。此时你只需要构建一个包含 Flink CDC 的自定义 Docker 镜像,其余资源由 Operator 管理。
3.1 构建自定义 Docker 镜像
从 release 页下载 Flink CDC tar 包与所需连接器 jar,放入镜像构建目录。假设构建目录为
/opt/docker/flink-cdc,其结构如下:/opt/docker/flink-cdc ├── flink-cdc-{{< param Version >}}-bin.tar.gz ├── flink-cdc-pipeline-connector-doris-{{< param Version >}}.jar ├── flink-cdc-pipeline-connector-mysql-{{< param Version >}}.jar ├── mysql-connector-java-8.0.27.jar └── ...基于
flink官方镜像创建 Dockerfile,添加 Flink CDC 依赖:FROM flink:1.18.0-java8 ADD *.jar $FLINK_HOME/lib/ ADD flink-cdc*.tar.gz $FLINK_HOME/ RUN mv $FLINK_HOME/flink-cdc-{{< param Version >}}/lib/flink-cdc-dist-{{< param Version >}}.jar $FLINK_HOME/lib/构建完成后目录结构为:
/opt/docker/flink-cdc ├── Dockerfile ├── flink-cdc-{{< param Version >}}-bin.tar.gz ├── flink-cdc-pipeline-connector-doris-{{< param Version >}}.jar ├── flink-cdc-pipeline-connector-mysql-{{< param Version >}}.jar ├── mysql-connector-java-8.0.27.jar └── ...注意镜像中
flink-cdc-dist-<版本>.jar必须位于$FLINK_HOME/lib下,这样 FlinkDeployment 的jarURI才能以local:///opt/flink/lib/flink-cdc-dist-{{< param Version >}}.jar引用它。仓库根目录还提供了一个 Dockerfile,展示了对应的镜像构建思路:解压 flink-cdc-dist 产物到/opt/flink-cdc、重命名 dist jar 为无版本号、并把 pipeline 连接器放入/opt/flink/usrlib。构建并推送镜像:
docker build -t flink-cdc-pipeline:{{< param Version >}} . docker push flink-cdc-pipeline:{{< param Version >}}
3.2 创建 ConfigMap 挂载配置文件
Flink CDC 的配置文件(全局配置 + 管道定义)通过 ConfigMap 挂载进 Pod。示例如下,请把连接参数替换为实际值:
--- apiVersion: v1 data: flink-cdc.yaml: |- parallelism: 4 schema.change.behavior: EVOLVE mysql-to-doris.yaml: |- 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: "" pipeline: name: Sync MySQL Database to Doris parallelism: 2 kind: ConfigMap metadata: name: flink-cdc-pipeline-configmap其中flink-cdc.yaml与发行包中 conf/flink-cdc.yaml 的语义一致(全局并行度与 schema 变更行为),管道定义文件的字段含义与 Session 模式中的mysql-to-doris.yaml完全相同。
3.3 创建 FlinkDeployment YAML
以下是一个示例文件flink-cdc-pipeline-job.yaml:
--- apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: flink-cdc-pipeline-job spec: flinkConfiguration: classloader.resolve-order: parent-first state.checkpoints.dir: 'file:///tmp/checkpoints' state.savepoints.dir: 'file:///tmp/savepoints' flinkVersion: v1_18 image: 'flink-cdc-pipeline:{{< param Version >}}' imagePullPolicy: Always job: args: - '--use-mini-cluster' - /opt/flink/flink-cdc-{{< param Version >}}/conf/mysql-to-doris.yaml entryClass: org.apache.flink.cdc.cli.CliFrontend jarURI: 'local:///opt/flink/lib/flink-cdc-dist-{{< param Version >}}.jar' parallelism: 1 state: running upgradeMode: savepoint jobManager: replicas: 1 resource: cpu: 1 memory: 1024m podTemplate: apiVersion: v1 kind: Pod spec: containers: # don't modify this name - name: flink-main-container volumeMounts: - mountPath: /opt/flink/flink-cdc-{{< param Version >}}/conf name: flink-cdc-pipeline-config volumes: - configMap: name: flink-cdc-pipeline-configmap name: flink-cdc-pipeline-config restartNonce: 0 serviceAccount: flink taskManager: resource: cpu: 1 memory: 1024m该 YAML 有两个必须理解的关键点(也与源码实现一一对应):
classloader.resolve-order必须为parent-first:这是由 Flink 类加载机制决定的,保证 CDC 连接器与 Flink 运行时共享父类加载器中的类;- 必须携带
--use-mini-cluster参数:如 2.5 节所述,Flink CDC 默认以remote目标提交任务到远程 Flink 集群;而在 Operator 模式下每个 Pod 是独立运行的,没有外部 Session 集群可提交,因此需要通过--use-mini-cluster让 CliFrontend 把部署目标切换为进程内 MiniCluster。
其他要点:job.args的第二个参数是挂载后的管道定义文件路径(由podTemplate中的 volumeMounts 将 ConfigMap 挂到该目录);job.parallelism: 1指 FlinkDeployment 层作业的并行度(CDC 管道自身的并行度由 YAML 中pipeline.parallelism或全局parallelism决定);job.state: running与upgradeMode: savepoint表示 Operator 会在作业更新时基于 savepoint 做无状态丢失的滚动升级。
3.4 提交任务
ConfigMap 与 FlinkDeployment YAML 就绪后,通过 kubectl 提交:
kubectl apply -f flink-cdc-pipeline-job.yaml成功返回:
flinkdeployment.flink.apache.org/flink-cdc-pipeline-job created如需追踪日志或暴露 Flink Web UI,请参考 Flink Kubernetes Operator 的官方文档(Operator 提供了flink.kubernetes.operator.expose等日志与 UI 配置能力)。
注意:目前不支持以 native application mode 提交 Flink CDC 任务,Operator 模式请使用上述 mini-cluster 方式。
四、两种模式对比与选型建议
| 维度 | 原生 Session 模式 | Operator 模式 |
|---|---|---|
| 集群形态 | Flink 原生 K8s 集成创建的 Session 集群 | Operator 按 FlinkDeployment 管理的 Pod |
| CDC 依赖安装 | 下载 tar 包 + 连接器 jar 放入本地lib | 构建自定义 Docker 镜像 |
| 配置管理 | 本地flink-cdc.yaml与管道 YAML | ConfigMap 挂载 |
| 任务提交 | ./bin/flink-cdc.sh xxx.yaml(remote目标) | kubectl apply,--use-mini-cluster运行 |
| 生命周期管理 | 手动(重启、升级、暴露 UI) | Operator 自动化(savepoint 升级、重启策略) |
| 适用场景 | 快速验证、轻量环境 | 生产环境、需要声明式与自动化运维 |
从源码结构看,两种模式最终都汇聚到同一个入口org.apache.flink.cdc.cli.CliFrontend:Session 模式经由 flink-cdc.sh 启动并默认走 remote 部署目标,Operator 模式则作为 FlinkDeployment 的entryClass以 mini-cluster 方式运行。理解这一统一入口及其 命令行参数,有助于在两种模式间灵活切换,并为后续扩展(savepoint 恢复、动态配置覆盖)打下基础。
【免费下载链接】flink-cdcFlink CDC is a streaming data integration tool项目地址: https://gitcode.com/GitHub_Trending/flin/flink-cdc
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考