基于 Flink + Paimon + CDC 的下一代实时数据湖平台技术方案
刘志敏
2026-06-20
109 分钟阅读
Flink
Paimon
CDC
数据湖
实时计算
湖仓一体
Kubernetes
基于 Flink + Paimon + CDC 的下一代实时数据湖平台技术方案
版本:v2.1 | 日期:2026-06-20 | 技术栈:Flink 2.2 / Paimon 1.0 / CDC 3.6 / Dinky 1.2
面向中小企业的流批一体实时数据湖架构设计,基于 Apache Flink 2.0+、Apache Paimon 1.0+ 与 Flink-CDC 3.6+ 构建下一代湖仓一体平台。
目录
- 01 项目背景与建设目标
- 02 技术选型与版本约束
- 03 整体架构设计
- 04 Paimon 数仓分层表设计
- 05 核心功能模块设计
- 06 核心 YAML 配置示例
- 07 关键 SQL 语句
- 08 冷热数据分层策略
- 09 业务侧表结构变更处理方案
- 10 数据安全与恢复方案
- 11 数据血缘与监控告警
- 12 部署与运维最佳实践
- 13 附录:版本兼容性矩阵
01 项目背景与建设目标
1.1 业务背景
随着数据驱动决策的深入,传统基于 Hive 的离线数仓面临数据延迟高(T+1)、运维复杂、流批割裂等痛点。企业需要一套实时数据湖平台,实现从数据采集、存储到分析的全链路实时化,同时控制中小公司的运维成本。
1.2 建设目标
| 目标 | 说明 |
|---|---|
| 秒级数据入湖 | 通过 CDC 自动捕获业务数据库变更,秒级延迟入湖,替代传统 ETL 批处理 |
| 流批一体存储 | 同一份 Paimon 数据同时支持流式消费和批量分析,消除数据孤岛 |
| 零代码采集 | YAML 配置驱动的 CDC 同步,无需编写 SQL,整库/分表一键入湖 |
| Schema 自动演进 | 上游表结构变更自动同步至数据湖,无需人工干预 DDL |
| 低成本运维 | 精简架构,降低运维复杂度,适合中小团队 2-3 人维护 |
| 近实时 OLAP | 对接 Trino/Presto 实现秒级延迟的交互式分析查询 |
| 冷热分层存储 | S3 冷存储 + 本地 HDD 热存储,降低约 30% 存储成本 |
| 弹性伸缩 | 基于 Kubernetes 部署,支持按需扩缩容,替代传统 YARN |
02 技术选型与版本约束
2.1 核心技术栈
| 组件 | 选型 | 版本 | 角色 |
|---|---|---|---|
| 计算引擎 | Apache Flink | 2.2.x | 流批一体计算、SQL 执行引擎 |
| 存储层 | Apache Paimon | 1.0.x | 湖仓一体存储、ACID 事务 |
| 数据采集 | Flink-CDC | 3.6.x | CDC 变更捕获、YAML Pipeline 编排 |
| SQL 开发平台 | Apache Dinky | 1.2.x | Flink SQL 一站式开发、CDC 任务管理 |
| OLAP 查询 | Trino | 435+ | 近实时交互式分析 |
| 存储底座(热) | 本地 HDD/SSD | - | 近期热数据存储 |
| 存储底座(冷) | MinIO / S3 | - | 历史冷数据归档存储 |
| 编排调度 | Kubernetes + Flink Operator | 1.15 | 容器化部署、弹性伸缩 |
| 监控 | Prometheus + Grafana | - | 指标采集与可视化 |
2.2 v2.0 版本改进要点
| 改进项 | v1.0 方案 | v2.0 方案 | 改进收益 |
|---|---|---|---|
| 数据链路 | MySQL → CDC → Paimon(直连) | MySQL → CDC → Paimon(CDC 3.6 原生支持) | Exactly-Once 断点续传,零中间件 |
| Catalog | Hive Metastore | JDBC Catalog | 中心化元数据 + 分布式锁,跨引擎一致 |
| 部署方式 | YARN | Kubernetes + Flink Operator | 弹性伸缩,云原生架构 |
| SQL 开发 | Flink SQL CLI | Apache Dinky | 可视化开发、血缘分析、任务管理 |
| 存储策略 | 单一存储 | 冷热分层(S3 + HDD) | 降低约 30% 存储成本 |
| 表设计 | 统一主键表 | 按分层设计不同表类型 | ODS/DWD/DWS/ADS 各取所需 |
2.3 版本选型依据
- Flink 2.2.x:当前最新稳定版(2025-12 发布),引入存算分离状态管理、Delta Join、自适应批处理执行等核心特性,最低要求 JDK 11 [1]。
- Paimon 1.0.x:首个正式大版本,支持 Caching Catalog、Deletion Vectors、Branch/Tag 时间旅行、Clone Table、嵌套类型 Schema Evolution [2]。
- Flink-CDC 3.6.x:同时兼容 Flink 1.20.x 和 2.2.x,新增 Oracle Source、Hudi Sink,PostgreSQL 支持 Schema Evolution [3]。CDC 3.6 内置 Paimon Sink,可直接将变更写入 Paimon 表,无需经过 Kafka 等中间件,实现真正零代码入湖。
- Apache Dinky 1.2.x:一站式 Flink SQL 开发平台,支持
EXECUTE PIPELINE WITHYAML语法直接提交 CDC 任务,支持 K8s Operator 集成 [4]。 - Flink Kubernetes Operator 1.15:支持 Flink v1.19 到 v2.2,内置 Job Autoscaler 自动扩缩容 [5]。
重要提示:Flink 2.0 起最低要求 JDK 11,推荐 JDK 17。Per-Job 部署模式已移除,统一使用 Application 模式。
03 整体架构设计
3.1 架构总览(v2.0)
┌─────────────────────────────────────────────────────────────────────────┐
│ 数据源层 │
│ ┌────────┐ ┌──────────────┐ ┌────────┐ │
│ │ MySQL │ │ PostgreSQL │ │ Oracle │ │
│ └───┬────┘ └──────┬───────┘ └───┬────┘ │
├────────┼───────────────┼───────────────┼──────────────────────────────┤
│ │ 数据集成层 - Flink CDC 3.6 │ │
│ ▼ ▼ ▼ │
│ ┌─────────────────────────────────────────────────────────┐ │
│ │ CDC Pipeline (YAML 编排) │ │
│ │ → Route 路由 (分表合并) │ │
│ │ → Transform (数据清洗/脱敏) │ │
│ │ → Schema Evolution (表结构自动同步) │ │
│ │ → Paimon Sink (直接入湖,无需中间件) │ │
│ └─────────────────────────┬───────────────────────────────┘ │
│ │ │
├──────────────────────────────┼──────────────────────────────────────────┤
│ ▼ │
│ 计算引擎层 - Flink 2.2 + Dinky 1.2 │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Dinky │ │ Flink SQL │ │ Streaming │ │
│ │ SQL 开发平台 │ │ Gateway │ │ / Batch Job │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
│ ┌──────────────────────────────────────────────────────┐ │
│ │ Materialized Table (物化表) │ │
│ └──────────────────────────────────────────────────────┘ │
├────────────────────────────────────────────────────────────────────────┤
│ 存储层 - Paimon 1.0 (JDBC Catalog) │
│ MySQL 元数据存储 + 本地热存储 / S3 冷存储 │
│ │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ 热存储 (本地 HDD/SSD) - 近 30 天数据 │ │
│ │ ├── ODS: Append-only 表 (无主键,原始日志) │ │
│ │ ├── DWD: 主键表 (动态桶 + deduplicate 去重) │ │
│ │ ├── DWS: 物化表 (低频刷新,聚合汇总) │ │
│ │ └── Branch & Tag (时间旅行) │ │
│ └──────────────────────────────────────────────────────────┘ │
│ ┌──────────────────────────────────────────────────────────┐ │
│ │ 冷存储 (S3/MinIO) - 30 天以上历史数据 │ │
│ │ └── Clone 归档表 (只读,按需查询) │ │
│ └──────────────────────────────────────────────────────────┘ │
├────────────────────────────────────────────────────────────────────────┤
│ 服务层 │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Trino │ │ BI Dashboard │ │ ADS 导出 │ │
│ │ OLAP 查询 │ │ 实时看板 │ │ MySQL/ES │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
├────────────────────────────────────────────────────────────────────────┤
│ 运维监控层 │
│ ┌──────────────┐ ┌──────────────┐ ┌──────────────┐ │
│ │ Kubernetes │ │ Prometheus │ │ Dinky │ │
│ │ Flink Operator│ │ + Grafana │ │ SQL 血缘 │ │
│ └──────────────┘ └──────────────┘ └──────────────┘ │
└─────────────────────────────────────────────────────────────────────────┘
图 1:实时数据湖平台整体架构图 v2.0
3.2 数据流转路径
┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐ ┌─────────────────┐
│ ODS 层 │ │ DWD 层 │ │ DWS 层 │ │ ADS 层 │
│ 原始数据 │ ──▶ │ 明细数据 │ ──▶ │ 汇总数据 │ ──▶ │ 应用数据 │
│ │ │ │ │ │ │ │
│ • Append-only │ │ • 主键表 │ │ • 物化表 │ │ • 导出表 │
│ • 无主键 │ │ • 动态桶 │ │ • 低频刷新 │ │ • 不存湖内 │
│ • 原始日志 │ │ deduplicate │ │ • 聚合计算 │ │ • MySQL / ES │
│ • CDC → Paimon │ │ • 去重/更新 │ │ • 窗口统计 │ │ • API 服务 │
└─────────────────┘ └─────────────────┘ └─────────────────┘ └─────────────────┘
图 2:数仓分层与数据流转
04 Paimon 数仓分层表设计
4.1 ODS 层 — Append-only 表(无主键)
设计原则:ODS 层保留原始 CDC 数据,不做去重,不做更新。使用 Append-only 表,写入性能最优,存储成本最低。
-- ODS 层:Append-only 表,无主键 CREATE TABLE ods.mysql_orders ( order_id BIGINT, user_id BIGINT, product_name STRING, amount DECIMAL(10, 2), status STRING, order_time TIMESTAMP(3), op_ts TIMESTAMP(3), -- CDC 操作时间戳 op_type STRING, -- CDC 操作类型 (INSERT/UPDATE/DELETE) dt STRING -- 数据日期分区 ) PARTITIONED BY (dt) WITH ( 'bucket' = '8', 'file.format' = 'parquet', 'file.compression' = 'zstd', -- ODS/DWD 基础明细层不设置物理过期,防止历史数据回溯写入被丢弃 -- 数据生命周期管理交由 DWS 层查询过滤 + 冷热归档策略(见第 8 章) -- 快照保留策略 'snapshot.num-retained.min' = '5', 'snapshot.num-retained.max' = '20', 'snapshot.time-retained' = '7 d' );
写入方式:CDC Pipeline 直接写入 Paimon ODS 表(无需 Kafka)
-- CDC 3.6 直接写入 ODS Append-only 表 INSERT INTO ods.mysql_orders SELECT order_id, user_id, product_name, amount, status, order_time, op_ts, op_type, DATE_FORMAT(order_time, 'yyyy-MM-dd') AS dt FROM mysql_cdc_source /*+ OPTIONS('scan.startup.mode'='initial') */;
4.2 DWD 层 — 主键表(lookup changelog)
设计原则:DWD 层对 ODS 数据进行清洗、去重、关联维度,生成高质量明细数据。使用主键表 + 动态桶(Dedicated Bucket),主键仅保留业务唯一标识
order_id,不包含分区字段 dt,防止订单跨天更新导致主键去重失效。-- DWD 层:主键表,动态桶,支持去重和更新 CREATE TABLE dwd.order_detail ( order_id BIGINT, user_id BIGINT, product_name STRING, category STRING, unit_price DECIMAL(10, 2), quantity INT, total_amount DECIMAL(12, 2), status STRING, updated_at TIMESTAMP(3), dt STRING, PRIMARY KEY (order_id) NOT ENFORCED -- 主键仅含业务ID,不含分区字段 dt ) PARTITIONED BY (dt) -- 仍按天分区,但分区字段不在 PK 中 WITH ( 'bucket' = '-1', -- 动态桶模式,系统自动根据数据量扩容 'dynamic-bucket.target-row-num' = '2000000', -- 每个桶控制在 200 万行 'merge-engine' = 'deduplicate', 'changelog-producer' = 'lookup', 'lookup.cache-size' = '1000000', 'deletion-vectors.enabled' = 'true', 'write-buffer-size' = '256 MB' -- 不设置 partition.expiration-time,防止历史数据回溯丢失 );
写入方式:从 ODS Append-only 表流式读取,按主键去重写入 DWD
-- ODS → DWD:流式去重写入 INSERT INTO dwd.order_detail SELECT order_id, user_id, product_name, 'default' AS category, amount AS unit_price, 1 AS quantity, amount AS total_amount, status, order_time AS updated_at, DATE_FORMAT(order_time, 'yyyy-MM-dd') AS dt FROM ods.mysql_orders /*+ OPTIONS('read-mode'='log') */ WHERE op_type IN ('INSERT', 'UPDATE_AFTER');
4.3 DWS 层 — 物化表(低频刷新 + 资源倾斜防范)
设计原则:DWS 层存储预聚合的汇总指标,使用 Paimon 物化表自动管理刷新。为避免频繁的局部聚合引发分布式计算倾斜(特别是
COUNT(DISTINCT)),建议采用低频刷新策略。刷新频率建议:freshness设为'6 h'(每 6 小时)或深夜 T+1 调度。对于中小企业的 SLA 要求,日级或半日级汇总完全满足业务报表需求,且能大幅降低全天候 CPU 损耗。
-- DWS 层:物化表,低频刷新减少资源消耗 CREATE MATERIALIZED TABLE dws.daily_sales_summary PARTITIONED BY (dt) WITH ( 'pipeline-name' = 'daily_sales_pipeline', 'freshness' = '6 h', -- 每 6 小时刷新一次(替代原 1h),降幅 83% 'auto-refresh' = 'true' ) AS SELECT dt, product_name, COUNT(DISTINCT order_id) AS order_count, SUM(total_amount) AS total_sales, AVG(total_amount) AS avg_order_amount FROM dwd.order_detail GROUP BY dt, product_name;
4.4 ADS 层 — 导出表(不存湖内)
设计原则:ADS 层面向业务应用,数据直接导出到 MySQL/ES 等业务系统,不存储在 Paimon 湖内,避免湖内数据膨胀。
-- ADS 导出:Paimon → MySQL(不存湖内) INSERT INTO mysql_report.daily_top_products SELECT dt, product_name, order_count, total_sales FROM dws.daily_sales_summary WHERE dt = DATE_FORMAT(CURRENT_TIMESTAMP, 'yyyy-MM-dd') ORDER BY total_sales DESC LIMIT 100; -- ADS 导出:Paimon → Elasticsearch(实时搜索) INSERT INTO es_index.order_realtime SELECT order_id, user_id, total_amount, status, updated_at FROM dwd.order_detail /*+ OPTIONS('read-mode'='log') */;
4.5 分层设计总结
| 分层 | 表类型 | 主键 | Bucket | 存储位置 | 保留策略 |
|---|---|---|---|---|---|
| ODS | Append-only | 无 | 8 | 热存储 (HDD) | 不物理过期,由冷热归档管理 |
| DWD | Primary-key | order_id(不含 dt) | -1 动态桶 | 热存储 (HDD) | 不物理过期,由冷热归档管理 |
| DWS | Materialized | 有 | -1 动态桶 | 热存储 (HDD) | 查询层 30 天过滤 |
| ADS | 外部表 | - | - | MySQL / ES | 由业务系统管理 |
| 归档 | Clone 只读 | 有 | -1 动态桶 | 冷存储 (S3) | 长期保留 |
05 核心功能模块设计
5.1 数据链路:MySQL → CDC → Paimon(直连)
v2.0 采用奥卡姆剃刀原则精简架构,CDC 3.6 原生支持 Paimon Sink,内部基于 Flink Checkpoint 机制实现 Exactly-Once 断点续传,无需 Kafka 中间件:
MySQL Binlog → Flink CDC Pipeline → Paimon 表
(数据源) (采集+计算) (存储层)
直连架构的核心收益:
- 简化运维:省去 Kafka 集群的部署与维护,降低中小团队运维负担
- Schema 演进畅通:CDC → Paimon 直连,上游加列自动同步到 Paimon,无需担心中间层断流
- Exactly-Once 语义:CDC 3.6 基于 Flink Checkpoint 实现端到端精确一次,与 Kafka 的至少一次语义相比数据一致性更强
- 更低资源消耗:省去 Kafka 存储 I/O 和副本开销,CDC 作业无需同时维护两套 Flink 集群(采集 + 计算)
- 零代码入湖:YAML 中 sink 直接配置为 Paimon,一行配置即可完成整库同步与全量+增量自动切换
5.2 Catalog 管理:JDBC Catalog(带中心化锁)
v2.0 使用 JDBC Catalog 替代 Filesystem Catalog,引入中心化元数据存储和分布式锁机制:
-- JDBC Catalog(推荐,利用现有 MySQL 实例存储元数据) CREATE CATALOG paimon WITH ( 'type' = 'paimon', 'metastore' = 'jdbc', 'uri' = 'jdbc:mysql://localhost:3306/paimon_metadata', 'warehouse' = 'file:///data/paimon/warehouse', 'lock.enabled' = 'true', -- 开启 Catalog 级锁机制 'lock.flavor' = 'jdbc' -- 跨引擎多并发读写安全 );
JDBC Catalog 与 Filesystem Catalog / Hive Metastore 对比:
| 维度 | Filesystem Catalog | JDBC Catalog(本方案) | Hive Metastore |
|---|---|---|---|
| 外部依赖 | 无(仅文件系统) | MySQL 实例 | 需部署 HMS + MySQL |
| 元数据一致性 | ❌ 无锁机制,对象存储易冲突 | ✅ JDBC 锁保障跨引擎一致 | ✅ 中心化元数据 |
| 并发 Compaction | ❌ 高并发下快照元数据冲突 | ✅ 锁机制防冲突 | ✅ HMS 锁 |
| 多租户隔离 | 通过 warehouse 路径 | 通过 database + MySQL | 通过 database |
| 企业级权限 | ❌ 无 | ✅ 可基于 MySQL 权限控制 | ✅ Ranger/Sentry |
| 运维复杂度 | 极低 | 低(复用现有 MySQL) | 中 |
| 适用场景 | 单机/开发测试 | 中小团队生产环境 | 大规模多团队 |
Filesystem Catalog 局限:在对象存储(S3/MinIO)或非强一致性文件系统上,缺乏中心化锁机制。Flink 多个 TaskManager 在高并发 Compaction 或 Trino 并发读取时,易发生快照元数据冲突或不可见。JDBC Catalog 很好地解决了这一问题,且对中小团队而言只需复用现有 MySQL 实例即可。
5.3 SQL 开发平台:Apache Dinky
使用 Dinky 替代 Flink SQL CLI,降低 Flink SQL 维护成本:
Dinky 核心能力:
- 可视化 SQL 开发:智能代码补全、语法校验、在线调试
- CDC Pipeline 集成:通过
EXECUTE PIPELINE WITHYAML语法直接提交 CDC 任务 - 数据血缘分析:自动追踪 SQL 作业上下游依赖关系
- UDF 管理:UDF 注册、版本管理、跨作业复用
- 全局变量:支持环境变量、日期变量等跨作业共享
- K8s 集成:支持将任务提交到 Kubernetes 集群
Dinky 中提交 CDC Pipeline 任务:
-- 在 Dinky 中使用 Flink CDC YAML Pipeline 同步 MySQL 到 Paimon EXECUTE PIPELINE WITHYAML ( source: type: mysql hostname: 192.168.1.100 port: 3306 username: cdc_reader password: Cdc@Secure2026 tables: app_db\..* server-id: 5400-5408 server-time-zone: Asia/Shanghai sink: type: paimon name: Paimon Sink catalog.properties.metastore: jdbc catalog.properties.uri: jdbc:mysql://localhost:3306/paimon_metadata catalog.properties.warehouse: file:///data/paimon/warehouse catalog.properties.lock.enabled: true catalog.properties.lock.flavor: jdbc pipeline: name: MySQL-to-Paimon-CDC parallelism: 4 schema.change.behavior: lenient )
5.4 CDC 同步任务管理
整库同步(Whole-database Sync)
通过 YAML 的
tables 字段使用正则表达式实现整库同步。分库分表合并(Sharding Merge)
通过
route 模块实现分表路由合并,支持正则捕获组高级路由。Schema Evolution
Flink-CDC 3.6 提供五种 Schema 变更行为模式:
| 模式 | 行为 | 适用场景 |
|---|---|---|
exception | 禁止任何变更,直接报错 | 目标端不支持变更 |
evolve | 尝试应用所有变更,失败报错 | 目标端完全兼容 |
try_evolve | 尝试应用,失败容忍 | 部分兼容目标端 |
lenient | 转换变更确保不丢数据 | 通用场景(推荐) |
ignore | 忽略所有变更 | 只读历史数据 |
5.5 流批一体
- Streaming Read:通过
read-mode=log选项,持续消费新写入的数据变更 - Batch Read:直接读取当前快照的全量数据,适用于 T+1 报表
- Branch 隔离:通过
scan.fallback-branch实现流批读写分支隔离
5.6 数据更新(Upsert)与 Changelog
| 模式 | 说明 | 延迟 | 推荐场景 |
|---|---|---|---|
none | 不产生 changelog(默认) | - | 仅入湖,不下发 |
input | 直接透传输入 changelog | 最低 | 上游已有完整 changelog |
full-compaction | Full Compaction 时产生 | 较高 | 对延迟不敏感 |
lookup | 点查历史文件生成 changelog | 中等 | 需要完整 changelog(推荐) |
5.7 时间旅行(Time Travel)
- Snapshot 读取:通过
scan.snapshot-id查询任意历史快照 - Tag 标记:为重要时间点创建命名标签
- Branch 分支:从 Tag 创建分支进行独立读写,支持
fast_forward合并
06 核心 YAML 配置示例
6.1 MySQL CDC → Paimon(直连入湖)
# mysql-to-paimon-cdc.yaml # MySQL CDC 直接写入 Paimon(零中间件,Exactly-Once) source: type: mysql name: MySQL Source hostname: 192.168.1.100 port: 3306 username: cdc_reader password: Cdc@Secure2026 tables: app_db\..* server-id: 5400-5408 server-time-zone: Asia/Shanghai schema-change.enabled: true include-comments.enabled: true transform: - source-table: app_db\.users projection: > id, username, concat(left(phone, 3), '****', right(phone, 4)) as phone_masked, email, status, created_at filter: "status = 'ACTIVE'" description: 用户表脱敏 + 过滤 route: - source-table: app_db\.order_[0-9]+ sink-table: app_db.orders_merged description: 分表合并:order_0 ~ order_N → orders_merged sink: type: paimon name: Paimon Sink catalog.properties.metastore: jdbc catalog.properties.uri: jdbc:mysql://localhost:3306/paimon_metadata catalog.properties.warehouse: file:///data/paimon/warehouse catalog.properties.lock.enabled: true catalog.properties.lock.flavor: jdbc pipeline: name: MySQL-to-Paimon-CDC parallelism: 4 schema.change.behavior: lenient
6.2 PostgreSQL CDC → Paimon
# pg-to-paimon-cdc.yaml source: type: postgres name: PostgreSQL Source hostname: 192.168.1.101 port: 5432 username: cdc_reader password: Cdc@Secure2026 tables: public\..* schema-change.enabled: true decoding.plugin.name: pgoutput sink: type: paimon name: Paimon Sink catalog.properties.metastore: jdbc catalog.properties.uri: jdbc:mysql://localhost:3306/paimon_metadata catalog.properties.warehouse: file:///data/paimon/warehouse pipeline: name: PG-to-Paimon-CDC parallelism: 2 schema.change.behavior: evolve
6.3 Oracle CDC → Paimon(3.6 新增)
# oracle-to-paimon-cdc.yaml source: type: oracle name: Oracle Source hostname: 192.168.1.102 port: 1521 username: cdc_reader password: Cdc@Secure2026 tables: SCOTT\..* schema-change.enabled: true sink: type: paimon name: Paimon Sink catalog.properties.metastore: jdbc catalog.properties.uri: jdbc:mysql://localhost:3306/paimon_metadata catalog.properties.warehouse: file:///data/paimon/warehouse pipeline: name: Oracle-to-Paimon-CDC parallelism: 2 schema.change.behavior: lenient
6.4 提交 CDC 作业命令
# 方式一:CDC CLI 直接运行(测试) sh bin/flink-cdc.sh mysql-to-paimon-cdc.yaml # 方式二:提交到 Flink K8s 集群(生产) flink run \ --target kubernetes-application \ -Dkubernetes.cluster-id=mysql-cdc-pipeline \ -Dkubernetes.container.image=flink:2.2.0-cdc \ -c org.apache.flink.cdc.cli.CliFrontend \ lib/flink-cdc-dist-3.6.0.jar \ --config mysql-to-paimon-cdc.yaml # 方式三:通过 Dinky 提交(推荐) # 在 Dinky Web UI 中使用 EXECUTE PIPELINE WITHYAML 语法
07 关键 SQL 语句
7.1 创建 JDBC Catalog
-- 生产环境 JDBC Catalog(利用现有 MySQL 实例存储元数据) CREATE CATALOG paimon WITH ( 'type' = 'paimon', 'metastore' = 'jdbc', 'uri' = 'jdbc:mysql://localhost:3306/paimon_metadata', 'warehouse' = 'file:///data/paimon/warehouse', 'lock.enabled' = 'true', 'lock.flavor' = 'jdbc' ); USE CATALOG paimon;
7.2 ODS Append-only 表
CREATE TABLE ods.mysql_orders ( order_id BIGINT, user_id BIGINT, product_name STRING, amount DECIMAL(10, 2), status STRING, order_time TIMESTAMP(3), op_ts TIMESTAMP(3), op_type STRING, dt STRING ) PARTITIONED BY (dt) WITH ( 'bucket' = '8', 'file.format' = 'parquet', 'file.compression' = 'zstd', -- ODS/DWD 基础明细层不设置物理过期,防止历史数据回溯写入被丢弃 'snapshot.num-retained.min' = '5', 'snapshot.num-retained.max' = '20', 'snapshot.time-retained' = '7 d' );
7.3 DWD 主键表(lookup changelog)
CREATE TABLE dwd.order_detail ( order_id BIGINT, user_id BIGINT, product_name STRING, category STRING, unit_price DECIMAL(10, 2), quantity INT, total_amount DECIMAL(12, 2), status STRING, updated_at TIMESTAMP(3), dt STRING, PRIMARY KEY (order_id) NOT ENFORCED -- 主键仅含order_id,不含dt,防止跨天更新失效 ) PARTITIONED BY (dt) WITH ( 'bucket' = '-1', -- 动态桶模式 'dynamic-bucket.target-row-num' = '2000000', 'merge-engine' = 'deduplicate', 'changelog-producer' = 'lookup', 'lookup.cache-size' = '1000000', 'deletion-vectors.enabled' = 'true', 'write-buffer-size' = '256 MB' -- 基础明细层不设置partition.expiration-time,由冷热分离策略管理生命周期 );
7.4 DWS 物化表(低频刷新)
CREATE MATERIALIZED TABLE dws.daily_sales_summary PARTITIONED BY (dt) WITH ( 'pipeline-name' = 'daily_sales_pipeline', 'freshness' = '6 h', -- 每 6 小时刷新,降低频繁聚合导致的资源倾斜 'auto-refresh' = 'true' ) AS SELECT dt, product_name, COUNT(DISTINCT order_id) AS order_count, SUM(total_amount) AS total_sales, AVG(total_amount) AS avg_order_amount FROM dwd.order_detail GROUP BY dt, product_name;
7.5 CDC Pipeline 直接写入 Paimon
CDC 3.6 直接写入 Paimon 表,无需编写 Kafka 消费 SQL,Streaming SQL 变得简洁直观:
-- CDC 3.6 直接写入 ODS Append-only 表(无需 Kafka) INSERT INTO ods.mysql_orders SELECT order_id, user_id, product_name, amount, status, order_time, op_ts, op_type, DATE_FORMAT(order_time, 'yyyy-MM-dd') AS dt FROM mysql_cdc_source /*+ OPTIONS('scan.startup.mode'='initial') */;
Schema 演进优势:MySQL 上游新增列后,CDC Pipeline 的
lenient 模式自动将新字段写入 Paimon。
下游 Flink SQL 的 INSERT INTO ... SELECT * 会自动感知新列,无需人工修改 SQL 中的 JSON_VALUE 解析逻辑。端到端数据一致性:
- 依赖 Flink Checkpoint 机制实现 Exactly-Once 语义
- 故障恢复时从最近 Checkpoint 断点续传,不会重复或丢失数据
- 无需 Kafka 的至少一次语义和下游去重逻辑
7.6 时间旅行查询
-- 查询指定 Snapshot ID 的历史数据 SELECT * FROM dwd.order_detail /*+ OPTIONS('scan.snapshot-id'='5') */; -- 创建 Tag CALL sys.create_tag('dwd.order_detail', 'tag_20260620', 5); -- 创建 Branch CALL sys.create_branch('dwd.order_detail', 'branch_test', 'tag_20260620'); -- 读取分支数据 SELECT * FROM `dwd.order_detail$branch_branch_test`; -- 合并分支 CALL sys.fast_forward('dwd.order_detail', 'branch_test');
08 冷热数据分层策略
8.1 分层架构
┌──────────────────────────────────────────────────────────────────┐
│ 数据生命周期管理 │
│ │
│ 写入 ──▶ 热存储 (HDD/SSD) ──▶ 30天 ──▶ 冷存储 (S3/MinIO) │
│ 近期数据 历史归档 │
│ │
│ • 频繁读写 • 只读查询 │
│ • 本地低延迟 • S3 低成本 │
│ • 分区过期自动清理 • Clone 归档 │
│ │
│ 成本:约 0.10 元/GB/月 成本:约 0.024 元/GB/月 │
│ │
│ 预计节省存储成本:约 30%(历史数据占总量 70%+ 场景) │
└──────────────────────────────────────────────────────────────────┘
8.2 冷存储 Catalog 配置
-- 热存储 Catalog(本地 HDD/SSD,JDBC 元数据 + 本地存储) CREATE CATALOG paimon_hot WITH ( 'type' = 'paimon', 'metastore' = 'jdbc', 'uri' = 'jdbc:mysql://localhost:3306/paimon_metadata', 'warehouse' = 'file:///data/paimon/warehouse', 'lock.enabled' = 'true', 'lock.flavor' = 'jdbc' ); -- 冷存储 Catalog(S3 / MinIO) CREATE CATALOG paimon_cold WITH ( 'type' = 'paimon', 'metastore' = 'jdbc', 'uri' = 'jdbc:mysql://localhost:3306/paimon_metadata', 'warehouse' = 's3://data-lake-cold/paimon-archive/', 's3.endpoint' = 'https://s3.amazonaws.com', 's3.access-key' = '${S3_ACCESS_KEY}', 's3.secret-key' = '${S3_SECRET_KEY}', 'lock.enabled' = 'true', 'lock.flavor' = 'jdbc' );
8.3 归档流程:Clone 到冷存储
-- 步骤 1:为热表创建 Tag(标记当前快照) CALL sys.create_tag( `table` => 'dwd.order_detail', tag => 'archive_20260520', time_retained => '365 d' ); -- 步骤 2:Clone 到冷存储 CALL sys.clone( warehouse => 'file:///data/paimon/warehouse', database => 'dwd', table => 'order_detail', target_warehouse => 's3://data-lake-cold/paimon-archive/', target_database => 'archive', target_table => 'order_detail' ); -- 步骤 3:热存储分区自动过期(由 partition.expiration-time 控制) -- 30 天前的分区自动删除,释放热存储空间
8.4 查询冷存储数据
-- 切换到冷存储 Catalog 查询历史数据 USE CATALOG paimon_cold; SELECT * FROM archive.order_detail WHERE dt BETWEEN '2025-01-01' AND '2025-12-31';
8.5 成本对比
| 存储类型 | 单价(参考) | 数据占比 | 月成本占比 |
|---|---|---|---|
| 本地 HDD(热) | ~0.10 元/GB/月 | 30%(近期数据) | ~60% |
| S3(冷) | ~0.024 元/GB/月 | 70%(历史数据) | ~25% |
| 合计 | - | 100% | ~85%(节省 15%) |
实际节省比例取决于冷热数据比例。对于历史数据占比 70%+ 的场景,存储成本可降低 25-40%。
09 业务侧表结构变更处理方案
9.1 整体流程
业务侧 DDL 变更 → Flink-CDC 捕获 → Paimon 自动演进 → 下游分层影响评估
(MySQL/PG) (schema-change) (lenient/evolve) (DWD/DWS/ADS)
9.2 CDC 采集层:自动捕获
Flink-CDC 3.6 的 MySQL/PostgreSQL Source 默认开启
schema-change.enabled: true,会自动捕获以下 DDL 事件:| 变更类型 | CDC 是否捕获 | 说明 |
|---|---|---|
| ADD COLUMN | 是 | 新增列自动同步 |
| CHANGE COLUMN TYPE | 是 | 类型变更(需向后兼容) |
| DROP COLUMN | 是 | 删除列事件 |
| RENAME COLUMN | 是 | 重命名列事件 |
| MODIFY COLUMN COMMENT | 是 | 需开启 include-comments.enabled: true |
| ADD INDEX | 否 | 索引变更不捕获 |
| DROP TABLE | 是 | 整表删除事件 |
9.3 五种 Schema 变更行为模式
当前方案使用
lenient 模式,各模式对业务侧变更的处理差异:| 模式 | ADD COLUMN | DROP COLUMN | CHANGE TYPE | RENAME COLUMN |
|---|---|---|---|---|
exception | ❌ 报错停止 | ❌ 报错停止 | ❌ 报错停止 | ❌ 报错停止 |
evolve | ✅ 自动加列 | ✅ 自动删列 | ✅ 兼容则成功 | ✅ 自动重命名 |
try_evolve | ✅ 尝试加列 | ⚠️ 失败则忽略 | ⚠️ 失败则忽略 | ⚠️ 失败则忽略 |
lenient | ✅ 自动加列 | ✅ 保留列不删 | ⚠️ 兼容则成功,否则保留 | ⚠️ 视为删旧+加新 |
ignore | ⚠️ 完全忽略 | ⚠️ 完全忽略 | ⚠️ 完全忽略 | ⚠️ 完全忽略 |
lenient 模式的核心价值:保证数据不丢失。DROP COLUMN 不会删除 Paimon 列,不兼容的类型变更不会报错。
9.4 各变更类型的详细处理
新增字段(ADD COLUMN)— 最常见,风险最低
业务侧: ALTER TABLE orders ADD COLUMN remark VARCHAR(255);
CDC 捕获: ✅ 自动捕获 ADD COLUMN 事件
Paimon ODS: ✅ Append-only 表自动加列,旧数据该字段为 NULL
Paimon DWD: ✅ 主键表自动加列,旧数据该字段为 NULL
Paimon DWS: ⚠️ 物化表不会自动加列,需要手动重建
下游 ADS: ⚠️ 导出目标表需要手动加列
无需人工干预,CDC + Paimon 自动完成,直连架构保障 Schema 演进不会在中间层(如 Kafka)断流。ODS/DWD 层新字段自动可用,Trino 查询也能直接看到新字段。
修改字段类型(CHANGE COLUMN TYPE)
业务侧: ALTER TABLE orders MODIFY amount DECIMAL(18,2); -- 从 DECIMAL(10,2) 扩宽
CDC 捕获: ✅ 自动捕获
Paimon 处理:
- 向后兼容(扩宽): ✅ VARCHAR(100) → VARCHAR(255)、INT → BIGINT → 自动成功
- 向前不兼容(缩窄): ⚠️ DECIMAL(18,2) → DECIMAL(10,2) 可能失败
推荐做法:业务侧只做向后兼容的类型变更(扩宽、精度增加)。如果必须缩窄,走灰度流程(见 9.6 节)。
删除字段(DROP COLUMN)— 有数据丢失风险
业务侧: ALTER TABLE orders DROP COLUMN temp_flag;
lenient 模式处理:
- Paimon 表不会删除该列 ✅(保证不丢数据)
- 后续新数据该字段写入 NULL
- 历史数据该字段值保持不变
重命名字段(RENAME COLUMN)— 需要灰度处理
业务侧: ALTER TABLE orders RENAME COLUMN amount TO total_amount;
lenient 模式处理: 视为"删旧列 + 加新列"
- amount 列保留(历史数据不变,新数据为 NULL)
- total_amount 列新增(新数据有值,历史数据为 NULL)
问题:数据被拆成两列。推荐走灰度流程(见 9.6 节)。
9.5 下游分层影响与处理
| 分层 | 是否自动演进 | 需要人工操作 | 说明 |
|---|---|---|---|
| ODS (Append-only) | ✅ 自动 | 无 | 新增列自动追加,旧数据为 NULL |
| DWD (主键表) | ✅ 自动 | 无 | 新增列自动追加 |
| DWS (物化表) | ❌ 不自动 | 需重建物化表 | 物化表定义 SQL 不变,新列不会自动加入聚合 |
| ADS (导出表) | ❌ 不自动 | 需手动加列 | MySQL/ES 目标表需要先执行 DDL |
DWS 物化表更新流程:
-- 步骤 1:删除旧物化表 DROP MATERIALIZED TABLE dws.daily_sales_summary; -- 步骤 2:重建物化表(包含新字段) CREATE MATERIALIZED TABLE dws.daily_sales_summary PARTITIONED BY (dt) WITH ( 'pipeline-name' = 'daily_sales_pipeline', 'freshness' = '1 h', 'auto-refresh' = 'true' ) AS SELECT dt, product_name, COUNT(DISTINCT order_id) AS order_count, SUM(total_amount) AS total_sales, AVG(total_amount) AS avg_order_amount, -- 新增字段 COUNT(*) AS total_records FROM dwd.order_detail GROUP BY dt, product_name;
9.6 灰度变更标准流程(推荐)
对于高风险变更(删列、缩窄类型、重命名),建议走完整灰度流程:
┌─────────────────────────────────────────────────────────────────┐
│ Schema 变更灰度流程 │
│ │
│ 1. 评估影响 │
│ ├── 确认变更类型(加列/改类型/删列/重命名) │
│ ├── 评估下游依赖(哪些 DWD/DWS/ADS 依赖该字段) │
│ └── 确认是否向后兼容 │
│ │
│ 2. 通知下游 │
│ ├── 通知数据分析师、BI 报表负责人 │
│ └── 通知下游应用开发(ADS 导出目标) │
│ │
│ 3. 执行变更(加列优先策略) │
│ ├── 先 ADD COLUMN 新列(不删旧列) │
│ ├── 业务代码双写(新旧列同时写入) │
│ ├── 等待 CDC 同步完成(Paimon 新列有数据) │
│ ├── 下游 DWD/DWS/ADS 切换到新列 │
│ └── 确认无问题后,再 DROP 旧列 │
│ │
│ 4. 验证 │
│ ├── 对比新旧列数据一致性 │
│ ├── 检查 DWS 物化表是否需要重建 │
│ └── 检查 ADS 导出目标表结构 │
│ │
│ 5. 清理 │
│ ├── 删除旧列(如需要) │
│ └── 更新 Dinky 中的 SQL 作业 │
└─────────────────────────────────────────────────────────────────┘
重命名字段灰度示例:
-- 步骤 1:新增目标列 ALTER TABLE orders ADD COLUMN total_amount DECIMAL(10, 2); -- 步骤 2:业务代码双写 amount + total_amount,跑 1-2 天 -- 步骤 3:CDC 自动同步,Paimon 两列都有完整数据 -- 步骤 4:下游 DWD/DWS 改为读取 total_amount 列 -- 更新 Dinky 中的 Flink SQL 作业 -- 步骤 5:业务代码切到只写 total_amount -- 步骤 6:确认无问题后,删除旧列 ALTER TABLE orders DROP COLUMN amount;
9.7 Schema 变更监控与告警
| 监控项 | 方式 | 告警条件 |
|---|---|---|
| DDL 事件捕获 | Flink CDC metrics | 收到 DDL 事件时发送通知 |
| Paimon Schema 变更 | Paimon table metrics | 表结构发生变化时告警 |
| 物化表失效 | Dinky 任务监控 | 上游表结构变更导致物化表报错 |
告警通知渠道:钉钉/飞书机器人,DDL 变更时自动推送变更详情(源表、变更类型、影响字段),便于数据团队及时评估下游影响。
10 数据安全与恢复方案
10.1 问题一:Binlog 过期导致全量重刷
场景:MySQL binlog 保留时间不足,CDC 作业断开后无法从断点续传,需要全量重刷。
解决方案:基于 Tag + Clone 的增量重同步
⚠️ 注意:传统MERGE INTO全量对撞在 Flink 批模式下会触发 Row-level Join State 大爆炸,千万级表极易 OOM。推荐采用sys.clone的时间旅行机制替代全量 MERGE。
-- 步骤 1:停止当前 CDC 作业,记录当前快照 ID -- 步骤 2:从 MySQL 拉取全量快照,写入临时表 CREATE TABLE mysql_full_load ( order_id BIGINT, user_id BIGINT, product_name STRING, amount DECIMAL(10, 2), status STRING, order_time TIMESTAMP(3), dt STRING, PRIMARY KEY (order_id) NOT ENFORCED ) WITH ( 'connector' = 'mysql-cdc', 'hostname' = '192.168.1.100', 'port' = '3306', 'username' = 'cdc_reader', 'password' = 'Cdc@Secure2026', 'database-name' = 'app_db', 'table-name' = 'orders', 'scan.startup.mode' = 'initial', 'scan.incremental.snapshot.enabled' = 'false' ); -- 步骤 3:写入临时 Paimon 表 INSERT INTO dwd.order_detail_full SELECT * FROM mysql_full_load; -- 步骤 4:创建 Tag 并调用 sys.clone 覆盖目标表 CALL sys.clone( warehouse => 'file:///data/paimon/warehouse', database => 'dwd', table => 'order_detail_full', target_warehouse => 'file:///data/paimon/warehouse', target_database => 'dwd', target_table => 'order_detail' ); -- 步骤 5:重新启动 CDC 增量同步作业
如果数据量较小(百万级以下),也可使用分区级
MERGE INTO 缩小状态范围:-- 逐分区 MERGE,减少单次状态量 CALL sys.merge_into( target_table => 'dwd.order_detail', source_table => 'mysql_full_load', merge_condition => 'dwd.order_detail.order_id = mysql_full_load.order_id', matched_upsert_setting => '*', not_matched_insert_values => '*' );
预防措施:
- MySQL binlog 保留时间至少 72 小时
- 配置 Flink Checkpoint 间隔 5-10 分钟,配合 CDC 3.6 的断点续传机制
- 监控 CDC Source 延迟,提前告警
- 大表优先使用 Clone + Tag 机制,避免全量 MERGE
10.2 问题二:Paimon 表误删恢复
场景:误执行
DROP TABLE 导致表被删除。预防措施:定期 Clone 备份
-- 定时任务:每日创建 Tag + Clone 到备份仓库 -- 每日凌晨执行 CALL sys.create_tag( `table` => 'dwd.order_detail', tag => 'daily_backup_20260620', time_retained => '365 d' ); CALL sys.clone( warehouse => 'file:///data/paimon/warehouse', database => 'dwd', table => 'order_detail', target_warehouse => 'file:///data/paimon/backup/', target_database => 'backup', target_table => 'order_detail' );
恢复步骤:
-- 方法 1:从备份 Clone 恢复 CALL sys.clone( warehouse => 'file:///data/paimon/backup/', database => 'backup', table => 'order_detail', target_warehouse => 'file:///data/paimon/warehouse/', target_database => 'dwd', target_table => 'order_detail' ); -- 方法 2:如果数据文件未被 purge,重新创建表结构 -- Paimon 会自动发现文件系统中的数据文件 CREATE TABLE dwd.order_detail ( -- 重建原表结构 order_id BIGINT, user_id BIGINT, product_name STRING, category STRING, unit_price DECIMAL(10, 2), quantity INT, total_amount DECIMAL(12, 2), status STRING, updated_at TIMESTAMP(3), dt STRING, PRIMARY KEY (order_id) NOT ENFORCED ) PARTITIONED BY (dt);
10.3 问题三:Paimon 表误更新恢复
场景:误操作导致 Paimon 表数据被错误更新。
解决方案:快照回滚
-- 方法 1:回滚到指定快照 ID CALL sys.rollback_to( `table` => 'dwd.order_detail', snapshot_id => 15 ); -- 方法 2:回滚到指定 Tag CALL sys.rollback_to( `table` => 'dwd.order_detail', tag => 'before_bad_update' ); -- 方法 3:回滚到指定时间戳 CALL sys.rollback_to_timestamp( `table` => 'dwd.order_detail', timestamp => 1718841600000 ); -- 回滚后清理旧快照 CALL sys.expire_snapshots( `table` => 'dwd.order_detail', retain_max => 5, older_than => '2026-06-20 00:00:00' );
注意:回滚操作不会删除数据文件,仅将快照指针指向历史版本。需配合expire_snapshots清理回滚后的废弃文件。
10.4 问题四:上游误 DROP TABLE 后的恢复
场景:MySQL 上游误执行
DROP TABLE,CDC 作业报错。恢复路径:
┌──────────────────────────────────────────────────────────────────┐
│ 上游误 DROP TABLE 恢复流程 │
│ │
│ 1. MySQL 恢复 │
│ ├── 从 MySQL 备份恢复表结构 + 数据 │
│ └── 或使用 MySQL flashback(如开启 recyclebin) │
│ │
│ 2. Paimon 侧处理 │
│ ├── CDC 作业会收到 DROP 事件(schema-change.enabled=true) │
│ ├── lenient 模式下:Paimon 不会删除表,仅忽略该事件 │
│ └── Paimon 中历史数据完整保留 │
│ │
│ 3. 重新同步 │
│ ├── MySQL 表恢复后,CDC 自动捕获新数据 │
│ ├── 如需全量重刷,参考 9.1 MERGE INTO 方案 │
│ └── Paimon 主键表自动去重,不会产生重复数据 │
│ │
│ 4. 验证 │
│ ├── 对比 Paimon 数据与 MySQL 恢复后数据 │
│ └── 确认数据一致性后恢复正常 CDC 同步 │
└──────────────────────────────────────────────────────────────────┘
关键配置确保安全:
# CDC Pipeline 中使用 lenient 模式 pipeline: schema.change.behavior: lenient # DROP TABLE 不会删除 Paimon 表
10.5 数据安全策略总结
| 风险场景 | 预防措施 | 恢复方案 |
|---|---|---|
| Binlog 过期 | binlog 保留 72h+,Checkpoint 3min | MERGE INTO 分区级重同步 |
| Paimon 表误删 | 每日 Tag + Clone 备份 | 从备份 Clone 恢复 |
| Paimon 数据误更新 | 关键操作前创建 Tag | rollback_to 快照回滚 |
| 上游误 DROP TABLE | lenient 模式 + MySQL 备份 | MySQL 恢复 + CDC 重同步 |
| CDC 作业故障 | Flink Checkpoint Exactly-Once 断点续传 | 从 Checkpoint 或 Savepoint 恢复 |
11 数据血缘与监控告警
11.1 元数据管理
使用 JDBC Catalog 后,元数据统一存储在 MySQL 实例中,提供中心化的 Schema 版本管理和快照指针存储。Trino 通过 Paimon Connector 连接 JDBC Catalog 读取元数据,借助 JDBC 锁机制保障跨引擎并发读写的一致性。
11.2 数据血缘追踪
通过 Dinky 内置血缘分析 实现数据血缘管理:
- CDC 入湖血缘:Dinky 自动解析
EXECUTE PIPELINE WITHYAML的 source → sink 映射 - 数仓分层血缘:Dinky 自动解析 Flink SQL 的 INSERT INTO ... SELECT 上下游依赖
- 物化表血缘:Materialized Table 的定义 SQL 自动记录数据来源
- 可视化展示:Dinky Web UI 提供血缘关系图可视化
11.3 监控告警体系
┌──────────────────┐ ┌──────────────────┐
│ Flink Metrics │ │ Paimon Metrics │
│ (Prometheus) │─────────────────▶│ (Table/Compact) │
└──────────────────┘ └────────┬─────────┘
│
▼
┌──────────────────┐
│ Prometheus │
└────────┬─────────┘
│
┌──────────────┼──────────────┐
▼ ▼ ▼
┌────────────┐ ┌──────────┐ ┌──────────┐
│ Grafana │ │AlertMgr │ │ │
│ 可视化 │ │ │ │ │
└────────────┘ └────┬─────┘ └──────────┘
│
┌────────────┼────────────┐
▼ ▼ ▼
┌────────┐ ┌────────┐ ┌──────┐
│ 钉钉 │ │ 飞书 │ │ 邮件 │
└────────┘ └────────┘ └──────┘
关键监控指标
| 维度 | 指标 | 告警阈值 |
|---|---|---|
| CDC 端到端延迟 | MySQL → Paimon 端到端延迟 | > 5 分钟 |
| Checkpoint | Checkpoint 完成时间 | > 10 分钟 |
| Backpressure | 算子反压比率 | > 50% 持续 5 分钟 |
| Compaction | Paimon Compaction 积压 | 待合并文件数 > 100 |
| 小文件 | 单分区文件数 | > 200 个 |
| 存储增长 | Warehouse 存储日增量 | > 500GB/天 |
| Binlog 保留 | MySQL binlog 最小剩余时间 | < 24 小时 |
12 部署与运维最佳实践
12.1 Kubernetes 部署架构(替代 YARN)
┌─────────────────────────────────────────────────────────────────┐
│ Kubernetes Cluster │
│ │
│ ┌─── Namespace: flink ──────────────────────────────────────┐ │
│ │ │ │
│ │ ┌─── Flink Kubernetes Operator ────────────────────────┐ │ │
│ │ │ • 管理 FlinkDeployment 生命周期 │ │ │
│ │ │ • Job Autoscaler 自动扩缩容 │ │ │
│ │ │ • Savepoint / Checkpoint 管理 │ │ │
│ │ └──────────────────────────────────────────────────────┘ │ │
│ │ │ │
│ │ ┌─── Flink Session Cluster ────────────────────────────┐ │ │
│ │ │ JobManager (1-3 replicas, HA) │ │ │
│ │ │ TaskManager (HPA: 2-10 pods, 弹性伸缩) │ │ │
│ │ └──────────────────────────────────────────────────────┘ │ │
│ │ │ │
│ │ ┌─── Dinky ────────────────────────────────────────────┐ │ │
│ │ │ Deployment (1-2 replicas) │ │ │
│ │ │ • SQL 开发 Web UI │ │ │
│ │ │ • 任务提交网关 │ │ │
│ │ └──────────────────────────────────────────────────────┘ │ │
│ │ │ │
│ └───────────────────────────────────────────────────────────┘ │
│ │
│ ┌─── Namespace: middleware ──────────────────────────────────┐│
│ │ MinIO / S3 (冷存储 / Checkpoint) ││
│ │ Prometheus + Grafana ││
│ │ Trino (1-2 workers) ││
│ │ MySQL (Paimon JDBC Catalog 元数据库) ││
│ └─────────────────────────────────────────────────────────────┘│
└─────────────────────────────────────────────────────────────────┘
12.2 Flink Kubernetes Operator 配置
# FlinkDeployment CR - CDC 采集作业 apiVersion: flink.apache.org/v1beta1 kind: FlinkDeployment metadata: name: mysql-cdc-pipeline namespace: flink spec: flinkVersion: v2_2 image: custom-flink:2.2.0-cdc serviceAccount: flink-service-account flinkConfiguration: taskmanager.numberOfTaskSlots: "4" state.backend: rocksdb state.backend.rocksdb.localdir: /data/flink/rocksdb state.backend.rocksdb.block.cache-size: 2048mb # 限制 RocksDB 块缓存,防止 OOM state.checkpoints.dir: s3://flink-checkpoints/ state.savepoints.dir: s3://flink-savepoints/ execution.checkpointing.interval: 5min # 放宽至 5min,减轻存储底座 I/O 压力 execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.timeout: 10min restart-strategy: fixed-delay restart-strategy.fixed-delay.attempts: "3" restart-strategy.fixed-delay.delay: 30s metrics.reporter.prom.factory.class: \ org.apache.flink.metrics.prometheus.PrometheusReporterFactory metrics.reporter.prom.port: "9250-9260" jobManager: resource: memory: "4096m" cpu: 2 replicas: 1 taskManager: resource: memory: "8192m" cpu: 4 job: jarURI: local:///opt/flink/usrlib/flink-cdc-dist-3.6.0.jar entryClass: org.apache.flink.cdc.cli.CliFrontend args: - "--config" - "/opt/flink/conf/mysql-to-paimon-cdc.yaml"
12.3 弹性伸缩配置
# FlinkDeployment 中配置 Autoscaler spec: job: autoscaler: enabled: true # 基于 CDC backlog 和 CPU 利用率自动扩缩 metrics.window: "5min" # 目标 lag 阈值 targets: - type: backlog value: 10000 scaleFactor: 2.0 - type: utilization value: 0.7 scaleFactor: 2.0 # 并行度范围 parallelism: min: 2 max: 16
12.4 硬件资源规划(K8s 环境)
| 角色 | Pod 数 | CPU/内存 | 存储 | 说明 |
|---|---|---|---|---|
| Flink JobManager | 1-3 (HA) | 2C / 4GB | 100GB SSD | 主节点 |
| Flink TaskManager | 2-6 (弹性) | 4C / 8GB | 500GB SSD | 计算节点 |
| Dinky | 1-2 | 2C / 4GB | 50GB | SQL 开发平台 |
| MinIO | 4 | 2C / 4GB | 2TB+ HDD | 冷存储 / Checkpoint |
| MySQL(元数据) | 1 | 2C / 4GB | 100GB SSD | Paimon JDBC Catalog |
| Trino | 1-2 | 4C / 8GB | 100GB | OLAP 查询 |
| Prometheus + Grafana | 1 | 2C / 4GB | 500GB | 监控 |
12.5 Flink 关键配置
# config.yaml - Flink 2.2 核心配置(K8s 环境) # Checkpoint 策略:对秒级入湖的平台,Checkpoint 仅用于恢复恢复 # 5-10 分钟间隔完全满足中小企业 SLA,且能大幅减轻存储底座 I/O 压力 execution.checkpointing.interval: 5min # 原 3min → 5min,降低高频增量 Checkpoint 带来的 I/O 飙高 execution.checkpointing.mode: EXACTLY_ONCE execution.checkpointing.timeout: 10min execution.checkpointing.min-pause: 2min execution.checkpointing.storage: filesystem # RocksDB 状态后端:显式限制内存,配合 Paimon 频繁的 LSM-Tree Compaction 时防止 OOM state.backend: rocksdb state.backend.rocksdb.localdir: /data/flink/rocksdb state.backend.rocksdb.block.cache-size: 2048mb # 显式限制 RocksDB 块缓存 state.checkpoints.dir: s3://flink-checkpoints/ state.savepoints.dir: s3://flink-savepoints/ restart-strategy: fixed-delay restart-strategy.fixed-delay.attempts: 3 restart-strategy.fixed-delay.delay: 30s taskmanager.memory.process.size: 8192mb taskmanager.numberOfTaskSlots: 4 parallelism.default: 4 # K8s HA 配置 high-availability.type: kubernetes high-availability.storageDir: s3://flink/ha/ kubernetes.cluster-id: flink-production
12.6 Paimon 表设计最佳实践
动态桶选择:主键表推荐使用bucket = '-1'动态桶模式,由系统根据数据量自动扩缩容,避免人工估算 bucket 数量。如需固定桶,bucket 数建议为 TaskManager Slot 总数的整数倍,单个 Bucket 数据量控制在 1-10GB 为宜。
- 分区策略:按天分区(
dt),避免过度分区 - 主键设计:主键仅含业务唯一标识(如
order_id),不包含分区字段(如dt),防止跨分区更新导致主键去重失效 - 文件格式:Paimon 1.0 默认 Parquet + ZSTD,无需额外配置
- Compaction 调优:开启
compaction.optimized-compactions-only - Deletion Vectors:SSD 环境强烈建议开启
- Catalog 选择:生产环境优先使用 JDBC Catalog 替代 Filesystem Catalog,启用
lock.enabled=true确保跨引擎并发安全 - 生命周期管理:基础明细层(ODS/DWD)不配置
partition.expiration-time,通过冷热归档策略(sys.clone 到 S3)+ 查询层过滤实现降本
12.7 CDC 运维要点
- Server-ID 管理:每个 MySQL CDC 作业分配独立的 server-id 范围
- Binlog 保留:至少 72 小时,防止全量同步期间 binlog 被清理
- Savepoint 管理:通过 K8s Operator 自动管理 Savepoint
- Schema 变更流程:lenient 模式自动同步,无需人工干预
- 监控告警:重点关注 Checkpoint 失败、CDC 端到端延迟、Binlog 解析异常
12.8 运维自动化脚本
#!/bin/bash # cdc-ops.sh - CDC 作业运维脚本(K8s 环境) # 列出所有 Flink 作业 list_jobs() { kubectl get flinkdeployment -n flink -o wide } # 触发 Savepoint savepoint_job() { local job_name=$1 kubectl create job manual-savepoint \ --from=flinkdeployment/$job_name \ -n flink \ --image=apache/flink:2.2.0 \ -- flink savepoint $job_name } # 更新 CDC Pipeline 配置(滚动升级) update_cdc_pipeline() { local job_name=$1 local yaml_path=$2 # 更新 ConfigMap 中的 YAML 配置 kubectl create configmap cdc-config-$job_name \ --from-file=pipeline.yaml=$yaml_path \ --dry-run=client -o yaml | kubectl apply -f - # 触发滚动重启 kubectl rollout restart flinkdeployment/$job_name -n flink } # 查看 Pod 状态 check_pods() { kubectl get pods -n flink -l app=flink -o wide } # 用法示例 list_jobs savepoint_job mysql-cdc-pipeline update_cdc_pipeline mysql-cdc-pipeline /path/to/mysql-to-paimon-cdc.yaml check_pods
12.9 成本估算(K8s 环境,中小公司参考)
| 项目 | 规格 | 数量 | 月成本估算 |
|---|---|---|---|
| K8s Worker 节点 | 16C32G + 1TB HDD | 3 台 | 约 4,500-9,000 元 |
| 冷存储 (S3) | 按量付费 | 5TB | 约 120 元 |
| 人力维护 | 大数据工程师 | 1-2 人 | - |
| 月总成本估算 | 约 4,620-9,120 元 |
对比 v1.0(YARN + HMS + Kafka):去掉了 HMS 依赖和 Kafka 集群(节省 3 台 4C8G 节点和 3TB HDD 存储),冷数据迁移到 S3 降低存储成本,K8s 弹性伸缩减少资源浪费。引入 JDBC Catalog 仅需复用现有 MySQL 实例,新增一个数据库即可。综合节省约 30-40%。
13 附录:版本兼容性矩阵
| 组件 | 推荐版本 | 最低版本 | JDK 要求 | 兼容说明 |
|---|---|---|---|---|
| Apache Flink | 2.2.1 | 2.0.0 | JDK 11+(推荐 17) | 2.0 起移除 Java 8 和 Per-Job 模式 |
| Apache Paimon | 1.0.x | 0.9.0 | JDK 11+ | 1.0 正式版,推荐生产使用 |
| Flink-CDC | 3.6.0 | 3.0.0 | JDK 11+ | 同时兼容 Flink 1.20.x 和 2.2.x |
| Apache Dinky | 1.2.5 | 1.2.0 | JDK 11+ | 支持 Flink 1.14~1.20,K8s Operator 集成 |
| Flink K8s Operator | 1.15 | 1.10 | - | 支持 Flink v1.19 到 v2.2 |
| MinIO | RELEASE.2024+ | - | - | S3 兼容冷存储 / Checkpoint |
| Trino | 435+ | 425+ | JDK 17+ | 需安装 Paimon Connector |
| MySQL | 5.7 / 8.0+ | 5.6 | - | CDC Source 需要 binlog 开启 |
| PostgreSQL | 14+ | 10+ | - | 需配置 WAL logical replication |
| Oracle | 19c+ | 12c | - | 3.6 新增 Pipeline Source |
参考资料
- Apache Flink 官方文档 - Flink 2.0 Release Notes
- Apache Paimon 1.0 发布公告
- Flink CDC 3.6.0 发布公告
- Apache Dinky 官方文档
- Flink Kubernetes Operator 文档
- Flink CDC 3.6 官方文档 - MySQL Pipeline Connector
- Flink CDC 3.6 官方文档 - Paimon Sink Connector
- Flink CDC 3.6 官方文档 - Route 模块
- Flink CDC 3.6 官方文档 - Transform 模块
- Apache Paimon 1.0 Procedures - Clone / Rollback