L刘志敏

基于 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 项目背景与建设目标

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 Flink2.2.x流批一体计算、SQL 执行引擎
存储层Apache Paimon1.0.x湖仓一体存储、ACID 事务
数据采集Flink-CDC3.6.xCDC 变更捕获、YAML Pipeline 编排
SQL 开发平台Apache Dinky1.2.xFlink SQL 一站式开发、CDC 任务管理
OLAP 查询Trino435+近实时交互式分析
存储底座(热)本地 HDD/SSD-近期热数据存储
存储底座(冷)MinIO / S3-历史冷数据归档存储
编排调度Kubernetes + Flink Operator1.15容器化部署、弹性伸缩
监控Prometheus + Grafana-指标采集与可视化

2.2 v2.0 版本改进要点

改进项v1.0 方案v2.0 方案改进收益
数据链路MySQL → CDC → Paimon(直连)MySQL → CDC → Paimon(CDC 3.6 原生支持)Exactly-Once 断点续传,零中间件
CatalogHive MetastoreJDBC Catalog中心化元数据 + 分布式锁,跨引擎一致
部署方式YARNKubernetes + Flink Operator弹性伸缩,云原生架构
SQL 开发Flink SQL CLIApache 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存储位置保留策略
ODSAppend-only8热存储 (HDD)不物理过期,由冷热归档管理
DWDPrimary-keyorder_id(不含 dt)-1 动态桶热存储 (HDD)不物理过期,由冷热归档管理
DWSMaterialized-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 CatalogJDBC 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-compactionFull 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 COLUMNDROP COLUMNCHANGE TYPERENAME 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 3minMERGE INTO 分区级重同步
Paimon 表误删每日 Tag + Clone 备份从备份 Clone 恢复
Paimon 数据误更新关键操作前创建 Tagrollback_to 快照回滚
上游误 DROP TABLElenient 模式 + 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 分钟
CheckpointCheckpoint 完成时间> 10 分钟
Backpressure算子反压比率> 50% 持续 5 分钟
CompactionPaimon 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 JobManager1-3 (HA)2C / 4GB100GB SSD主节点
Flink TaskManager2-6 (弹性)4C / 8GB500GB SSD计算节点
Dinky1-22C / 4GB50GBSQL 开发平台
MinIO42C / 4GB2TB+ HDD冷存储 / Checkpoint
MySQL(元数据)12C / 4GB100GB SSDPaimon JDBC Catalog
Trino1-24C / 8GB100GBOLAP 查询
Prometheus + Grafana12C / 4GB500GB监控

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 HDD3 台约 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 Flink2.2.12.0.0JDK 11+(推荐 17)2.0 起移除 Java 8 和 Per-Job 模式
Apache Paimon1.0.x0.9.0JDK 11+1.0 正式版,推荐生产使用
Flink-CDC3.6.03.0.0JDK 11+同时兼容 Flink 1.20.x 和 2.2.x
Apache Dinky1.2.51.2.0JDK 11+支持 Flink 1.14~1.20,K8s Operator 集成
Flink K8s Operator1.151.10-支持 Flink v1.19 到 v2.2
MinIORELEASE.2024+--S3 兼容冷存储 / Checkpoint
Trino435+425+JDK 17+需安装 Paimon Connector
MySQL5.7 / 8.0+5.6-CDC Source 需要 binlog 开启
PostgreSQL14+10+-需配置 WAL logical replication
Oracle19c+12c-3.6 新增 Pipeline Source

参考资料

  1. Apache Flink 官方文档 - Flink 2.0 Release Notes
  2. Apache Paimon 1.0 发布公告
  3. Flink CDC 3.6.0 发布公告
  4. Apache Dinky 官方文档
  5. Flink Kubernetes Operator 文档
  6. Flink CDC 3.6 官方文档 - MySQL Pipeline Connector
  7. Flink CDC 3.6 官方文档 - Paimon Sink Connector
  8. Flink CDC 3.6 官方文档 - Route 模块
  9. Flink CDC 3.6 官方文档 - Transform 模块
  10. Apache Paimon 1.0 Procedures - Clone / Rollback