Apache Flink 学习指南

写在前面:如果你是一名刚接触 Flink 的数据开发工程师,面对流处理、状态管理、Checkpoint、Watermark 等概念感到迷茫,这篇文章就是为你准备的。我们将从"什么是 Flink"一路讲到性能调优与实时数仓架构,配合大量代码示例和对比表格,帮你建立完整的知识体系。本文基于 Flink 1.18/1.19/1.20 版本,并展望 Flink 2.0 方向。


目录

  1. Flink 概述
  2. Flink 架构与运行原理
  3. 环境搭建
  4. DataStream API
  5. 时间语义与 Window
  6. 状态管理
  7. 容错机制
  8. Flink SQL
  9. Flink CDC
  10. Flink CEP
  11. Flink 性能调优
  12. Flink 常见问题与踩坑
  13. Flink 1.18-1.20 新特性与 2.0 展望
  14. Flink 在实时数仓中的典型架构
  15. 学习路线与实战建议

1. Flink 概述

1.1 什么是 Flink

Apache Flink 是一个分布式流处理引擎,同时支持批处理(流批一体)。它的核心定位是:有状态的实时计算,即在无限数据流上进行低延迟、高吞吐、精确一次(Exactly-Once)的计算。

有限数据集(批)          无限数据流(流)
┌──────────┐           ─ ─ ─ ─ ─ ─ ─ ─ ▶
│ 数据有边界 │         │ 数据无边界,持续到达 │
└──────────┘           ─ ─ ─ ─ ─ ─ ─ ─ ▶
   Flink 批处理              Flink 流处理(核心)

1.2 发展历史

时间里程碑
2010德国柏林理工大学研究项目 “Stratosphere”
2014捐赠 Apache,更名为 Flink,成为顶级项目
2016Flink 1.0 发布,引入 Event Time & CEP
2019Flink 1.9 合并 Table API,流批一体雏形
2021Flink 1.13/1.14,State Backend 重构,Unaligned Checkpoint
2023Flink 1.18,Flink CDC 3.0,Lakehouse 集成
2024Flink 1.19/1.20,自适应执行,Flink 2.0 路线图

1.3 为什么选 Flink

  • 流批一体:同一套 API 同时处理流和批,无需维护两套系统。
  • Exactly-Once:基于 Checkpoint 机制保证状态一致性。
  • 高性能低延迟:纯流处理,毫秒级延迟;JVM + 内存计算,高吞吐。
  • 强大的状态管理:支持超大状态(TB 级),RocksDB 后端。
  • 丰富的时间语义:Event Time + Watermark,完美处理乱序数据。
  • 完善的生态:Flink SQL、CDC、CEP、ML、Paimon 等。

1.4 Flink vs Spark Streaming

对比维度FlinkSpark Streaming
处理模型真流处理(事件驱动,逐条处理)微批处理(将流切分成小批次)
延迟毫秒级秒级(批次间隔决定)
吞吐高高
时间语义原生 Event Time + Watermark微批时间窗口,Event Time 支持较弱
状态管理原生支持,Checkpoint/SavepointDStream 状态管理较繁琐
Exactly-Once端到端精确一次至少一次(需手动保证精确一次)
SQL 支持Flink SQL(成熟)Spark SQL(更成熟于批)
适用场景实时风控、实时监控、实时数仓近实时 ETL、批流混合

一句话总结:Flink 是"来一条处理一条",Spark Streaming 是"攒一批处理一批"。

1.5 典型应用场景

  • 实时数仓:Kafka → Flink ETL → Doris/StarRocks/ClickHouse
  • 实时监控告警:日志/指标流 → CEP 规则匹配 → 告警
  • 实时风控:交易流 → 模式检测 → 拦截可疑交易
  • CEP 复杂事件处理:从事件流中检测特定模式
  • 机器学习:Flink ML 在线学习/实时特征工程
  • 实时推荐:用户行为流 → 实时特征更新 → 推荐系统

2. Flink 架构与运行原理

2.1 整体架构

Flink 运行时由主从架构组成:

在这里插入图片描述

2.2 核心组件

组件职责
JobManager主节点,负责作业调度、Checkpoint 协调、故障恢复。包含 Dispatcher、ResourceManager、JobMaster
TaskManager工作节点,执行具体的 Task,提供 Slot 资源,缓存和交换数据
ResourceManager管理 TaskManager 的 Slot 分配,支持 YARN/K8s/Standalone
Dispatcher提供 REST 接口,接收作业提交,启动 JobMaster
Client编译代码生成 JobGraph,提交给集群(不参与运行时)

2.3 TaskSlot 与 SlotSharing

TaskSlot 是 TaskManager 中的资源单位,每个 Slot 独占一份内存(主要是 Managed Memory),但共享 CPU。

TaskManager (3 Slots)
┌─────────────────────────────────┐
│  Slot 0    │  Slot 1   │ Slot 2 │
│  (并行度1)  │ (并行度1)  │(并行度1)│
│            │           │        │
│  source    │  source   │ source │
│    map     │    map    │   map  │
│   sink     │   sink    │  sink  │  ← SlotSharing:不同算子的子任务可共享同一Slot
└─────────────────────────────────┘

SlotSharing 允许同一个 Job 的不同算子的子任务共享同一个 Slot,好处:

  • 资源均衡:每个 Slot 持有完整的 Pipeline
  • 减少数据网络传输
  • 默认开启,可通过 slotSharingGroup() 控制

2.4 ExecutionGraph 生成

用户代码经过四层图转换:

StreamGraph (用户代码)
    │  优化(算子链合并)
    ▼
JobGraph (提交给集群)
    │  并行化展开
    ▼
ExecutionGraph (JobMaster 生成,并行执行版本)
    │  部署到 TaskManager
    ▼
物理执行图 (实际 Task 运行)

2.5 数据交换与 Buffer

TaskManager 之间通过 Network Buffer 传输数据:

  • 基于 Netty 的 TCP 长连接
  • BufferPool 管理,每个 Channel 有独占 Buffer + 浮动 Buffer
  • 反压机制:下游 Buffer 满 → 上游无法写入 → 逐级反压

2.6 部署模式

模式说明适用场景
Session集群预先启动,多个作业共享集群小规模作业、调试
Per-Job每个作业启动独立集群,作业结束集群释放生产环境(Flink 1.15 前主流)
Application每个作业独立 JobManager,main() 在集群执行生产推荐,资源隔离好

2.7 资源调度平台

  • Standalone:独立集群,手动管理
  • YARN:Hadoop 生态,经典企业级部署
  • Kubernetes:云原生方向,Flink 1.18+ 原生 Operator 支持,自动扩缩容

2.8 TaskManager 内存模型详解

Flink TaskManager 进程的内存配置直接影响作业稳定性和性能,理解内存模型是生产调优的基础。

内存全景图:

┌──────────────────────── TaskManager Process Memory ────────────────────────┐
│                                                                            │
│  ┌─────────────────────── Total Flink Memory ──────────────────────────┐   │
│  │                                                                     │   │
│  │  ┌── Framework Heap ──┐  ┌── Task Heap ────────────────────────┐    │   │
│  │  │  (Flink框架本身)    │  │  (用户代码/算子对象)                  │   │   │
│  │  │  默认 128MB         │  │  Managed Memory - 上述部分           │   │   │
│  │  └────────────────────┘  └──────────────────────────────────────┘   │   │
│  │                                                                     │   │
│  │  ┌── Framework Off-Heap ─┐  ┌── Task Off-Heap ─────────────────┐    │   │
│  │  │  (框架堆外,默认128MB)  │  │  (用户堆外内存,默认0)             │   │   │
│  │  └───────────────────────┘  └───────────────────────────────────┘   │   │
│  │                                                                     │   │
│  │  ┌── Network Memory ─────────────────────────────────────────────┐  │   │
│  │  │  (网络缓冲,Shuffle数据交换,默认占Total Flink 10%)             │  │   │
│  │  └──────────────────────────────────────────────────────────────┘  │   │
│  │                                                                     │   │
│  │  ┌── Managed Memory ────────────────────────────────────────────┐  │   │
│  │  │  (RocksDB状态/Python/排序/哈希表,默认占Total Flink 40%)       │  │   │
│  │  └──────────────────────────────────────────────────────────────┘  │   │
│  │                                                                     │   │
│  └─────────────────────────────────────────────────────────────────────┘   │
│                                                                            │
│  ┌── JVM Metaspace ──┐  ┌── JVM Overhead ──────────────────────────┐     │
│  │  (类元数据,默认256MB)│  │  (线程栈/IO/编译等,默认占总内存10%)     │     │
│  └────────────────────┘  └──────────────────────────────────────────┘     │
│                                                                            │
└────────────────────────────────────────────────────────────────────────────┘

各内存区域说明:

内存区域作用默认配置
Framework HeapFlink 框架本身使用的堆内存(不计入 Slot)128MB
Task Heap用户算子代码、对象分配的堆内存剩余自动推算
Framework Off-HeapFlink 框架堆外内存128MB
Task Off-Heap用户代码申请的堆外内存(如 Netty/NIO)0 bytes
Network Memory网络 Shuffle 缓冲池,每个 InputGate/ResultPartition 使用10% of Total Flink Memory
Managed MemoryRocksDB 状态后端、Python 进程、批处理排序/哈希表40% of Total Flink Memory
JVM MetaspaceJVM 类加载元数据256MB
JVM OverheadJVM 其他开销(线程栈、GC 结构、IO 缓冲区)10% of Process Memory

关键配置参数:

# flink-conf.yaml

# ===== 进程总内存(推荐直接配置此项,Flink自动推算各部分)=====
taskmanager.memory.process.size: 4096m    # TM 进程总内存(含JVM开销)
# 或配置 Flink 总内存(不含JVM开销),Flink 自动加 JVM 部分:
# taskmanager.memory.flink.size: 4096m

# ===== Framework 内存 =====
taskmanager.memory.framework.heap.size: 128mb
taskmanager.memory.framework.off-heap.size: 128mb

# ===== Task 内存 =====
taskmanager.memory.task.heap.size: 4096m       # 不设则自动推算
taskmanager.memory.task.off-heap.size: 0 bytes

# ===== Managed Memory =====
taskmanager.memory.managed.fraction: 0.4       # 占 Total Flink Memory 比例
taskmanager.memory.managed.size: 1024m         # 也可直接指定绝对值
# RocksDB 状态后端下,Managed Memory 全部用于 RocksDB 的块缓存和写缓冲区

# ===== Network Memory =====
taskmanager.memory.network.fraction: 0.10      # 占 Total Flink Memory 比例
taskmanager.memory.network.min: 64mb
taskmanager.memory.network.max: 1gb

# ===== JVM 部分 =====
taskmanager.memory.jvm-metaspace.size: 256mb
taskmanager.memory.jvm-overhead.fraction: 0.10
taskmanager.memory.jvm-overhead.min: 192mb
taskmanager.memory.jvm-overhead.max: 1gb

配置优先级:显式配置绝对值 > fraction 比例推算。配置 process.size 后,Flink 按 fraction 自动划分各区域;若某区域显式配置了绝对值,则从剩余部分中再按比例分配。

RocksDB 状态后端下 Managed Memory 的分配:

使用 RocksDB 时,Managed Memory 被划分为每个 Slot 独占的份额,用于:

  • Block Cache:RocksDB 读缓存(默认占 Managed Memory 的约 1/3)
  • Write Buffer Manager:控制 RocksDB MemTable 内存(默认约占 2/3)
Managed Memory (per TM)
├── Slot 0 Share
│   ├── RocksDB Block Cache (~1/3)
│   └── Write Buffer / MemTables (~2/3)
├── Slot 1 Share
│   ├── Block Cache
│   └── Write Buffer
└── ...

Flink 自动将 Managed Memory 转化为 RocksDB 的 write_buffer_manager 和 cache 配置,避免手动调参。可通过以下参数微调:

state.backend.rocksdb.memory.managed: true        # 默认开启,使用Managed Memory
state.backend.rocksdb.memory.write-buffer-ratio: 0.5  # 写缓冲占比
state.backend.rocksdb.memory.high-prio-pool-ratio: 0.1 # 高优先级池占比

常见内存问题与排查:

问题现象原因解决
容器 OOM KilledPod/Container 被 YARN/K8s 杀死,exit code 137JVM Overhead/Direct Memory 超出容器限制增大 jvm-overhead.fraction;检查用户代码是否有堆外泄漏;适当降低 process.size 留出余量
Managed Memory 不足RocksDB 频繁刷盘,Checkpoint 慢,吞吐下降Managed Memory fraction 太小增大 managed.fraction(大状态场景可设 0.5~0.6);减少每个 TM 的 Slot 数
Network Buffer 不足Insufficient number of network buffers 异常高并行度下网络缓冲不够增大 network.max;检查 taskmanager.network.memory.buffers-per-channel
Heap OOMjava.lang.OutOfMemoryError: Java heap space用户代码对象过多/状态过大增大 Task Heap;大状态切换 RocksDB 后端;优化代码减少对象
Direct Memory OOMOutOfMemoryError: Direct buffer memory堆外内存不足增大 task.off-heap.size 或 jvm-overhead

内存调优建议:

# ===== 流式作业(大状态 RocksDB)=====
taskmanager.memory.process.size: 8192m
taskmanager.memory.managed.fraction: 0.5     # 大状态给更多 Managed Memory
taskmanager.memory.network.fraction: 0.10
taskmanager.numberOfTaskSlots: 2             # 大状态场景减少 Slot,增加每 Slot 内存

# ===== 批式作业(排序/哈希密集)=====
taskmanager.memory.process.size: 8192m
taskmanager.memory.managed.fraction: 0.6     # 批处理 Managed Memory 用于排序/哈希
taskmanager.memory.network.fraction: 0.15    # 批 Shuffle 数据量大
taskmanager.numberOfTaskSlots: 4

# ===== 小状态流作业(HashMap 后端)=====
taskmanager.memory.process.size: 4096m
taskmanager.memory.managed.fraction: 0.2     # HashMap 不需要太多 Managed Memory
taskmanager.memory.network.fraction: 0.10
taskmanager.numberOfTaskSlots: 4

黄金法则:容器环境务必设置 taskmanager.memory.process.size,且该值应略小于容器内存限制(留 10%~15% 余量给 JVM 自身开销),否则容易被 OOM Killed。


3. 环境搭建

3.1 Local 模式(快速体验)

# 下载 Flink
wget https://archive.apache.org/dist/flink/flink-1.18.1/flink-1.18.1-bin-scala_2.12.tgz
tar -xzf flink-1.18.1-bin-scala_2.12.tgz
cd flink-1.18.1

# 启动本地集群
./bin/start-cluster.sh

# 访问 Web UI: http://localhost:8081

3.2 Standalone 集群

# conf/flink-conf.yaml
jobmanager.rpc.address: flink-master
jobmanager.memory.process.size: 1600m
taskmanager.memory.process.size: 4096m
taskmanager.numberOfTaskSlots: 4
parallelism.default: 4

# conf/workers
flink-worker1
flink-worker2
flink-worker3

# 分发到所有节点后启动
./bin/start-cluster.sh

3.3 On YARN

# Session 模式
./bin/yarn-session.sh -nm flink-session -n 4 -s 4 -jm 1024m -tm 4096m

# Application 模式(推荐)
./bin/flink run-application -t yarn-application \
  -Djobmanager.memory.process.size=1024m \
  -Dtaskmanager.memory.process.size=4096m \
  -Dtaskmanager.numberOfTaskSlots=4 \
  ./examples/streaming/TopSpeedWindowing.jar

3.4 Flink SQL Client

./bin/sql-client.sh

# 在 SQL Client 中
SET 'sql-client.execution.result-mode' = 'tableau';

CREATE TABLE kafka_source (
  id BIGINT,
  name STRING,
  ts TIMESTAMP(3),
  WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'test',
  'properties.bootstrap.servers' = 'localhost:9092',
  'format' = 'json',
  'scan.startup.mode' = 'latest-offset'
);

SELECT * FROM kafka_source;

3.5 PyFlink 环境

# 安装 PyFlink
pip install apache-flink==1.18.1

# 验证
python -c "from pyflink.datastream import StreamExecutionEnvironment; print('OK')"
# PyFlink 快速示例
from pyflink.datastream import StreamExecutionEnvironment

env = StreamExecutionEnvironment.get_execution_environment()
ds = env.from_collection([1, 2, 3, 4, 5])
ds.map(lambda x: x * 2).print()
env.execute("pyflink_demo")

3.6 Flink on Kubernetes

随着云原生普及,Kubernetes 已成为 Flink 生产部署的主流选择。Flink 原生支持在 K8s 上运行,JobManager 和 TaskManager 均以 Pod 形式调度。

Native K8s 部署模式原理:

┌─────────────────────────── Kubernetes Cluster ───────────────────────────┐
│                                                                          │
│  ┌─── Flink Client (kubectl/flink) ───┐                                  │
│  │  解析 flink-conf.yaml → 调用 K8s API│                                  │
│  └────────────────┬────────────────────┘                                 │
│                   │ 创建                                                  │
│                   ▼                                                       │
│  ┌─── JobManager Pod ──────────────────────────────────────────────┐     │
│  │  JM Container                                                   │     │
│  │  ┌──────────┐ ┌──────────────┐ ┌────────────┐                  │     │
│  │  │Dispatcher│ │ResourceManager│ │  JobMaster │                  │     │
│  │  └──────────┘ └──────┬───────┘ └────────────┘                  │     │
│  └──────────────────────┼──────────────────────────────────────────┘     │
│                         │ 申请Slot → 动态创建TM Pod                        │
│         ┌───────────────┼───────────────┐                                │
│         ▼               ▼               ▼                                │
│  ┌─── TM Pod ───┐ ┌─── TM Pod ───┐ ┌─── TM Pod ───┐                     │
│  │ TaskManager  │ │ TaskManager  │ │ TaskManager  │                     │
│  │ Slot ×N      │ │ Slot ×N      │ │ Slot ×N      │                     │
│  └──────────────┘ └──────────────┘ └──────────────┘                     │
│                                                                          │
└──────────────────────────────────────────────────────────────────────────┘

Flink 直接与 K8s API Server 交互:ResourceManager 按需申请/归还 Pod,无需预先部署 TaskManager。

三种部署模式对比(K8s 下):

模式SessionApplicationPer-Job
集群生命周期长运行,多作业共享每作业独立集群,作业结束销毁每作业独立集群(1.15 已废弃)
main() 执行位置Client 端JobManager 端(集群内)Client 端
资源隔离差(多作业共享 JM/TM)好(独立 JM)好
适用场景短期调试、小规模作业生产推荐已废弃,使用 Application 替代
JM 故障影响所有作业失败仅当前作业仅当前作业

Application 模式提交命令:

# Native K8s Application 模式
./bin/flink run-application \
  --target kubernetes-application \
  -Dkubernetes.cluster-id=flink-streaming-job \
  -Dkubernetes.container.image=flink:1.18.1-java17 \
  -Dkubernetes.namespace=flink-prod \
  -Djobmanager.memory.process.size=1024m \
  -Dtaskmanager.memory.process.size=4096m \
  -Dtaskmanager.numberOfTaskSlots=4 \
  -Dkubernetes.taskmanager.cpu=2.0 \
  -Dkubernetes.jobmanager.cpu=1.0 \
  -Dkubernetes.rest-service.exposed.type=LoadBalancer \
  local:///opt/flink/usrlib/my-streaming-job.jar

# 查看作业
./bin/flink list --target kubernetes-application \
  -Dkubernetes.cluster-id=flink-streaming-job

# 取消作业
./bin/flink cancel --target kubernetes-application \
  -Dkubernetes.cluster-id=flink-streaming-job <jobId>

Docker 镜像构建:

# Dockerfile - 基于官方镜像叠加业务依赖
FROM flink:1.18.1-java17

# 复制用户 JAR 包(SQL Connector/UDX 等)
COPY ./lib/flink-sql-connector-kafka-3.1.0-1.18.jar /opt/flink/lib/
COPY ./lib/flink-connector-jdbc-3.2.0-1.18.jar /opt/flink/lib/
COPY ./lib/mysql-connector-j-8.0.33.jar /opt/flink/lib/
COPY ./lib/flink-sql-connector-mysql-cdc-3.0.1.jar /opt/flink/lib/
COPY ./target/my-streaming-job.jar /opt/flink/usrlib/

# 自定义日志配置
COPY ./conf/log4j2.xml /opt/flink/conf/log4j2.xml

# 设置时区
ENV TZ=Asia/Shanghai
RUN ln -snf /usr/share/zoneinfo/$TZ /etc/localtime && echo $TZ > /etc/timezone
# 构建并推送镜像
docker build -t my-registry/flink:1.18.1-myjob .
docker push my-registry/flink:1.18.1-myjob

K8s 关键配置参数:

参数说明示例
kubernetes.cluster-id集群唯一标识,用于资源命名flink-streaming-job
kubernetes.container.imageDocker 镜像flink:1.18.1-java17
kubernetes.container.image.pull-policy镜像拉取策略IfNotPresent / Always
kubernetes.namespaceK8s 命名空间flink-prod
kubernetes.service-accountRBAC 服务账号flink-service-account
kubernetes.taskmanager.cpuTM CPU 请求2.0
kubernetes.jobmanager.cpuJM CPU 请求1.0
taskmanager.memory.process.sizeTM 进程内存4096m
jobmanager.memory.process.sizeJM 进程内存1024m
taskmanager.numberOfTaskSlots每 TM Slot 数4
kubernetes.rest-service.exposed.typeREST 暴露方式ClusterIP / NodePort / LoadBalancer
kubernetes.pod-template-filePod 模板文件(高级配置)/etc/flink/pod-template.yaml

Pod Template 示例(高级配置):

# pod-template.yaml - 挂载 ConfigMap/Secret/卷
apiVersion: v1
kind: Pod
metadata:
  name: pod-template
spec:
  containers:
    - name: flink-main-container
      volumeMounts:
        - name: flink-config-volume
          mountPath: /opt/flink/conf
        - name: hdfs-config
          mountPath: /etc/hadoop/conf
      env:
        - name: TZ
          value: Asia/Shanghai
  volumes:
    - name: flink-config-volume
      configMap:
        name: flink-config
    - name: hdfs-config
      configMap:
        name: hdfs-config

Flink K8s Operator 简介:

Flink K8s Operator 是 Apache 官方项目,通过 CRD(Custom Resource Definition)以声明式方式管理 Flink 作业,支持自动扩缩容、滚动升级、Savepoint 管理。

# 使用 Helm 安装 Flink K8s Operator
helm repo add flink-operator-repo https://downloads.apache.org/flink/flink-kubernetes-operator-1.9.0/
helm install flink-kubernetes-operator flink-operator-repo/flink-kubernetes-operator \
  --namespace flink-system \
  --create-namespace
# flinkdeployment.yaml - 声明式定义 Flink 作业
apiVersion: flink.apache.org/v1beta1
kind: FlinkDeployment
metadata:
  name: flink-streaming-job
  namespace: flink-prod
spec:
  image: my-registry/flink:1.18.1-myjob
  flinkVersion: v1_18
  flinkConfiguration:
    taskmanager.numberOfTaskSlots: "4"
    state.backend: rocksdb
    state.checkpoints.dir: s3://flink/checkpoints
    state.savepoints.dir: s3://flink/savepoints
    execution.checkpointing.interval: "60s"
  jobManager:
    resource:
      memory: "1024m"
      cpu: 1
  taskManager:
    resource:
      memory: "4096m"
      cpu: 2
  podTemplate:
    spec:
      containers:
        - name: flink-main-container
          env:
            - name: TZ
              value: Asia/Shanghai
  job:
    jarURI: local:///opt/flink/usrlib/my-streaming-job.jar
    parallelism: 8
    upgradeMode: savepoint          # 升级时自动做 Savepoint
    state: running
# 部署
kubectl apply -f flinkdeployment.yaml

# 查看状态
kubectl get flinkdeployment -n flink-prod
kubectl describe flinkdeployment flink-streaming-job -n flink-prod

K8s vs YARN 对比:

对比维度KubernetesYARN
资源模型容器化(Pod),CPU/Memory 显式声明Container(Slot),内存为主
弹性扩缩原生 HPA/Operator Autoscaler,秒级伸缩需配置 NodeManager,弹性较弱
多租户Namespace + RBAC + ResourceQuota 隔离完善Queue 隔离,粒度较粗
镜像管理Docker 镜像版本化,依赖打包一致性好依赖散落在节点,易出现环境不一致
运维生态kubectl/Helm/Prometheus 原生集成YARN UI + ResourceManager
存储集成PVC/HostPath/S3/OSS,CSI 插件丰富HDFS 深度集成,外部存储需额外配置
启动速度Pod 拉起较快(镜像缓存后秒级)Container 启动较慢
学习成本需掌握 K8s 概念(Pod/Service/Ingress)Hadoop 生态团队上手快
适用场景云原生/混合云/新建集群已有 Hadoop 生态的企业

4. DataStream API

4.1 执行环境

// Java
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setParallelism(4);
env.enableCheckpointing(60000); // 60秒一次Checkpoint
# PyFlink
from pyflink.datastream import StreamExecutionEnvironment
env = StreamExecutionEnvironment.get_execution_environment()
env.set_parallelism(4)

4.2 Source

Kafka Source(Flink 1.14+ 推荐使用 KafkaSource):

KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers("localhost:9092")
    .setTopics("orders")
    .setGroupId("flink-group")
    .setStartingOffsets(OffsetsInitializer.latest())
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

DataStream<String> stream = env.fromSource(
    source, WatermarkStrategy.noWatermarks(), "Kafka Source");

File Source:

FileSource<String> fileSource = FileSource.forRecordStreamFormat(
    new TextLineInputFormat(), new Path("/data/input/"))
    .monitorContinuously(Duration.ofSeconds(10)) // 流式监控新文件
    .build();

DataStream<String> lines = env.fromSource(
    fileSource, WatermarkStrategy.noWatermarks(), "File Source");

Socket Source:

DataStream<String> socketStream = env.socketTextStream("localhost", 9999);

自定义 Source:

public class CustomSource implements SourceFunction<Order> {
    private volatile boolean running = true;

    @Override
    public void run(SourceContext<Order> ctx) throws Exception {
        Random random = new Random();
        while (running) {
            Order order = new Order(
                random.nextLong(),
                "item_" + random.nextInt(100),
                random.nextDouble() * 1000,
                System.currentTimeMillis()
            );
            ctx.collect(order);
            Thread.sleep(100);
        }
    }

    @Override
    public void cancel() {
        running = false;
    }
}

DataStream<Order> customStream = env.addSource(new CustomSource());

注意:Flink 推荐使用新的 Source 接口(FLIP-27),旧的 SourceFunction 仍可用但逐渐被取代。

4.3 Transformation

// map: 一对一转换
DataStream<Order> enriched = orders.map(order -> {
    order.setDiscounted(order.getAmount() * 0.9);
    return order;
});

// flatMap: 一对多展开
DataStream<String> words = lines.flatMap((String line, Collector<String> out) -> {
    for (String word : line.split(" ")) {
        out.collect(word);
    }
}).returns(Types.STRING);

// filter: 过滤
DataStream<Order> largeOrders = orders.filter(o -> o.getAmount() > 1000);

// keyBy: 按 key 分区(注意:返回 KeyedStream)
KeyedStream<Order, String> keyedByItem = orders.keyBy(Order::getItemId);

// reduce: 增量聚合
DataStream<Order> maxOrder = keyedByItem
    .reduce((a, b) -> a.getAmount() > b.getAmount() ? a : b);

// aggregations: sum/min/max
DataStream<Order> sumByItem = keyedByItem
    .sum("amount")  // 按字段名聚合
    .name("sum-by-item");

// window: 在 keyBy 后开窗
WindowedStream<Order, String, TimeWindow> tumblingWindow = keyedByItem
    .window(TumblingProcessingTimeWindows.of(Time.seconds(10)));

// connect: 合流(保留各自类型)
DataStream<OrderA> streamA = ...;
DataStream<OrderB> streamB = ...;
ConnectedStreams<OrderA, OrderB> connected = streamA.connect(streamB);
DataStream<String> result = connected.map(
    new CoMapFunction<OrderA, OrderB, String>() {
        @Override
        public String map1(OrderA value) { return "A:" + value; }
        @Override
        public String map2(OrderB value) { return "B:" + value; }
    });

// union: 合流(相同类型,可多个流合并)
DataStream<Order> allOrders = streamA.union(streamB, streamC);

// side outputs: 侧输出(处理迟到数据等)
final OutputTag<Order> lateTag = new OutputTag<Order>("late-events"){};
SingleOutputStreamOperator<Order> mainStream = orders
    .process(new ProcessFunction<Order, Order>() {
        @Override
        public void processElement(Order order, Context ctx, Collector<Order> out) {
            if (order.getAmount() < 0) {
                ctx.output(lateTag, order); // 侧输出
            } else {
                out.collect(order);
            }
        }
    });
DataStream<Order> lateStream = mainStream.getSideOutput(lateTag);

4.4 Sink

Kafka Sink:

KafkaSink<String> kafkaSink = KafkaSink.<String>builder()
    .setBootstrapServers("localhost:9092")
    .setRecordSerializer(KafkaRecordSerializationSchema.builder()
        .setTopic("output-topic")
        .setValueSerializationSchema(new SimpleStringSchema())
        .build())
    .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
    .setTransactionalIdPrefix("flink-txn-")
    .build();

stream.sinkTo(kafkaSink);

JDBC Sink:

JdbcSink.sink(
    "INSERT INTO orders (id, item_id, amount) VALUES (?, ?, ?)",
    (ps, order) -> {
        ps.setLong(1, order.getId());
        ps.setString(2, order.getItemId());
        ps.setDouble(3, order.getAmount());
    },
    JdbcExecutionOptions.builder()
        .withBatchSize(1000)
        .withBatchIntervalMs(200)
        .build(),
    new JdbcConnectionOptions.JdbcConnectionOptionsBuilder()
        .withUrl("jdbc:mysql://localhost:3306/analytics")
        .withDriverName("com.mysql.cj.jdbc.Driver")
        .withUsername("root")
        .withPassword("password")
        .build()
);

文件 Sink(StreamingFileSink → FileSink):

FileSink<String> fileSink = FileSink.forRowFormat(
    new Path("/data/output/"),
    new SimpleStringEncoder<String>("UTF-8"))
    .withRollingPolicy(DefaultRollingPolicy.builder()
        .withRolloverInterval(Duration.ofMinutes(5))
        .withInactivityInterval(Duration.ofMinutes(2))
        .withMaxPartSize(MemorySize.parse("128MB"))
        .build())
    .build();

stream.sinkTo(fileSink);

4.5 并行度设置

并行度优先级(从高到低):

  1. 算子级别:.setParallelism(8)
  2. 执行环境级别:env.setParallelism(4)
  3. 提交参数:-p 4
  4. 配置文件:parallelism.default: 1

4.6 异步 I/O(AsyncDataStream)

在流处理中,经常需要与外部系统交互(如查询 Redis/MySQL 维表、调用 HTTP 接口)。同步请求会阻塞算子线程,导致吞吐极低。异步 I/O 允许并发处理多个请求,大幅提升吞吐。

为什么需要异步 I/O:

同步请求(串行,瓶颈明显):        异步请求(并发,吞吐倍增):
───[查询DB]──▶───[查询DB]──▶      ┌──[查询DB]──┐
   等待...     等待...             ├──[查询DB]──┤──▶ 批量回调
   10ms/条    10ms/条             ├──[查询DB]──┤
   100条/秒                       └──[查询DB]──┘
                                   1000条/秒+

前提:外部系统客户端必须支持异步 API(如 Lettuce(Redis)、Async MySQL Driver、AsyncHttpClient)。若客户端仅支持同步,可使用线程池模拟异步,但效果有限。

AsyncFunction 实现:

// 异步查询 Redis 维表(使用 Lettuce 异步客户端)
public class AsyncRedisDimFunction extends RichAsyncFunction<Order, Order> {

    private transient RedisAsyncCommands<String, String> asyncCommands;
    private transient ExecutorService executor;

    @Override
    public void open(Configuration parameters) {
        RedisClient redisClient = RedisClient.create("redis://localhost:6379");
        StatefulRedisConnection<String, String> connection =
            redisClient.connect();
        asyncCommands = connection.async();

        // 对于仅支持同步的客户端,用线程池模拟
        executor = Executors.newFixedThreadPool(20);
    }

    @Override
    public void asyncInvoke(Order order, ResultFuture<Order> resultFuture) {
        // 异步查询 Redis
        RedisFuture<String> future = asyncCommands.hget(
            "dim:item", order.getItemId());

        future.thenAccept(dimJson -> {
            if (dimJson != null) {
                ItemDim dim = parseJson(dimJson);
                order.setItemName(dim.getName());
                order.setCategory(dim.getCategory());
            }
            // 完成回调,输出结果
            resultFuture.complete(Collections.singletonList(order));
        }).exceptionally(ex -> {
            // 异常处理:记录日志,原样输出或发送侧输出
            LOG.warn("Async Redis query failed for item: {}",
                order.getItemId(), ex);
            resultFuture.complete(Collections.singletonList(order));
            return null;
        });
    }

    @Override
    public void timeout(Order order, ResultFuture<Order> resultFuture) {
        // 超时处理:可输出默认值或侧输出
        LOG.warn("Async query timed out for order: {}", order.getId());
        resultFuture.complete(Collections.singletonList(order));
    }

    @Override
    public void close() {
        if (executor != null) executor.shutdown();
    }
}

使用 AsyncDataStream 调用:

DataStream<Order> enrichedStream = AsyncDataStream
    .unorderedWait(
        orders,                          // 输入流
        new AsyncRedisDimFunction(),     // AsyncFunction
        5000,                            // 超时时间 5s
        TimeUnit.MILLISECONDS,
        100                              // 最大并发请求数 capacity
    )
    .name("async-redis-dim-lookup")
    .setParallelism(4);

// 有序输出版本(保证事件顺序,吞吐略低)
DataStream<Order> orderedStream = AsyncDataStream
    .orderedWait(
        orders,
        new AsyncMysqlDimFunction(),
        3, TimeUnit.SECONDS,
        50
    );

orderedWait vs unorderedWait:

模式顺序保证吞吐延迟适用场景
unorderedWait不保证顺序,先完成先输出高低大部分维表关联场景
orderedWait严格按输入顺序输出较低较高(需等待前面的完成)需要严格顺序的业务

注意:使用 EventTime 时,unorderedWait 会在 Watermark 处保序,即先完成的记录可以先输出,但不会越过 Watermark;orderedWait 则完全按输入顺序。

关键参数:

参数说明建议
timeout单个异步请求最大等待时间,超时触发 timeout()根据下游 SLA 设置,通常 3~10s
capacity最大并发异步请求数(同时 in-flight 的请求数)根据外部系统承载能力,通常 50~200
backpressure当并发请求达到 capacity 时,自动反压上游自动机制,无需手动配置

生产注意事项:

  1. 连接池管理:AsyncFunction 中使用连接池,避免每条请求新建连接
  2. 超时与重试:设置合理的超时;失败重试需注意幂等性,建议最多重试 2~3 次
  3. 缓存:在 AsyncFunction 中加本地缓存(Caffeine/Guava Cache),减少对外部系统的请求
  4. 降级:查询失败时使用默认值或侧输出,避免作业失败
  5. 监控:关注异步请求延迟、成功率、超时率
// 带 Caffeine 缓存的 AsyncFunction 示例
public class CachedAsyncDimFunction extends RichAsyncFunction<Order, Order> {
    private transient Cache<String, ItemDim> cache;
    private transient AsyncHttpClient httpClient;

    @Override
    public void open(Configuration parameters) {
        cache = Caffeine.newBuilder()
            .maximumSize(10_000)
            .expireAfterWrite(Duration.ofMinutes(10))
            .build();
        httpClient = Dsl.asyncHttpClient();
    }

    @Override
    public void asyncInvoke(Order order, ResultFuture<Order> resultFuture) {
        ItemDim cached = cache.getIfPresent(order.getItemId());
        if (cached != null) {
            order.setItemName(cached.getName());
            resultFuture.complete(Collections.singletonList(order));
            return;
        }

        // 异步 HTTP 查询
        httpClient.prepareGet("http://dim-service/item/" + order.getItemId())
            .execute()
            .toCompletableFuture()
            .thenApply(resp -> parseJson(resp.getResponseBody()))
            .thenAccept(dim -> {
                cache.put(order.getItemId(), dim);
                order.setItemName(dim.getName());
                resultFuture.complete(Collections.singletonList(order));
            })
            .exceptionally(ex -> {
                resultFuture.complete(Collections.singletonList(order));
                return null;
            });
    }
}

同步 vs 异步性能对比(参考):

维度同步查询异步 I/O
请求模型串行阻塞并发非阻塞
单并行度吞吐100500 QPS(取决于 RT)500020000 QPS
CPU 利用率低(大量等待)高(I/O 等待期间处理其他请求)
外部系统压力低(请求慢但并发低)需控制 capacity 防止打挂下游
代码复杂度简单需异步客户端 + 回调

5. 时间语义与 Window

5.1 三种时间语义

事件发生 ──── 进入Flink ──── 被处理
   │              │              │
   ▼              ▼              ▼
Event Time   Ingestion Time  Processing Time
(事件时间)    (摄入时间)       (处理时间)
时间类型说明适用场景
Event Time事件自身携带的时间戳生产推荐,正确处理乱序数据
Ingestion Time事件进入 Flink Source 的时间较少使用
Processing Time算子处理事件的系统时间无需正确性保证的低延迟场景

5.2 Watermark 机制

Watermark 是一种衡量事件时间进展的机制,它告诉算子"小于这个时间戳的事件应该都已经到了"。

事件流(乱序):
  10:01  10:03  10:02  10:05  10:04  10:07  10:06
    │      │      │      │      │      │      │
    ▼      ▼      ▼      ▼      ▼      ▼      ▼
Watermark: W(10:00:55) → W(10:02:55) → W(10:04:55) → W(10:06:55)
                     (maxEventTime - 5s 允许延迟)

生成 Watermark:

// 方式一:forBoundedOutOfOrderness(固定延迟)
WatermarkStrategy<Order> strategy = WatermarkStrategy
    .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
    .withTimestampAssigner((order, ts) -> order.getEventTime())
    .withIdleness(Duration.ofMinutes(1)); // 空闲源检测

DataStream<Order> withTimestamps = orders
    .assignTimestampsAndWatermarks(strategy);

// 方式二:Kafka Source 中直接指定
KafkaSource<Order> kafkaSource = KafkaSource.<Order>builder()
    // ...
    .build();
DataStream<Order> stream = env.fromSource(
    kafkaSource,
    WatermarkStrategy.forBoundedOutOfOrderness(Duration.ofSeconds(5)),
    "Kafka Source");

Watermark 传递:

  • 多输入算子(如 keyBy 后 union)取所有输入 Watermark 的最小值
  • Watermark 周期性生成(默认 200ms,由 pipeline.auto-watermark-interval 控制)
  • 空闲分区/源用 withIdleness() 防止 Watermark 不推进

5.3 Window 类型

Tumbling Window(滚动窗口,无重叠)
│──────│──────│──────│──────│
0     10     20     30     40

Sliding Window(滑动窗口,有重叠)
│──────│
  │──────│
    │──────│
      │──────│
0     10     20     30     40
(窗口大小10s,滑动步长5s)

Session Window(会话窗口,间隙触发)
│───│    │────│        │──│
     gap  gap           gap
(相邻事件间隔超过gap则开新窗口)

Global Window(全局窗口,需自定义Trigger)
│──────────────────────────────│
(所有数据到同一个窗口,由Trigger决定何时触发)
// 滚动事件时间窗口
.window(TumblingEventTimeWindows.of(Time.seconds(10)))

// 滑动事件时间窗口
.window(SlidingEventTimeWindows.of(Time.seconds(10), Time.seconds(5)))

// 会话窗口
.window(EventTimeSessionWindows.withGap(Time.minutes(1)))

// 全局窗口
.window(GlobalWindows.create())
.trigger(CountTrigger.of(100)) // 每100条触发

5.4 Window Function

增量聚合(效率高,窗口内只存中间结果):

// ReduceFunction
DataStream<Order> result = keyedOrders
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .reduce((a, b) -> {
        a.setAmount(a.getAmount() + b.getAmount());
        return a;
    });

// AggregateFunction
DataStream<Double> avg = keyedOrders
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .aggregate(new AggregateFunction<Order, Tuple2<Double, Integer>, Double>() {
        @Override
        public Tuple2<Double, Integer> createAccumulator() {
            return Tuple2.of(0.0, 0);
        }
        @Override
        public Tuple2<Double, Integer> add(Order value, Tuple2<Double, Integer> acc) {
            return Tuple2.of(acc.f0 + value.getAmount(), acc.f1 + 1);
        }
        @Override
        public Double getResult(Tuple2<Double, Integer> acc) {
            return acc.f0 / acc.f1;
        }
        @Override
        public Tuple2<Double, Integer> merge(Tuple2<Double, Integer> a,
                                              Tuple2<Double, Integer> b) {
            return Tuple2.of(a.f0 + b.f0, a.f1 + b.f1);
        }
    });

全量聚合(ProcessWindowFunction,可访问窗口元数据):

DataStream<String> result = keyedOrders
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .process(new ProcessWindowFunction<Order, String, String, TimeWindow>() {
        @Override
        public void process(String key, Context ctx,
                            Iterable<Order> values,
                            Collector<String> out) {
            long count = 0;
            double sum = 0;
            for (Order o : values) {
                count++;
                sum += o.getAmount();
            }
            out.collect(String.format("Window [%d - %d] key=%s count=%d avg=%.2f",
                ctx.window().getStart(), ctx.window().getEnd(),
                key, count, sum / count));
        }
    });

最佳实践:增量聚合 + 全量窗口函数结合使用,兼顾性能和灵活性:

.reduce(new MyReduceFunction(), new MyProcessWindowFunction());

5.5 Allowed Lateness 与侧输出迟到数据

OutputTag<Order> lateTag = new OutputTag<Order>("late-orders"){};

SingleOutputStreamOperator<Order> result = keyedOrders
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .allowedLateness(Time.seconds(30))  // 窗口触发后再等30秒
    .sideOutputLateData(lateTag)        // 仍然迟到的数据进入侧输出
    .reduce(new MyReduceFunction());

DataStream<Order> lateData = result.getSideOutput(lateTag);
lateData.addSink(new LateDataSink()); // 单独处理迟到数据

6. 状态管理

6.1 状态的分类

Flink State
├── Keyed State(按 key 分区,只能在 KeyedStream 上使用)
│   ├── ValueState<T>        单个值
│   ├── ListState<T>         列表
│   ├── MapState<K,V>        Map(推荐,性能好)
│   ├── ReducingState<T>     增量聚合
│   └── AggregatingState<I,O> 聚合
├── Operator State(算子级别,与 key 无关)
│   └── ListState<T>
└── Broadcast State(广播状态,动态配置/规则下发)

6.2 Keyed State 示例

public class CountWithKeyedState extends RichFlatMapFunction<Order, String> {
    private transient ValueState<Long> countState;
    private transient MapState<String, Double> amountByItem;

    @Override
    public void open(Configuration parameters) {
        // ValueState
        ValueStateDescriptor<Long> countDesc = new ValueStateDescriptor<>(
            "count", Long.class);
        countState = getRuntimeContext().getState(countDesc);

        // MapState(推荐替代 ListState 做分类统计)
        MapStateDescriptor<String, Double> mapDesc = new MapStateDescriptor<>(
            "amount-by-item", String.class, Double.class);
        amountByItem = getRuntimeContext().getMapState(mapDesc);
    }

    @Override
    public void flatMap(Order order, Collector<String> out) throws Exception {
        Long count = countState.value();
        if (count == null) count = 0L;
        count++;
        countState.update(count);

        Double current = amountByItem.get(order.getItemId());
        if (current == null) current = 0.0;
        amountByItem.put(order.getItemId(), current + order.getAmount());

        out.collect("Total orders: " + count);
    }
}

6.3 State Backend

State Backend特点适用场景
HashMapStateBackend状态存堆内存,读写快,Checkpoint 到文件小状态(GB 级以内),低延迟
RocksDBStateBackend状态存本地 RocksDB(磁盘),支持超大状态大状态(TB 级),生产推荐
EmbeddedRocksDBStateBackendFlink 1.18+ 官方推荐配置名同上
// Flink 1.18+ 配置方式
Configuration config = new Configuration();
config.set(StateBackendOptions.STATE_BACKEND, "hashmap");
// 或 RocksDB:
config.set(StateBackendOptions.STATE_BACKEND, "rocksdb");
config.set(CheckpointingOptions.CHECKPOINT_STORAGE, "filesystem");
config.set(CheckpointingOptions.CHECKPOINTS_DIRECTORY,
    "hdfs:///flink/checkpoints");

env.configure(config);

6.4 状态 TTL

状态 TTL 可以自动清理过期状态,防止状态无限增长:

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(24))              // 24小时过期
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite) // 写入时刷新
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .cleanupInRocksdbCompactFilter(1000)      // RocksDB 压缩时清理
    .build();

ValueStateDescriptor<Long> desc = new ValueStateDescriptor<>("count", Long.class);
desc.enableTimeToLive(ttlConfig);

6.5 Broadcast State

用于低吞吐量的规则流广播到所有 Task:

// 规则流(如风控规则更新)
MapStateDescriptor<String, Rule> ruleStateDesc = new MapStateDescriptor<>(
    "rules", String.class, Rule.class);

BroadcastStream<Rule> ruleBroadcastStream = rules
    .broadcast(ruleStateDesc);

// 连接主流和广播流
BroadcastConnectedStream<Order, Rule> connected = orders
    .connect(ruleBroadcastStream);

connected.process(new BroadcastProcessFunction<Order, Rule, Alert>() {
    @Override
    public void processElement(Order order, ReadOnlyContext ctx,
                                Collector<Alert> out) throws Exception {
        ReadOnlyBroadcastState<String, Rule> state =
            ctx.getBroadcastState(ruleStateDesc);
        // 只读访问规则,进行匹配
        Rule rule = state.get(order.getItemId());
        if (rule != null && rule.matches(order)) {
            out.collect(new Alert(order));
        }
    }

    @Override
    public void processBroadcastElement(Rule rule, Context ctx,
                                         Collector<Alert> out) throws Exception {
        BroadcastState<String, Rule> state = ctx.getBroadcastState(ruleStateDesc);
        state.put(rule.getId(), rule); // 更新规则
    }
});

6.6 状态原语选择建议

需求推荐原因
存单个值ValueState最简单
按类别累加MapState比 ListState 遍历高效
窗口内累加ReducingState/AggregatingState增量计算,不存原始数据
大状态RocksDB + MapState磁盘存储,增量 Checkpoint

7. 容错机制

7.1 Checkpoint 原理

Checkpoint 是 Flink 实现容错的核心机制,基于 Chandy-Lamport 分布式快照算法:

Source 注入 Barrier
    │
    ▼
┌─────┐ Barrier_n
│Task1│─────────▶┌─────┐ Barrier_n
└─────┘           │Task2│─────────▶ Sink
                  └─────┘
每个 Task 收到 Barrier 后:
1. 对当前状态做快照(异步)
2. 将 Barrier 转发给下游
3. 快照完成后向 JobManager 确认
所有 Task 确认 → Checkpoint 完成

配置:

env.enableCheckpointing(60000); // 间隔60秒

CheckpointConfig ckConfig = env.getCheckpointConfig();
ckConfig.setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
ckConfig.setCheckpointTimeout(120000);        // 超时2分钟
ckConfig.setMinPauseBetweenCheckpoints(30000); // 两次CK最小间隔30秒
ckConfig.setMaxConcurrentCheckpoints(1);       // 最大并发1
ckConfig.setTolerableCheckpointFailureNumber(3); // 容忍3次失败
ckConfig.enableExternalizedCheckpoints(
    CheckpointConfig.ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

7.2 Checkpoint vs Savepoint

维度CheckpointSavepoint
触发方式自动(定时)手动
用途故障恢复升级/迁移/扩缩容
格式可能增量(RocksDB)全量标准格式
生命周期作业运行时自动管理手动管理,永久保留
状态兼容严格支持重命名/删除等
性能轻量,频繁重量,偶尔

7.3 Barrier 对齐与非对齐检查点

对齐检查点(Aligned Checkpoint):

  • 算子收到第一个 Barrier 后,阻塞该 Channel,等待其他 Channel 的 Barrier 到达
  • 问题:反压时对齐时间长,Checkpoint 超时

非对齐检查点(Unaligned Checkpoint, UC):

  • Barrier 直接穿越缓冲区,不等齐
  • 将 in-flight 数据一起做快照
  • 优势:反压场景下 Checkpoint 更快
  • 代价:Checkpoint 更大(包含缓冲数据),恢复时间可能更长
// 启用非对齐检查点
ckConfig.enableUnalignedCheckpoints();
// UC 触发超时(反压超过此时间后自动使用 UC)
ckConfig.setAlignedCheckpointTimeout(Duration.ofSeconds(30));

7.4 Exactly-Once vs At-Least-Once

语义说明代价
Exactly-Once每条数据恰好处理一次,不重不漏Checkpoint 对齐,延迟略高
At-Least-Once至少处理一次,可能重复不等 Barrier 对齐,延迟低

7.5 两阶段提交(2PC)与端到端 Exactly-Once

Flink 的端到端 Exactly-Once 需要 Source、Flink 引擎、Sink 三者配合:

阶段一(Pre-commit):
  Checkpoint Barrier 到达 Sink → Sink 开启事务写入外部系统
  → 提交事务预提交 → 通知 JobManager

阶段二(Commit):
  JobManager 收到所有 Task 的 CK 确认 → 通知所有 Task CK 完成
  → Sink 正式提交事务 → 数据对外可见

若失败:
  → Sink 回滚事务 → 从上一个 Checkpoint 恢复

支持 2PC 的 Sink: Kafka(事务写入)、MySQL(XA)、FileSink(原子重命名)、Paimon/Iceberg/Hudi。

// Kafka Sink 端到端 Exactly-Once
KafkaSink<String> sink = KafkaSink.<String>builder()
    .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
    .setTransactionalIdPrefix("flink-e2e-")
    // ...
    .build();

8. Flink SQL

8.1 Table API & SQL 简介

Flink SQL 基于 Apache Calcite,将 SQL 查询编译为 DataStream/DataSet 程序:

-- 纯 SQL 方式
CREATE TABLE orders (
  order_id BIGINT,
  item_id STRING,
  amount DECIMAL(10,2),
  order_time TIMESTAMP(3),
  WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'orders',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format' = 'json'
);

-- 每10秒滚动窗口统计每个商品的GMV
SELECT
  window_start,
  window_end,
  item_id,
  SUM(amount) AS gmv,
  COUNT(*) AS order_cnt
FROM TABLE(
  TUMBLE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '10' SECONDS)
)
GROUP BY window_start, window_end, item_id;

8.2 动态表与持续查询

Flink SQL 的核心抽象是动态表(Dynamic Table):流 → 动态表 → 持续查询 → 动态表 → 流。

Kafka Stream
    │
    ▼
┌──────────────┐
│  Dynamic Table│  (持续追加/更新)
└──────┬───────┘
       │ Continuous Query
       ▼
┌──────────────┐
│ Result Table │  (持续变化)
└──────┬───────┘
       │ toChangelogStream
       ▼
  Kafka/MySQL/ES Sink

8.3 时间属性

-- 处理时间(系统时间,无需 WATERMARK)
CREATE TABLE orders (
  order_id BIGINT,
  item_id STRING,
  amount DECIMAL(10,2),
  proc_time AS PROCTIME()  -- 计算列
) WITH (...);

-- 事件时间(需 WATERMARK)
CREATE TABLE orders (
  order_id BIGINT,
  item_id STRING,
  amount DECIMAL(10,2),
  order_time TIMESTAMP(3),
  WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (...);

8.4 Window TVF(Flink 1.13+ 推荐)

-- 滚动窗口
SELECT window_start, window_end, item_id, SUM(amount)
FROM TABLE(TUMBLE(TABLE orders, DESCRIPTOR(order_time), INTERVAL '10' SECONDS))
GROUP BY window_start, window_end, item_id;

-- 滑动窗口
SELECT window_start, window_end, item_id, SUM(amount)
FROM TABLE(HOP(TABLE orders, DESCRIPTOR(order_time),
               INTERVAL '5' SECONDS, INTERVAL '10' SECONDS))
GROUP BY window_start, window_end, item_id;

-- 累计窗口(Cumulate,Flink 1.13+)
SELECT window_start, window_end, item_id, SUM(amount)
FROM TABLE(CUMULATE(TABLE orders, DESCRIPTOR(order_time),
                    INTERVAL '1' MINUTES, INTERVAL '1' HOURS))
GROUP BY window_start, window_end, item_id;

-- 会话窗口(TVF 暂不支持,使用旧语法)
SELECT item_id, SESSION_START(order_time, INTERVAL '1' MINUTE), SUM(amount)
FROM orders
GROUP BY item_id, SESSION(order_time, INTERVAL '1' MINUTE);

8.5 Join 类型

Join 类型说明典型场景
Regular Join双流 Join,两边状态持续保留双流关联(需注意状态大小)
Interval Join时间区间内 Join,自动清理状态订单+支付关联
Temporal Join关联版本表(随时间变化的维表)汇率/费率历史关联
Lookup Join查询外部数据库维表MySQL/Redis 维度关联
Lateral Join表函数展开数组/JSON 展开
-- Interval Join
SELECT o.order_id, o.amount, p.pay_time
FROM orders o, payments p
WHERE o.order_id = p.order_id
  AND p.pay_time BETWEEN o.order_time AND o.order_time + INTERVAL '30' MINUTE;

-- Temporal Join(版本表)
SELECT o.order_id, o.amount, r.rate, o.amount * r.rate AS usd_amount
FROM orders o
JOIN rates FOR SYSTEM_TIME AS OF o.order_time AS r
  ON o.currency = r.currency;

-- Lookup Join(MySQL 维表)
CREATE TABLE dim_user (
  user_id BIGINT,
  user_name STRING,
  PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:mysql://localhost:3306/dim',
  'table-name' = 'user',
  'username' = 'root',
  'password' = 'xxx'
);

SELECT o.order_id, u.user_name, o.amount
FROM orders o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS u
  ON o.user_id = u.user_id;

8.6 Top-N 与去重

-- Top-N(每组取前N)
SELECT * FROM (
  SELECT *,
    ROW_NUMBER() OVER (
      PARTITION BY window_start, window_end
      ORDER BY gmv DESC
    ) AS rn
  FROM window_gmv
) WHERE rn <= 10;

-- 去重(保留最新一条)
SELECT * FROM (
  SELECT *,
    ROW_NUMBER() OVER (
      PARTITION BY order_id
      ORDER BY update_time DESC
    ) AS rn
  FROM orders
) WHERE rn = 1;

8.7 CDC 入湖

-- MySQL CDC Source
CREATE TABLE mysql_orders (
  id BIGINT,
  item_id STRING,
  amount DECIMAL(10,2),
  updated_at TIMESTAMP(3),
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = 'mysql',
  'port' = '3306',
  'username' = 'root',
  'password' = 'xxx',
  'database-name' = 'shop',
  'table-name' = 'orders',
  'server-time-zone' = 'Asia/Shanghai'
);

-- Paimon/Iceberg Sink(入湖)
CREATE TABLE lake_orders (
  id BIGINT,
  item_id STRING,
  amount DECIMAL(10,2),
  updated_at TIMESTAMP(3),
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'paimon',
  'warehouse' = 'hdfs:///lake/paimon',
  'database' = 'ods',
  'table' = 'orders',
  'primary-key' = 'id'
);

INSERT INTO lake_orders SELECT * FROM mysql_orders;

8.8 Catalog

-- Hive Catalog(持久化元数据)
CREATE CATALOG my_hive WITH (
  'type' = 'hive',
  'hive-conf-dir' = '/etc/hive/conf'
);

USE CATALOG my_hive;
SHOW DATABASES;
USE my_db;
SHOW TABLES;

8.9 Flink SQL 性能调优

Flink SQL 作业在生产中常遇到数据倾斜、状态过大、维表 Join 慢等问题。以下是系统化的调优手段。

8.9.1 MiniBatch 聚合

MiniBatch 将一批数据缓存后再触发聚合,大幅减少 State 访问次数,显著提升吞吐。

SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '5 s';   -- 最多等待5秒
SET 'table.exec.mini-batch.size' = '5000';            -- 最多缓存5000条

MiniBatch 会引入最多 allow-latency 的延迟,需在吞吐和延迟间权衡。

8.9.2 Local-Global 聚合(两阶段预聚合)

Local-Global 在 MiniBatch 基础上自动实现 MapReduce 式两阶段聚合,解决数据倾斜:

原始数据(key分布不均)
    │
    ▼
┌─────────────────┐
│  Local Aggregate │  每个并行度局部预聚合(加局部缓冲)
│  (Map端Combiner) │  输出: (key+localSalt, partialSum)
└────────┬────────┘
         │ Shuffle by key
         ▼
┌─────────────────┐
│ Global Aggregate │  全局最终聚合
│  (Reduce端)      │  输出: (key, finalSum)
└─────────────────┘
SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '5 s';
SET 'table.exec.mini-batch.size' = '5000';
SET 'table.optimizer.agg-phase-strategy' = 'TWO_PHASE';  -- 强制两阶段
-- AUTO: 自动选择; ONE_PHASE: 仅一阶段; TWO_PHASE: 强制两阶段
8.9.3 Split Distinct Aggregation

针对 COUNT(DISTINCT) 场景,将去重聚合拆分为两层,解决热点 key 倾斜:

-- 原始写法(热点 key 时单个 subtask 压力大)
SELECT item_id, COUNT(DISTINCT user_id) AS uv
FROM orders
GROUP BY item_id;

-- 开启 Split Distinct(自动改写为两阶段)
SET 'table.optimizer.distinct-agg.split.enabled' = 'true';
SET 'table.optimizer.distinct-agg.split.bucket-num' = '1024';
-- 第一层:按 item_id + (user_id % bucket) 分组去重
-- 第二层:按 item_id 聚合 SUM(distinct_count)
8.9.4 自动去重重写(Filtering Distinct Agg)

当聚合中同时包含 COUNT(*) 和 COUNT(DISTINCT) 且有相同过滤条件时,Flink 自动优化:

SET 'table.optimizer.filter-merge-acc.enabled' = 'true';
8.9.5 Top-N 缓存优化
-- 开启 Rank 重写优化,减少状态大小
SET 'table.optimizer.rank-rewrite.enabled' = 'true';

对 Top-N 查询,Flink 会优化状态访问策略:对排名靠后的数据不存入状态(因为几乎不会进入 Top-N),显著减少状态量。

8.9.6 维表 Join 优化
-- 异步 I/O Lookup Join(需 Connector 支持异步)
CREATE TABLE dim_user (
  user_id BIGINT,
  user_name STRING,
  PRIMARY KEY (user_id) NOT ENFORCED
) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:mysql://localhost:3306/dim',
  'table-name' = 'user',
  'username' = 'root',
  'password' = 'xxx',
  'lookup.async' = 'true',                    -- 开启异步查询
  'lookup.cache' = 'PARTIAL',                 -- 启用缓存
  'lookup.partial-cache.max-rows' = '10000',  -- 最大缓存行数
  'lookup.partial-cache.expire-after-write' = '10 min',  -- 写入后过期
  'lookup.partial-cache.expire-after-access' = '5 min',  -- 访问后过期
  'lookup.max-retries' = '3'                  -- 查询失败重试次数
);

-- Lookup Join 自动使用异步+缓存
SELECT o.order_id, u.user_name, o.amount
FROM orders AS o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS u
  ON o.user_id = u.user_id;
8.9.7 聚合下发(Projection Pushdown / Partition Prune)
-- Source Connector 支持谓词下推和投影下推时,自动减少读取量
SET 'table.optimizer.source.predicate-pushdown-enabled' = 'true';
SET 'table.optimizer.source.aggregate-pushdown-enabled' = 'true';
8.9.8 SQL 作业关键配置参数表
参数默认值说明
table.local-time-zonedefault本地时区,国内设为 Asia/Shanghai
table.exec.state.ttl0(不过期)聚合/Join 状态保留时间,如 24 h
table.exec.idle-state-retention0空闲状态保留时间(防止空闲 key 状态泄漏)
table.exec.mini-batch.enabledfalse是否开启 MiniBatch
table.exec.mini-batch.allow-latency0MiniBatch 延迟
table.exec.mini-batch.size-1MiniBatch 最大条数
table.optimizer.agg-phase-strategyAUTO聚合策略(AUTO/ONE_PHASE/TWO_PHASE)
table.optimizer.distinct-agg.split.enabledfalse拆分 COUNT(DISTINCT)
table.exec.resource.default-parallelism-1SQL 作业默认并行度
table.exec.sink.upsert-materializeAUTOUpsert Sink 是否物化(TABLE/SCAN/NONE/AUTO)
pipeline.name-作业名称
-- 生产环境推荐的基础配置
SET 'table.local-time-zone' = 'Asia/Shanghai';
SET 'table.exec.state.ttl' = '24 h';
SET 'table.exec.idle-state-retention' = '30 min';
SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '5 s';
SET 'table.exec.mini-batch.size' = '5000';
SET 'table.optimizer.agg-phase-strategy' = 'TWO_PHASE';

8.10 常用 Connector 详解

Flink SQL 通过 Connector 与外部系统交互,以下是生产常用 Connector 对比。

8.10.1 Connector 特性对比
ConnectorSourceSinkExactly-OnceCDC流批一体特点
Kafka✅✅✅ (2PC)❌流为主高吞吐消息队列,流处理标配
MySQL CDC✅❌✅✅流为主基于 Binlog,全量+增量
JDBC✅(Lookup)✅⚠️ 幂等/XA❌批为主通用关系库,批量写入
FileSystem✅✅✅ (原子提交)❌✅本地/HDFS/S3/OSS
Hive✅✅⚠️❌✅数仓迁移,批为主
Elasticsearch❌✅⚠️ 幂等❌流为主全文检索,幂等 Upsert
Redis✅(Lookup)✅❌ (At-Least-Once)❌流为主低延迟 KV 缓存/查询
Doris✅✅✅ (Stream Load 2PC)✅流+批MPP OLAP,实时导入
StarRocks✅✅✅ (Stream Load 2PC)✅流+批MPP OLAP,主键模型
ClickHouse❌✅⚠️ 幂等❌批为主列式 OLAP,分析性能强
Paimon✅✅✅✅✅流式数据湖,Flink 子项目
Iceberg✅✅✅⚠️✅数据湖格式,Schema 演进强
8.10.2 Flink CDC 入 Doris 完整示例
-- ODS: MySQL CDC Source
CREATE TABLE ods_orders (
  id BIGINT,
  order_no STRING,
  user_id BIGINT,
  item_id BIGINT,
  amount DECIMAL(10,2),
  order_status INT,
  create_time TIMESTAMP(3),
  update_time TIMESTAMP(3),
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = 'mysql',
  'port' = '3306',
  'username' = 'root',
  'password' = 'xxx',
  'database-name' = 'shop',
  'table-name' = 'orders',
  'server-time-zone' = 'Asia/Shanghai'
);

-- DWS: Doris Sink(Stream Load 2PC 精确一次)
CREATE TABLE dws_order_summary (
  item_id BIGINT,
  dt DATE,
  gmv DECIMAL(18,2),
  order_cnt BIGINT,
  PRIMARY KEY (item_id, dt) NOT ENFORCED
) WITH (
  'connector' = 'doris',
  'fenodes' = 'doris-fe:8030',
  'table.identifier' = 'dws.order_summary',
  'username' = 'root',
  'password' = 'xxx',
  'sink.label-prefix' = 'flink-doris',
  'sink.enable-2pc' = 'true',        -- 开启两阶段提交
  'sink.properties.format' = 'json',
  'sink.properties.read_json_by_line' = 'true'
);

-- 实时聚合写入 Doris
INSERT INTO dws_order_summary
SELECT
  item_id,
  DATE(create_time) AS dt,
  SUM(amount) AS gmv,
  COUNT(*) AS order_cnt
FROM ods_orders
WHERE order_status = 1
GROUP BY item_id, DATE(create_time);
8.10.3 Flink 写 StarRocks
CREATE TABLE dws_user_stats (
  user_id BIGINT,
  dt DATE,
  total_amount DECIMAL(18,2),
  order_count BIGINT,
  PRIMARY KEY (user_id, dt) NOT ENFORCED
) WITH (
  'connector' = 'starrocks',
  'jdbc-url' = 'jdbc:mysql://starrocks-fe:9030',
  'load-url' = 'starrocks-fe:8030;starrocks-be:8040',
  'database-name' = 'dws',
  'table-name' = 'user_stats',
  'username' = 'root',
  'password' = 'xxx',
  'sink.version' = 'V2',             -- V2 使用 Stream Load
  'sink.properties.format' = 'json',
  'sink.properties.strip_outer_array' = 'true'
);
8.10.4 Flink 写 ClickHouse
-- 方式一:通过 JDBC Connector 写入
CREATE TABLE ck_sink (
  event_date Date,
  item_id String,
  pv UInt64,
  uv UInt64
) WITH (
  'connector' = 'jdbc',
  'url' = 'jdbc:clickhouse://clickhouse:8123/analytics',
  'table-name' = 'item_stats',
  'username' = 'default',
  'password' = 'xxx',
  'sink.buffer-flush.max-rows' = '10000',
  'sink.buffer-flush.interval' = '5s',
  'sink.max-retries' = '3'
);

-- 方式二:使用官方 ClickHouse Connector(支持 Sink)
CREATE TABLE ck_sink_v2 (
  event_date Date,
  item_id String,
  pv UInt64,
  uv UInt64
) WITH (
  'connector' = 'clickhouse',
  'url' = 'clickhouse:8123',
  'database-name' = 'analytics',
  'table-name' = 'item_stats',
  'username' = 'default',
  'password' = 'xxx',
  'sink.batch-size' = '10000',
  'sink.flush-interval' = '5s'
);
8.10.5 Flink 写 Redis
-- Redis Sink
CREATE TABLE redis_sink (
  item_id STRING,
  item_name STRING,
  category STRING,
  PRIMARY KEY (item_id) NOT ENFORCED
) WITH (
  'connector' = 'redis',
  'mode' = 'single',                -- single / cluster / sentinel
  'host' = 'redis',
  'port' = '6379',
  'password' = 'xxx',
  'command' = 'HMSET',              -- 写入命令: SET/HSET/ LPUSH 等
  'redis.value.data-type' = 'hash'  -- string / hash / list / set / zset
);

-- Redis Lookup Source(维表关联)
CREATE TABLE redis_dim (
  item_id STRING,
  item_name STRING,
  category STRING,
  PRIMARY KEY (item_id) NOT ENFORCED
) WITH (
  'connector' = 'redis',
  'mode' = 'single',
  'host' = 'redis',
  'port' = '6379',
  'command' = 'HGETALL',
  'lookup.cache.max-rows' = '5000',
  'lookup.cache.ttl' = '300s'
);
8.10.6 Flink 写 Elasticsearch
CREATE TABLE es_sink (
  order_id STRING,
  user_id BIGINT,
  item_name STRING,
  amount DECIMAL(10,2),
  order_time TIMESTAMP(3),
  PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
  'connector' = 'elasticsearch-7',
  'hosts' = 'http://es-node1:9200;http://es-node2:9200',
  'index' = 'orders',
  'document-id.key-delimiter' = '_',
  'sink.bulk-flush.max-actions' = '1000',
  'sink.bulk-flush.max-size' = '5mb',
  'sink.bulk-flush.interval' = '2s',
  'sink.bulk-flush.backoff.strategy' = 'EXPONENTIAL',
  'sink.bulk-flush.backoff.max-retries' = '5',
  'connection.username' = 'elastic',
  'connection.password' = 'xxx'
);
8.10.7 Sink Exactly-Once 支持情况
Sink一致性保证机制
KafkaExactly-Once事务(Kafka 0.11+),2PC 提交
DorisExactly-OnceStream Load 2PC,Checkpoint 时预提交,完成时正式提交
StarRocksExactly-OnceStream Load 事务标签(label),2PC 模式
Paimon/Iceberg/HudiExactly-Once数据湖原子提交(Snapshot 提交)
FileSystemExactly-OncePending 文件 → Finished 文件的原子重命名
ElasticsearchAt-Least-Once幂等 Upsert(按文档 ID),失败重试不产生重复
JDBC(MySQL)At-Least-Once幂等 Upsert(INSERT ... ON DUPLICATE KEY UPDATE)或 XA 事务
RedisAt-Least-Once无事务,需业务层保证幂等(如 SET 覆盖写)
ClickHouseAt-Least-OnceReplacingMergeTree 去重或幂等批量写入

关键结论:端到端 Exactly-Once 需要 Source 支持重放、Flink 引擎 Checkpoint、Sink 支持 2PC 或幂等写入三者配合。


9. Flink CDC

9.1 什么是 CDC

CDC(Change Data Capture) 即变更数据捕获,实时捕获数据库的增删改操作,用于数据同步、缓存失效、实时数仓等。

CDC 方案原理代表
基于查询定时 SELECT,无法捕获删除,延迟高DataX 离线同步
基于日志解析 Binlog/WAL,实时、可靠Debezium、Flink CDC、Canal

9.2 Flink CDC 演进

版本核心特性
1.x基于 Debezium,单并发,全量锁表
2.x无锁算法、并行读取、断点续传、整库同步
3.x架构升级(Pipeline 连接器)、Schema Evolution 支持、增量快照框架、整库同步到 Kafka/Paimon

9.3 全量 + 增量同步

Flink CDC 2.x 的增量快照算法:

  1. 全量阶段:按主键分片(Chunk),并行读取,无锁(使用可重复读事务 + Binlog 位点位)
  2. 增量阶段:全量完成后,无缝切换到 Binlog 流式读取
  3. 断点续传:Chunk 粒度的 Checkpoint,失败可从上次位置恢复
全量阶段(并行)              增量阶段(单并发)
┌────┬────┬────┬────┐        Binlog Stream
│ C1 │ C2 │ C3 │ C4 │  ───▶  INSERT/UPDATE/DELETE
└────┴────┴────┴────┘        (实时流)
  无锁,并行读取,CK记录已完成Chunk

9.4 整库同步(Flink CDC 3.x)

# pipeline.yaml(Flink CDC 3.x)
source:
  type: mysql
  name: MySQL Source
  hostname: localhost
  port: 3306
  username: root
  password: xxx
  tables: shop\..*

sink:
  type: paimon
  name: Paimon Sink
  warehouse: hdfs:///lake/paimon

pipeline:
  name: Sync MySQL to Paimon
  parallelism: 4
# 提交整库同步任务
./bin/flink-cdc.sh pipeline.yaml

10. Flink CEP

10.1 CEP 概念

CEP(Complex Event Processing) 复杂事件处理,在事件流中检测满足特定模式的事件序列,广泛用于风控、监控、欺诈检测。

10.2 Pattern API

// 场景:10秒内同一用户连续两次大额交易(>10000)触发告警
Pattern<Order, ?> pattern = Pattern.<Order>begin("first")
    .where(new SimpleCondition<Order>() {
        @Override
        public boolean filter(Order order) {
            return order.getAmount() > 10000;
        }
    })
    .next("second")  // 严格紧邻(使用 followedBy 为非紧邻)
    .where(new SimpleCondition<Order>() {
        @Override
        public boolean filter(Order order) {
            return order.getAmount() > 10000;
        }
    })
    .within(Time.seconds(10)); // 10秒内

PatternStream<Order> patternStream = CEP.pattern(
    orders.keyBy(Order::getUserId), pattern);

patternStream.select(new PatternSelectFunction<Order, Alert>() {
    @Override
    public Alert select(Map<String, List<Order>> pattern) {
        Order first = pattern.get("first").get(0);
        Order second = pattern.get("second").get(0);
        return new Alert("可疑交易: 用户 " + first.getUserId()
            + " 在10秒内两次大额交易");
    }
}).addSink(new AlertSink());

10.3 Pattern 量词与循环

// times: 出现 N 次
pattern.times(3);                    // 恰好3次
pattern.timesOrMore(3);              // 3次或更多
pattern.times(2, 4);                 // 2-4次

// 连续性
pattern.begin("a").where(...)
    .next("b")        // 严格紧邻:a 和 b 之间不能有其他事件
    .followedBy("c")  // 宽松紧邻:a 和 b 之间可以有其他事件
    .followedByAny("d") // 更宽松,允许跳过已匹配事件
    .notNext("e")     // 不出现(严格)
    .notFollowedBy("f"); // 不出现(宽松)

// within: 时间约束
pattern.within(Time.minutes(5));

10.4 超时事件处理

PatternStream<Order> patternStream = CEP.pattern(keyedOrders, pattern);

OutputTag<String> timeoutTag = new OutputTag<String>("timeout"){};

SingleOutputStreamOperator<Alert> result = patternStream
    .flatSelect(
        timeoutTag,
        new PatternFlatTimeoutFunction<Order, String>() {
            @Override
            public void timeout(Map<String, List<Order>> pattern,
                                long timeoutTimestamp,
                                Collector<String> out) {
                Order first = pattern.get("first").get(0);
                out.collect("用户 " + first.getUserId() + " 交易模式超时未完成");
            }
        },
        new PatternFlatSelectFunction<Order, Alert>() {
            @Override
            public void flatSelect(Map<String, List<Order>> pattern,
                                   Collector<Alert> out) {
                // 正常匹配
            }
        }
    );

DataStream<String> timeoutAlerts = result.getSideOutput(timeoutTag);

11. Flink 性能调优

11.1 反压定位与处理

什么是反压: 下游处理速度跟不上上游,导致 Buffer 填满,逐级向上传播。

定位手段:

手段方法
Web UI查看每个算子的 BackPressured 和 Busy 指标,红色为反压
MetricsoutPoolUsage、inPoolUsage、numRecordsInPerSecond
线程栈jstack <TaskManager PID>,查看线程在做什么(GC?锁?慢方法?)
火焰图jfr 或 async-profiler 生成火焰图定位热点方法

处理方法:

  1. 提高慢算子并行度
  2. 优化慢算子逻辑(数据库查询、外部 HTTP 调用异步化)
  3. 使用异步 I/O(AsyncDataStream)
  4. 增加网络 Buffer
  5. 数据倾斜处理

11.2 并行度设置

Source 并行度 = Kafka Partition 数(或其约数)
算子并行度 = 根据吞吐和单并行度处理能力调整
Sink 并行度 = 根据下游写入能力

经验公式: 并行度 = 总数据量 / (单并行度处理速率 × 时间窗口)

建议 Source 并行度与 Kafka Partition 数一致,避免消费不均。

11.3 Slot 配置

# flink-conf.yaml
taskmanager.numberOfTaskSlots: 4          # 每个 TM 的 Slot 数
taskmanager.memory.process.size: 4096m    # TM 总内存
taskmanager.memory.managed.fraction: 0.4  # Managed Memory 占比
taskmanager.memory.network.fraction: 0.15 # Network Memory 占比

Slot 数一般设为 CPU 核数;大状态场景可适当减少 Slot 数以增加每 Slot 内存。

11.4 Checkpoint 调优

ckConfig.setCheckpointInterval(60000);  // 不要太短,1-5分钟
ckConfig.setMinPauseBetweenCheckpoints(30000);
ckConfig.enableUnalignedCheckpoints();   // 反压场景开启
ckConfig.setAlignedCheckpointTimeout(Duration.ofSeconds(30));

// RocksDB 增量 Checkpoint
config.set(CheckpointingOptions.INCREMENTAL_CHECKPOINTS, true);
// RocksDB 压缩过滤清理 TTL
config.set(StateBackendOptions.ROCKSDB_COMPRESSION, "zstd");

11.5 数据倾斜处理

现象: 某些 SubTask 的数据量远大于其他,导致整体吞吐瓶颈。

方案一:加随机前缀两阶段聚合

// 第一阶段:加随机前缀,局部聚合
DataStream<Tuple2<String, Double>> withSalt = keyedStream
    .map(order -> {
        int salt = new Random().nextInt(10); // 0-9 随机前缀
        return Tuple2.of(salt + "_" + order.getItemId(), order.getAmount());
    })
    .keyBy(t -> t.f0)
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .reduce((a, b) -> Tuple2.of(a.f0, a.f1 + b.f1));

// 第二阶段:去掉前缀,全局聚合
DataStream<Tuple2<String, Double>> result = withSalt
    .map(t -> {
        String key = t.f0.substring(t.f0.indexOf("_") + 1);
        return Tuple2.of(key, t.f1);
    })
    .keyBy(t -> t.f0)
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .reduce((a, b) -> Tuple2.of(a.f0, a.f1 + b.f1));

方案二:Local-Global 聚合(Flink SQL 内置)

SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '5s';
SET 'table.exec.mini-batch.size' = '5000';
SET 'table.optimizer.agg-phase-strategy' = 'TWO_PHASE';
-- 自动启用 Local-Global 两阶段聚合

方案三:MiniBatch

-- 开启 MiniBatch,减少 State 访问次数
SET 'table.exec.mini-batch.enabled' = 'true';
SET 'table.exec.mini-batch.allow-latency' = '1 s';
SET 'table.exec.mini-batch.size' = '5000';

11.6 JVM 调优

# flink-conf.yaml
env.java.opts.taskmanager: >
  -XX:+UseG1GC
  -XX:MaxGCPauseMillis=200
  -XX:InitiatingHeapOccupancyPercent=45
  -XX:+ParallelRefProcEnabled
  -XX:+UnlockExperimentalVMOptions
  -XX:G1NewSizePercent=10
  -XX:G1MaxNewSizePercent=20

大状态 + RocksDB 场景,关注堆外内存配置,避免 Direct Memory OOM。

11.7 反压深度原理

反压(Backpressure)是流处理系统中最常见的问题之一。当某个算子处理速度跟不上上游发送速度时,反压会沿数据流图向上传播,最终导致整条链路吞吐下降。

Credit-based 流控机制:

Flink 1.5+ 引入 Credit-based 流控,替代了原来的纯 TCP 反压,实现更精细的缓冲区管理:

上游 Task (ResultSubpartition)          下游 Task (InputChannel)
┌──────────────────────┐              ┌──────────────────────┐
│  Buffer Pool         │   1. 发送     │  Buffer Pool         │
│  ┌──┐┌──┐┌──┐       │  Credit(n)   │  ┌──┐┌──┐           │
│  │B1││B2││B3│  ─────┼─────────────▶│  │  ││  │ (有n个空闲  │
│  └──┘└──┘└──┘       │              │  └──┘└──┘  Buffer)   │
│  等待Credit...       │   2. 消费完  │                      │
│         ▲            │◀─────────────│  归还 Credit + 反馈   │
│         │            │   Credit(n)  │  可用Buffer数         │
│  有Credit才发送Buffer│              │                      │
└──────────────────────┘              └──────────────────────┘

核心流程:

  1. 下游 InputChannel 有空闲 Buffer 时,向上游发送 Credit(表示"我能接收 n 个 Buffer")
  2. 上游只有收到 Credit 后才发送数据 Buffer;Credit 用完则停止发送
  3. 下游消费完 Buffer 后归还 Buffer 到 BufferPool,并发送新的 Credit
  4. 如果下游处理慢,Credit 长期不归还,上游自然停止发送 → 反压逐级传播

反压传播链路:

Sink (慢,写DB瓶颈)
  │ Buffer满,不发Credit
  ▼
Map/Window (无法往下写,Buffer满)
  │ Buffer满,不发Credit
  ▼
KeyBy/Shuffle (无法往下写)
  │ Buffer满,不发Credit
  ▼
Source (无法往下写,消费Kafka变慢)
  │ Kafka Consumer.poll() 阻塞或暂停
  ▼
Kafka Lag 增长

反压从最下游瓶颈点开始,沿拓扑反向传播至 Source。

Web UI 定位反压:

在 Flink Web UI 的 Job 页面,点击任意算子顶点查看 BackPressure 标签:

状态含义阈值
OK无反压,运行正常反压时间 < 10%
LOW轻微反压10% ~ 50%
HIGH严重反压> 50%

Web UI 还展示每个 SubTask 的 backPressured / busy / idle 时间占比,快速定位瓶颈 SubTask。

Metrics 定位:

Metric类型说明
outputQueueLengthGauge输出缓冲区队列长度,持续高说明下游反压
inputQueueLengthGauge输入缓冲区队列长度
backPressuredTimeMsPerSecondMeter每秒反压时间(ms),接近 1000 表示完全反压
busyTimeMsPerSecondMeter每秒繁忙时间,接近 1000 表示算子满负荷
idleTimeMsPerSecondMeter每秒空闲时间
credit (Flink 1.15+)Gauge可用 Credit 数,为 0 表示被反压
availableMemorySegmentsGauge可用网络内存段数
numRecordsInPerSecondMeter每秒输入记录数
numRecordsOutPerSecondMeter每秒输出记录数

Flink 1.15+ Buffer Pool 反压指标:

Flink 1.15 重构了网络栈指标,新增基于 Buffer Pool 的细粒度监控:

  • outputBufferPoolUsage:输出缓冲池使用率
  • inputBufferPoolUsage:输入缓冲池使用率
  • floatingBuffersUsage / exclusiveBuffersUsage:浮动/独占 Buffer 使用情况

常见原因与解决:

原因识别特征解决方法
数据倾斜个别 SubTask busy=100%,其他空闲Local-Global 聚合、加随机前缀两阶段聚合、Key 重分布
GC 问题周期性反压,GC 日志停顿时间长调优 G1GC、增大堆内存、切换 RocksDB 减少堆内存
慢 SinkSink 算子 busy=100%,反压从 Sink 开始批量写入、异步写入、增加 Sink 并行度、优化下游系统
Checkpoint 对齐CK 期间出现周期性反压开启 Unaligned Checkpoint、增大 CK 间隔
资源不足所有 SubTask 都 busy=100%增加并行度、扩容 TaskManager
外部系统慢维表查询/HTTP 调用延迟高异步 I/O + 缓存、连接池、降级

11.8 Watermark Alignment(1.16+)

在多源或多分区场景中,不同 Source/分区的 Watermark 推进速度可能差异巨大,导致算子需要缓存大量"等慢源"的状态,造成状态膨胀和内存压力。Watermark Alignment 通过对齐各源的 Watermark 来解决此问题。

问题场景:

Source A (快): WM = 10:30:00  ──────────────────────▶
                                              Union/Join
Source B (慢): WM = 10:00:00  ──────────────────────▶

结果:全局 WM = min(10:30, 10:00) = 10:00
问题:A 已到达 10:30 的数据需要在状态中缓存30分钟等待B
      → 状态膨胀、Checkpoint 增大、内存压力

对齐原理:

Watermark Alignment 会暂停消费"过快"的 Source/分区,等待慢 Source 追上来,使各源的 Watermark 偏差控制在阈值内:

Source A: WM=10:05  暂停读取(等待B)     ▶ 等待...
Source B: WM=10:00  正常读取追赶        ▶ WM=10:04 → 10:05
                                      对齐!
Source A: 恢复读取  ──────────────────▶ WM=10:06
Source B: 正常读取  ──────────────────▶ WM=10:06

DataStream API 配置:

WatermarkStrategy<Order> strategy = WatermarkStrategy
    .<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
    .withTimestampAssigner((order, ts) -> order.getEventTime())
    .withWatermarkAlignment(
        "alignment-group-1",            // 对齐组名(同组内对齐)
        Duration.ofSeconds(20),          // maxAllowedWatermarkDrift: 最大允许偏差
        Duration.ofSeconds(5)            // updateInterval: 检查间隔(默认1s)
    );

DataStream<Order> stream = env.fromSource(kafkaSource, strategy, "Kafka Source");

参数说明:

参数说明建议
maxAllowedWatermarkDrift各源/分区之间 Watermark 的最大允许偏差根据业务延迟容忍度,通常 10~60s
updateInterval对齐检查和调整的频率默认 1s,通常无需修改

不同 Source 配置相同的 alignment group 名称即可跨源对齐;同一 Kafka Source 的不同分区也会自动对齐。

SQL 中的配置:

-- 表级别配置 Watermark Alignment
CREATE TABLE orders (
  order_id BIGINT,
  order_time TIMESTAMP(3),
  WATERMARK FOR order_time AS order_time - INTERVAL '5' SECOND
) WITH (
  'connector' = 'kafka',
  'topic' = 'orders',
  'properties.bootstrap.servers' = 'kafka:9092',
  'format' = 'json',
  'scan.watermark.alignment.group' = 'alignment-group-1',
  'scan.watermark.alignment.max-drift' = '20 s',
  'scan.watermark.alignment.update-interval' = '5 s'
);

-- 全局默认配置(flink-conf.yaml)
# table.exec.source.watermark-alignment.group: default-group
# table.exec.source.watermark-alignment.max-drift: 20 s

适用场景与注意事项:

  • ✅ 多源 Join / Union 时各源延迟差异大
  • ✅ Kafka 分区数据倾斜(部分分区数据量远大于其他)
  • ✅ 状态膨胀严重且与 Watermark 不同步直接相关
  • ⚠️ 会降低快源的吞吐量(暂停消费),需权衡
  • ⚠️ 不适合各源本身延迟差异就应该存在的场景(如事实流 + 慢维表流)

11.9 Metrics 与监控

完善的监控是生产 Flink 作业稳定运行的保障。Flink 提供了丰富的 Metrics 体系,可对接 Prometheus + Grafana 实现可视化和告警。

Metric 类型:

类型说明示例
Counter计数器,单调递增/递减numRecordsIn、numRecordsOut
Gauge瞬时值,任意时刻的数值outputQueueLength、currentInputWatermark
Histogram分布统计(均值/分位数/最大值)记录处理延迟分布
Meter速率统计(每秒平均值)numRecordsInPerSecond、numRecordsOutPerSecond

关键 Metrics 列表:

分类Metric说明
吞吐numRecordsInPerSecond每秒输入记录数
numRecordsOutPerSecond每秒输出记录数
numBytesInPerSecond每秒输入字节数(Kafka Source)
反压backPressuredTimeMsPerSecond每秒反压时间(ms)
busyTimeMsPerSecond每秒繁忙时间
outputQueueLength输出队列长度
CheckpointlastCheckpointDuration最近一次 CK 耗时(ms)
lastCheckpointSize最近一次 CK 大小(字节)
numberOfCompletedCheckpoints已完成 CK 数
numberOfFailedCheckpoints失败 CK 数
lastCheckpointRestoreTimestamp最近恢复时间
状态stateSize (RocksDB)状态大小
rocksdb.block-cache-usageRocksDB 块缓存使用量
rocksdb.mem-table-flush-pending待刷新 MemTable 数
资源Status.JVM.Memory.Heap.Used堆内存使用量
Status.JVM.Memory.NonHeap.Used非堆内存使用量
Status.JVM.GarbageCollector.*.TimeGC 时间
outOfMemoryErrorOOM 计数器
Kafka 消费consumer-lagKafka 消费积压(消息数)
consumer-records-lag-max最大消费延迟
WatermarkcurrentInputWatermark当前输入 Watermark
currentOutputWatermark当前输出 Watermark

Prometheus + Grafana 集成:

# flink-conf.yaml - Prometheus Reporter 配置
metrics.reporter.prom.factory.class: org.apache.flink.metrics.prometheus.PrometheusReporterFactory
metrics.reporter.prom.port: 9999
metrics.reporter.prom.filterLabelValueCharacters: true

# 或使用 Prometheus PushGateway(主动推送)
metrics.reporter.promgateway.factory.class: org.apache.flink.metrics.prometheus.PrometheusPushGatewayReporterFactory
metrics.reporter.promgateway.host: pushgateway
metrics.reporter.promgateway.port: 9091
metrics.reporter.promgateway.jobName: flink-job
metrics.reporter.promgateway.interval: 15 SECONDS
# prometheus.yml - Prometheus 抓取配置
scrape_configs:
  - job_name: 'flink-taskmanager'
    kubernetes_sd_configs:
      - role: pod
    relabel_configs:
      - source_labels: [__meta_kubernetes_pod_label_app]
        regex: flink-taskmanager
        action: keep
      - source_labels: [__address__]
        regex: '([^:]+):(\d+)'
        target_label: __address__
        replacement: '${1}:9999'

  - job_name: 'flink-jobmanager'
    static_configs:
      - targets: ['flink-jobmanager:9999']

推荐监控大盘面板(Grafana):

面板指标展示方式
作业概览作业状态、运行时长、并行度状态卡片
吞吐量numRecordsInPerSecond / numRecordsOutPerSecond折线图
反压backPressuredTimeMsPerSecond / busyTimeMsPerSecond热力图/折线图
CheckpointCK 时长、CK 大小、成功/失败次数柱状图+状态
状态大小stateSize / RocksDB 指标折线图
内存Heap/NonHeap/Managed/Direct 使用率面积图
GCGC 时间/次数折线图
Kafka Lagconsumer-lag折线图
Watermark 延迟System.time() - currentInputWatermark折线图
容器资源CPU/内存使用率(K8s 指标)面积图

告警规则建议:

# Prometheus AlertManager 规则示例
groups:
  - name: flink-alerts
    rules:
      # Checkpoint 连续失败
      - alert: FlinkCheckpointFailure
        expr: increase(flink_jobmanager_job_numberOfFailedCheckpoints[5m]) > 3
        for: 2m
        labels:
          severity: critical
        annotations:
          summary: "Flink作业 {{ $labels.job_name }} Checkpoint连续失败"

      # 反压严重
      - alert: FlinkBackpressureHigh
        expr: flink_taskmanager_job_task_backPressuredTimeMsPerSecond > 800
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Flink任务 {{ $labels.task_name }} 反压严重"

      # Kafka 消费积压
      - alert: FlinkKafkaLagHigh
        expr: flink_taskmanager_job_task_operator_kafka_consumer_lag > 100000
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "Kafka消费积压超过10万条"

      # GC 时间过长
      - alert: FlinkGCHigh
        expr: rate(jvm_gc_collection_seconds_sum[5m]) > 0.3
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "TaskManager GC时间占比超过30%"

      # 容器内存使用率高
      - alert: FlinkContainerMemoryHigh
        expr: container_memory_working_set_bytes / container_spec_memory_limit_bytes > 0.85
        for: 5m
        labels:
          severity: warning
        annotations:
          summary: "TaskManager内存使用率超过85%"

      # 作业失败
      - alert: FlinkJobFailed
        expr: flink_jobmanager_job_status == 0
        for: 1m
        labels:
          severity: critical
        annotations:
          summary: "Flink作业 {{ $labels.job_name }} 已失败"

Flink Web UI 监控要点:

  1. Overview 页:查看作业状态、并行度、运行时长
  2. Vertices 页:每个算子的 Records Received/Sent、Bytes Received/Sent、BackPressured、Busy
  3. BackPressure 页:查看每个 SubTask 的反压状态(OK/LOW/HIGH)
  4. Checkpoints 页:CK 历史、时长、大小、对齐时间;检查 Savepoint 路径
  5. Task Managers 页:每个 TM 的 Slot 使用、内存、GC 情况
  6. Metrics 页:可直接在 Web UI 添加自定义 Metric 图表,无需 Grafana 也能快速排查

12. Flink 常见问题与踩坑

12.1 Watermark 不推进

原因:

  • 某个 Source 分区/SubTask 长时间无数据(空闲源)
  • 多输入取最小值,一个慢输入拖慢全局

解决:

WatermarkStrategy.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(5))
    .withIdleness(Duration.ofMinutes(1)); // 检测空闲源

12.2 Checkpoint 超时/失败

排查方向:

  1. 反压导致 Barrier 对齐慢 → 开启 Unaligned Checkpoint
  2. 状态太大,快照时间长 → 增量 Checkpoint、RocksDB、调大超时
  3. 存储(HDFS/S3)写入慢 → 检查存储性能
  4. 数据倾斜导致个别 Task 状态过大 → 解决倾斜

12.3 状态过大 OOM

  • 使用 RocksDB 替代 HashMap 后端
  • 配置状态 TTL,自动清理
  • 减少 Keyed State 中的字段
  • Regular Join 设置状态 TTL(Flink SQL)
SET 'table.exec.state.ttl' = '24 h';

12.4 Kafka 消费积压

  1. 提高 Source 并行度(不超过 Partition 数)
  2. 增加 Partition 数
  3. 优化下游处理逻辑
  4. 检查是否有数据倾斜
  5. 临时扩容 TaskManager

12.5 时区问题

-- Flink SQL 中统一使用 UTC 还是北京时间?
-- 建议在 Source 中指定时区
CREATE TABLE kafka_source (
  ts TIMESTAMP(3)
) WITH (
  'connector' = 'kafka',
  ...
  'properties.server-time-zone' = 'Asia/Shanghai'  -- MySQL CDC
);

-- 时间转换
SELECT
  TO_TIMESTAMP_LTZ(UNIX_TIMESTAMP(ts) * 1000, 3) AS bj_time
FROM source;

12.6 序列化问题

Flink 优先使用 POJO 序列化(高效、支持 Schema 演进),否则回退到 Kryo(慢、不支持状态迁移)。

POJO 规则:

  • Public 类
  • 无参构造函数
  • 所有字段 Public 或有 getter/setter
  • 字段类型均受 Flink 支持
// ✅ 正确的 POJO
public class Order {
    public long id;
    public String itemId;
    public double amount;
    public long eventTime;

    public Order() {} // 必须有无参构造

    // getters/setters...
}

// ❌ 避免使用匿名类、Lambda 中的复杂捕获

12.7 并行度变化状态不兼容

  • 使用 uid() 为算子显式指定 UID,确保 Savepoint 恢复时能匹配
stream.keyBy(...).window(...).reduce(...).uid("window-reducer");
  • 改变并行度时,RocksDB 后端支持 Rescale(自动重新分配 KeyGroup)
  • HashMap 后端在 1.18+ 也支持 Rescale

13. Flink 1.18-1.20 新特性与 2.0 展望

13.1 Flink 1.18 新特性

特性说明
Flink CDC 3.0Pipeline 架构、Schema Evolution、整库同步
Lakehouse 集成Paimon/Iceberg/Hudi 原生支持增强
Runtime非对齐检查点改进,Watermark 优化
SQLWindow TVF 增强,Top-N 优化
Python APIPyFlink 性能提升,更多 SQL 功能

13.2 Flink 1.19 新特性

特性说明
自适应查询执行(AQE)运行时根据数据量动态调整 Join 策略、聚合策略
动态表配置运行时动态修改参数,无需重启
RocksDB 优化内存管理改进,Checkpoint 更快
SQL 增强更多函数、NOT NULL 约束、Materialized Table
Kubernetes Operator自动扩缩容、滚动更新改进

13.3 Flink 1.20 新特性

特性说明
Disaggregated State Management 预览状态与计算分离,状态存远端存储(ForSt DB)
流式数仓增强Paimon 深度集成,增量物化视图
Python API 增强PyFlink 与 SQL 功能对齐
性能优化向量化计算、JVM 优化

13.4 Flink 2.0 方向

┌─────────────────────────────────────────────────┐
│                 Flink 2.0 愿景                   │
├─────────────────────────────────────────────────┤
│  1. Flink SQL 成为一等公民                       │
│     - Streaming Lakehouse 统一 API               │
│     - 批流一体在 SQL 层完全融合                   │
│                                                 │
│  2. Disaggregated State Management              │
│     - 计算与状态分离,状态存分布式存储             │
│     - 秒级扩缩容,不再依赖本地磁盘                │
│                                                 │
│  3. 动态扩缩容(Autoscaling)                    │
│     - 根据负载自动调整并行度                      │
│     - K8s Operator 原生支持                      │
│                                                 │
│  4. 流式数仓(Streaming Lakehouse)              │
│     - Flink + Paimon 构建实时湖仓               │
│     - 增量物化视图、CDC 入湖一体                 │
│                                                 │
│  5. 性能与稳定性                                 │
│     - 向量化执行引擎                             │
│     - 更好的资源利用率                           │
└─────────────────────────────────────────────────┘

13.5 Lakehouse 集成

Flink 与三大数据湖格式的集成:

格式Flink 集成特点
Paimon原生深度集成(Flink 子项目)流批一体、主键表、CDC 入湖、LSM 架构
IcebergFlink Connector 完善Schema 演进强、ACID、适合批为主
HudiFlink Connector 支持增量视图、主键更新成熟、MoR 表

13.6 Flink + Paimon 深度集成

Apache Paimon(原 Flink Table Store)是 Flink 官方子项目,定位为流式数据湖存储,将湖存储的流批一体能力发挥到极致。

Paimon 核心定位:

┌──────────────────────────────────────────────────┐
│                  Paimon 定位                      │
├──────────────────────────────────────────────────┤
│  流式数据湖存储(Streaming Lakehouse)            │
│  - 主键表实时 Upsert(LSM 架构,毫秒级可见)       │
│  - 原生 Changelog 产出(支持 CDC 入湖/出湖)      │
│  - 流批一体读写(流式写入 + 批量查询)             │
│  - Lookup Join 高性能点查(主键索引)              │
│  - 与 Flink SQL 深度集成                          │
└──────────────────────────────────────────────────┘

核心概念:

概念说明
Primary Key Table主键表,支持 Upsert/Delete,底层 LSM Tree,支持 Changelog 产出
Append Table追加表,仅支持追加写入,适合日志/事件类数据,无主键
Bucket数据分桶,主键表按 Bucket 分布(类似分区+哈希),每个 Bucket 是一个 LSM 引擎
Changelog ProducerChangelog 生成方式:none(仅写入)、input(依赖输入 Changelog)、lookup(Lookup 补全)、full-compaction(全量压缩产出)
Snapshot快照,每次提交生成一个 Snapshot,支持 Time Travel 和增量读取
CompactionLSM 压缩,合并 Sorted Run,产生 Changelog,支持异步/Full Compaction

Flink SQL 建表/写入/查询示例:

-- 创建 Catalog
CREATE CATALOG paimon_catalog WITH (
  'type' = 'paimon',
  'warehouse' = 'hdfs:///lake/paimon'
);
USE CATALOG paimon_catalog;

-- 创建数据库
CREATE DATABASE IF NOT EXISTS ods;

-- 1. 创建主键表(支持 CDC Upsert)
CREATE TABLE IF NOT EXISTS ods.orders (
  id BIGINT,
  order_no STRING,
  user_id BIGINT,
  item_id BIGINT,
  amount DECIMAL(10,2),
  order_status INT,
  create_time TIMESTAMP(3),
  update_time TIMESTAMP(3),
  dt STRING,
  PRIMARY KEY (id, dt) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
  'bucket' = '4',
  'changelog-producer' = 'lookup',   -- 产出完整 Changelog
  'snapshot.time-retained' = '24 h',
  'compaction.min.file-num' = '5',
  'compaction.max.file-num' = '50'
);

-- 2. 从 Kafka 流式写入(CDC Changelog 自动传播)
INSERT INTO ods.orders
SELECT id, order_no, user_id, item_id, amount, order_status,
       create_time, update_time, DATE_FORMAT(create_time, 'yyyy-MM-dd') AS dt
FROM kafka_orders;

-- 3. 从 MySQL CDC 整库同步写入
INSERT INTO paimon_catalog.ods.orders
SELECT * FROM mysql_cdc_orders;

-- 4. 流式读取(消费 Changelog)
SELECT * FROM paimon_catalog.ods.orders /*+ OPTIONS('scan.mode'='latest') */;

-- 5. 批量查询(Time Travel)
SELECT * FROM paimon_catalog.ods.orders
/*+ OPTIONS('scan.snapshot-id'='10') */;

SELECT * FROM paimon_catalog.ods.orders
FOR SYSTEM_TIME AS OF TIMESTAMP '2024-01-15 10:00:00';

-- 6. 聚合写入(批流一体)
CREATE TABLE dws.item_gmv (
  item_id BIGINT,
  dt STRING,
  gmv DECIMAL(18,2),
  order_cnt BIGINT,
  PRIMARY KEY (item_id, dt) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
  'bucket' = '2',
  'changelog-producer' = 'full-compaction',
  'full-compaction.delta-commits' = '10'
);

INSERT INTO dws.item_gmv
SELECT item_id, dt, SUM(amount) AS gmv, COUNT(*) AS order_cnt
FROM ods.orders
WHERE order_status = 1
GROUP BY item_id, dt;

-- 7. Lookup Join(Paimon 高性能点查)
SELECT o.order_id, i.item_name, o.amount
FROM kafka_orders AS o
LEFT JOIN paimon_catalog.dim.item FOR SYSTEM_TIME AS OF o.proc_time AS i
  ON o.item_id = i.id;

Paimon vs Iceberg vs Hudi 对比:

维度PaimonIcebergHudi
定位流式数据湖通用表格式流式数据湖
主键更新✅ LSM 原生,毫秒级⚠️ 需 Merge-on-Read,延迟较高✅ MoR/CoW 支持
Changelog 产出✅ 原生支持(lookup/full-compaction)❌ 不支持⚠️ 部分支持(增量视图)
Lookup Join✅ 高性能主键索引❌ 不支持⚠️ 有限支持
CDC 入湖✅ 原生 CDC Pipeline⚠️ 需额外处理✅ 支持 Delta Streamer
流批一体✅ 深度优化⚠️ 批为主,流支持弱⚠️ 流支持较好
Schema 演进✅✅ 最强✅
Flink 集成✅ 官方子项目,最深✅ Connector 完善✅ Connector 支持
引擎支持Flink(最强)、Spark/TrinoSpark/Flink/Trino/Hive 全面Spark/Flink/Hive/Trino
写入吞吐高(LSM 批量写)高(CoW)/中(MoR)中(索引开销)

典型架构:Kafka → Flink → Paimon → StarRocks/Doris

┌──────────┐     ┌──────────┐     ┌──────────────┐     ┌──────────────┐
│  MySQL   │     │  Kafka   │     │    Paimon    │     │  StarRocks/  │
│  业务库   │────▶│  日志/  │────▶│  (湖存储)    │────▶│  Doris       │
│  (CDC)   │     │  埋点    │     │ ODS/DWD/DWS  │     │  (OLAP查询)   │
└──────────┘     └──────────┘     └──────┬───────┘     └──────────────┘
      │                     ▲            │                    ▲
      │                     │            │ 批读/流读           │
      └─────────────────────┘            ▼                    │
   Flink CDC整库同步           ┌──────────────────┐           │
                              │  Flink SQL       │───────────┘
                              │  ETL/聚合/Join   │  实时导入
                              └──────────────────┘
  • ODS 层:Flink CDC → Paimon 主键表,保留原始 Changelog
  • DWD 层:Flink SQL ETL 清洗、维表 Lookup Join → Paimon
  • DWS 层:Flink SQL 聚合 → Paimon 主键表
  • ADS 层:Paimon → StarRocks/Doris(批量/实时导入),或直接查询 Paimon

Compaction 与 Expire Snapshot 配置:

-- 建表时配置 Compaction
CREATE TABLE dws.orders (
  ...
) WITH (
  -- 异步 Compaction(写入时自动触发)
  'compaction.min.file-num' = '5',
  'compaction.max.file-num' = '50',
  'compaction.target-file-size' = '128mb',

  -- Full Compaction(定期全量压缩,产出完整 Changelog)
  'changelog-producer' = 'full-compaction',
  'full-compaction.delta-commits' = '10',  -- 每10次提交触发一次

  -- Snapshot 过期(控制存储大小)
  'snapshot.time-retained' = '24 h',       -- 保留最近24小时快照
  'snapshot.num-retained.max' = '100',     -- 最多保留100个快照

  -- 分区过期(自动删除旧分区)
  'partition.expiration-time' = '7 d',
  'partition.expiration-check-interval' = '1 h',
  'partition.timestamp-formatter' = 'yyyy-MM-dd',
  'partition.timestamp-pattern' = '$dt'
);

-- 手动触发 Compaction
CALL sys.compact('paimon_catalog.dws.orders');

-- 手动触发 Full Compaction
CALL sys.compact('paimon_catalog.dws.orders', 'pt=2024-01-15');

-- 手动过期 Snapshot
CALL sys.expire_snapshots('paimon_catalog.dws.orders', '2024-01-14 00:00:00');

-- 删除分区
ALTER TABLE dws.orders DROP PARTITION (dt = '2024-01-01');

14. Flink 在实时数仓中的典型架构

14.1 经典分层架构

┌─────────────────────────────────────────────────────────────┐
│                        数据源                                │
│   MySQL/PostgreSQL  │  Kafka  │  日志  │  API  │  IoT       │
└──────────┬──────────────────────────────────────────────────┘
           │ CDC / Flume / Logstash
           ▼
┌──────────────────────────────────────────────────────────────┐
│  ODS 层(原始数据层)  ── Kafka / Paimon                      │
│  保留原始数据,不做修改                                        │
└──────────┬───────────────────────────────────────────────────┘
           │ Flink ETL(清洗、脱敏、维度关联)
           ▼
┌──────────────────────────────────────────────────────────────┐
│  DWD 层(明细数据层)  ── Kafka / Paimon                      │
│  标准化明细事实数据,统一维度                                   │
└──────────┬───────────────────────────────────────────────────┘
           │ Flink 聚合
           ▼
┌──────────────────────────────────────────────────────────────┐
│  DWS 层(汇总数据层)  ── Paimon / Doris                      │
│  轻度聚合(按天/小时/分钟),公共汇总                           │
└──────────┬───────────────────────────────────────────────────┘
           │ 导入/查询
           ▼
┌──────────────────────────────────────────────────────────────┐
│  ADS 层(应用数据层)  ── Doris / StarRocks / ClickHouse      │
│  面向业务的指标、报表、推荐特征                                 │
└──────────┬───────────────────────────────────────────────────┘
           │
           ▼
┌──────────────────────────────────────────────────────────────┐
│  应用层:BI 报表 / 实时大屏 / 风控 / 推荐 / 告警               │
└──────────────────────────────────────────────────────────────┘

14.2 Lambda vs Kappa

维度Lambda 架构Kappa 架构
核心思想批处理层(全量)+ 速度层(实时)+ 服务层一切皆流,用流处理统一批和实时
技术栈Hadoop/Spark(批)+ Storm/Flink(流)Flink(流批一体)
代码维护两套代码(批+流),需保持逻辑一致一套代码
数据一致性批和流结果可能不一致天然一致
历史数据重算批处理层重跑消费 Kafka 历史数据重跑
适用场景早期大数据平台Flink 时代的推荐架构

14.3 流批一体实践

-- 同一 SQL,既可流执行也可批执行
SET 'execution.runtime-mode' = 'streaming';  -- 流式
-- SET 'execution.runtime-mode' = 'batch';   -- 批式

INSERT INTO dws_order_summary
SELECT
  item_id,
  DATE(order_time) AS dt,
  SUM(amount) AS gmv,
  COUNT(*) AS cnt
FROM dwd_orders
GROUP BY item_id, DATE(order_time);
  • Paimon 作为流批一体存储层:流写批读、批写流读
  • Flink SQL 统一计算层:同一 SQL 切换运行模式
  • 历史数据回刷用批模式,实时增量用流模式

14.4 端到端实时数仓实战项目

项目背景

以电商实时数仓为场景,构建从业务库 CDC 到实时大屏的端到端数据管道,覆盖订单、用户、商品三大核心域,实现秒级延迟的实时指标。

数据源
数据源类型表/Topic说明
MySQL业务库 CDCorders / order_detail / user / item / category订单/用户/商品,Binlog 实时采集
Kafka日志/埋点page_view / click_event / cart_event用户行为日志,JSON 格式
Kafka业务事件pay_success / refund支付成功/退款事件
技术栈
Flink CDC 3.x + Kafka 3.x + Flink 1.18 SQL + Paimon 0.8 + Doris 2.0 + Redis
分层架构
┌─────────────────────────────────────────────────────────────────────┐
│                           数据源层                                   │
│   MySQL(订单/用户/商品)          Kafka(埋点/支付/退款日志)           │
└───────────┬──────────────────────────────┬──────────────────────────┘
            │ Flink CDC                    │ Flink Kafka Source
            ▼                              ▼
┌─────────────────────────────────────────────────────────────────────┐
│  ODS 层(贴源层)  ── Paimon 主键表 / Kafka Topic                    │
│  ods_orders / ods_order_detail / ods_user / ods_item                │
│  ods_page_view / ods_pay_success / ods_refund                       │
│  保留原始 CDC Changelog(+I/-U/+U/-D)                               │
└───────────┬─────────────────────────────────────────────────────────┘
            │ Flink SQL ETL(清洗/脱敏/维度关联/时间统一)
            ▼
┌─────────────────────────────────────────────────────────────────────┐
│  DWD 层(明细层)  ── Paimon 主键表 / Kafka Topic                    │
│  dwd_order_detail(订单宽表,关联用户/商品维度)                      │
│  dwd_page_view(标准化行为日志)                                      │
│  dwd_pay_success(支付明细,Interval Join 订单关联)                  │
│  dwd_refund(退款明细)                                              │
└───────────┬─────────────────────────────────────────────────────────┘
            │ Flink SQL 实时聚合(TUMBLE/HOP/CUMULATE 窗口)
            ▼
┌─────────────────────────────────────────────────────────────────────┐
│  DWS 层(汇总层)  ── Paimon / Doris                                 │
│  dws_item_gmv_min(商品分钟级GMV)                                    │
│  dws_user_order_hour(用户小时级订单汇总)                            │
│  dws_category_gmv_day(品类天级GMV)                                  │
│  dws_trade_overview(交易总览:GMV/订单量/客单价/UV)                  │
└───────────┬─────────────────────────────────────────────────────────┘
            │ 实时写入 / 批量导入
            ▼
┌─────────────────────────────────────────────────────────────────────┐
│  ADS 层(应用层)  ── Doris / Redis                                  │
│  ads_realtime_dashboard(大屏查询:实时GMV/TopN/趋势)                │
│  ads_user_profile(用户画像标签 → Redis)                             │
│  ads_item_recommend(商品推荐特征 → Redis)                           │
└───────────┬─────────────────────────────────────────────────────────┘
            │
            ▼
┌─────────────────────────────────────────────────────────────────────┐
│  应用:实时大屏 / BI报表 / 风控告警 / 推荐系统                        │
└─────────────────────────────────────────────────────────────────────┘
核心代码

1. ODS:Flink CDC 整库同步 MySQL → Kafka(Paimon)

-- ODS: MySQL CDC Source(订单表)
CREATE TABLE ods_orders_source (
  id BIGINT,
  order_no STRING,
  user_id BIGINT,
  item_id BIGINT,
  category_id BIGINT,
  amount DECIMAL(10,2),
  order_status INT,
  create_time TIMESTAMP(3),
  update_time TIMESTAMP(3),
  PRIMARY KEY (id) NOT ENFORCED
) WITH (
  'connector' = 'mysql-cdc',
  'hostname' = 'mysql',
  'port' = '3306',
  'username' = 'root',
  'password' = 'xxx',
  'database-name' = 'shop',
  'table-name' = 'orders',
  'server-time-zone' = 'Asia/Shanghai',
  'scan.incremental.snapshot.enabled' = 'true',
  'scan.incremental.snapshot.chunk.size' = '8096'
);

-- ODS: Paimon Sink
CREATE TABLE ods_orders (
  id BIGINT,
  order_no STRING,
  user_id BIGINT,
  item_id BIGINT,
  category_id BIGINT,
  amount DECIMAL(10,2),
  order_status INT,
  create_time TIMESTAMP(3),
  update_time TIMESTAMP(3),
  dt STRING,
  PRIMARY KEY (id, dt) NOT ENFORCED
) PARTITIONED BY (dt) WITH (
  'connector' = 'paimon',
  'warehouse' = 'hdfs:///lake/paimon',
  'database' = 'ods',
  'table' = 'orders',
  'bucket' = '4',
  'changelog-producer' = 'input',
  'snapshot.time-retained' = '24 h'
);

-- 整库同步(CDC Changelog 直接传播到 Paimon)
INSERT INTO ods_orders
SELECT id, order_no, user_id, item_id, category_id, amount,
       order_status, create_time, update_time,
       DATE_FORMAT(create_time, 'yyyy-MM-dd') AS dt
FROM ods_orders_source;

2. DWD:Flink SQL ETL 清洗 + 维表 Lookup Join

-- DWD: 订单明细宽表(关联用户/商品/品类维度)
CREATE TABLE dwd_order_detail (
  order_id BIGINT,
  order_no STRING,
  user_id BIGINT,
  user_name STRING,
  user_level INT,
  item_id BIGINT,
  item_name STRING,
  category_id BIGINT,
  category_name STRING,
  amount DECIMAL(10,2),
  order_status INT,
  order_status_name STRING,
  create_time TIMESTAMP(3),
  WATERMARK FOR create_time AS create_time - INTERVAL '5' SECOND,
  PRIMARY KEY (order_id) NOT ENFORCED
) WITH (
  'connector' = 'paimon',
  'warehouse' = 'hdfs:///lake/paimon',
  'database' = 'dwd',
  'table' = 'order_detail',
  'bucket' = '4',
  'changelog-producer' = 'lookup'
);

INSERT INTO dwd_order_detail
SELECT
  o.id AS order_id,
  o.order_no,
  o.user_id,
  u.user_name,
  u.user_level,
  o.item_id,
  i.item_name,
  o.category_id,
  c.category_name,
  o.amount,
  o.order_status,
  CASE o.order_status
    WHEN 1 THEN '待支付'
    WHEN 2 THEN '已支付'
    WHEN 3 THEN '已发货'
    WHEN 4 THEN '已完成'
    WHEN 5 THEN '已取消'
    WHEN 6 THEN '已退款'
  END AS order_status_name,
  o.create_time
FROM ods_orders AS o
LEFT JOIN dim_user FOR SYSTEM_TIME AS OF o.proc_time AS u
  ON o.user_id = u.id
LEFT JOIN dim_item FOR SYSTEM_TIME AS OF o.proc_time AS i
  ON o.item_id = i.id
LEFT JOIN dim_category FOR SYSTEM_TIME AS OF o.proc_time AS c
  ON o.category_id = c.id;

3. DWS:实时聚合写入 Paimon/Doris

-- DWS: 商品分钟级 GMV(滑动/滚动窗口聚合)
CREATE TABLE dws_item_gmv_min (
  item_id BIGINT,
  item_name STRING,
  window_start TIMESTAMP(3),
  window_end TIMESTAMP(3),
  gmv DECIMAL(18,2),
  order_cnt BIGINT,
  user_cnt BIGINT,
  PRIMARY KEY (item_id, window_start) NOT ENFORCED
) WITH (
  'connector' = 'doris',
  'fenodes' = 'doris-fe:8030',
  'table.identifier' = 'dws.item_gmv_min',
  'username' = 'root',
  'password' = 'xxx',
  'sink.label-prefix' = 'dws-gmv',
  'sink.enable-2pc' = 'true',
  'sink.properties.format' = 'json'
);

INSERT INTO dws_item_gmv_min
SELECT
  item_id,
  item_name,
  window_start,
  window_end,
  SUM(amount) AS gmv,
  COUNT(*) AS order_cnt,
  COUNT(DISTINCT user_id) AS user_cnt
FROM TABLE(
  TUMBLE(TABLE dwd_order_detail, DESCRIPTOR(create_time), INTERVAL '1' MINUTE)
)
WHERE order_status IN (2, 3, 4)
GROUP BY item_id, item_name, window_start, window_end;

-- DWS: 交易总览(累计指标)
CREATE TABLE dws_trade_overview (
  stat_date DATE,
  total_gmv DECIMAL(18,2),
  total_orders BIGINT,
  total_users BIGINT,
  avg_order_amount DECIMAL(10,2),
  PRIMARY KEY (stat_date) NOT ENFORCED
) WITH (
  'connector' = 'doris',
  'fenodes' = 'doris-fe:8030',
  'table.identifier' = 'dws.trade_overview',
  'username' = 'root',
  'password' = 'xxx',
  'sink.enable-2pc' = 'true'
);

INSERT INTO dws_trade_overview
SELECT
  DATE(create_time) AS stat_date,
  SUM(amount) AS total_gmv,
  COUNT(*) AS total_orders,
  COUNT(DISTINCT user_id) AS total_users,
  AVG(amount) AS avg_order_amount
FROM dwd_order_detail
WHERE order_status IN (2, 3, 4)
GROUP BY DATE(create_time);

4. ADS:结果输出到 Redis 供大屏查询

-- ADS: 实时大屏 TopN 商品 → Redis
CREATE TABLE ads_topn_items (
  rank_key STRING,
  item_info STRING,
  PRIMARY KEY (rank_key) NOT ENFORCED
) WITH (
  'connector' = 'redis',
  'mode' = 'single',
  'host' = 'redis',
  'port' = '6379',
  'command' = 'SET',
  'redis.value.data-type' = 'string'
);

INSERT INTO ads_topn_items
SELECT
  CONCAT('topn:item:', CAST(window_start AS STRING)) AS rank_key,
  item_name AS item_info
FROM (
  SELECT *,
    ROW_NUMBER() OVER (
      PARTITION BY window_start
      ORDER BY gmv DESC
    ) AS rn
  FROM dws_item_gmv_min
) WHERE rn <= 10;
关键设计

数据流向与 CDC Changelog 传播:

MySQL Binlog (+I/-U/+U/-D)
    │ Flink CDC 解析为 RowData
    ▼
Paimon ODS 主键表(保留完整 Changelog)
    │ 流式读取(changelog-producer=input 直接传播)
    ▼
DWD 宽表(Lookup Join 补全维度,CDC 语义保持)
    │ 聚合(+I 产生新结果,-U/+U 更新之前结果)
    ▼
DWS/ADS(Upsert Sink 幂等写入)

关键:Flink SQL 原生支持 Changelog 传播,CDC 的 UPDATE/DELETE 会在整个管道中正确传递,最终通过 Upsert Sink 写入外部系统。

状态 TTL 与 Idle State Retention:

-- 生产环境必须配置,防止状态无限增长
SET 'table.exec.state.ttl' = '24 h';
SET 'table.exec.idle-state-retention' = '1 h';
  • Regular Join 的两侧状态在 TTL 后自动清理
  • 空闲超过 idle-state-retention 的 key 状态会被清理
  • 窗口状态在窗口结束 + Allowed Lateness 后自动清理

多流 Join:

-- Interval Join:订单流 + 支付流(30分钟内支付)
SELECT o.order_id, o.amount, p.pay_time, p.pay_channel
FROM dwd_order_detail AS o, dwd_pay_success AS p
WHERE o.order_id = p.order_id
  AND p.pay_time BETWEEN o.create_time AND o.create_time + INTERVAL '30' MINUTE;

-- Temporal Join:关联历史版本汇率
SELECT o.order_id, o.amount * r.rate AS usd_amount
FROM dwd_order_detail AS o
JOIN dim_exchange_rate FOR SYSTEM_TIME AS OF o.create_time AS r
  ON o.currency = r.currency;

幂等写入与 Exactly-Once:

层Sink一致性机制
ODSPaimon原子 Snapshot 提交(2PC)
DWDPaimon原子 Snapshot 提交(2PC)
DWSDorisStream Load 2PC
ADSRedis幂等 SET 覆盖写(At-Least-Once,结果一致)

时区统一:

-- 全局配置
SET 'table.local-time-zone' = 'Asia/Shanghai';

-- MySQL CDC 必须指定时区
'server-time-zone' = 'Asia/Shanghai'

-- Kafka 时间戳统一转换
CREATE TABLE kafka_source (
  event_time BIGINT,
  ts AS TO_TIMESTAMP_LTZ(event_time, 3),
  WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
) WITH (...);

数据质量监控:

  • Flink DataStream 中添加校验算子(Null 值/异常范围/枚举校验),异常数据发侧输出
  • 接入 Great Expectations 或 Apache Griffin 做离线数据质量校验
  • 监控核心指标:空值率、主键重复率、数值异常率

离线补数与版本切换:

-- 批模式回刷历史数据
SET 'execution.runtime-mode' = 'batch';

INSERT INTO dws_item_gmv_min
SELECT ... FROM ods_orders WHERE dt BETWEEN '2024-01-01' AND '2024-01-31';

-- 使用 Savepoint 升级作业(状态兼容)
./bin/flink stop --savepointPath hdfs:///flink/savepoints <jobId>
./bin/flink run -s hdfs:///flink/savepoints/savepoint-xxx ./new-job.jar

监控告警:

  • Checkpoint:成功率、时长、大小(参见 11.9 节)
  • Kafka Lag:各 Source 消费积压
  • 反压:backPressuredTimeMsPerSecond 告警
  • 数据延迟:System.time() - currentOutputWatermark
  • 资源:CPU/内存/GC

资源隔离与多租户:

  • K8s Namespace 隔离不同业务线
  • Flink K8s Operator 的 FlinkDeployment 按团队独立部署
  • Paimon Database 级权限控制
  • Kafka Topic 前缀隔离 + ACL
  • Doris Database/Table 级 RBAC

15. 学习路线与实战建议

15.1 学习阶段划分

阶段一(1-2周):基础入门
├── Flink 概念与架构
├── Local/Standalone 环境搭建
├── DataStream API 基础
└── 简单 WordCount / Kafka 消费

阶段二(2-4周):核心概念
├── 时间语义与 Watermark
├── Window 机制
├── 状态管理
├── Checkpoint 与容错
├── TaskManager 内存模型(Framework/Task/Network/Managed/JVM)
└── Flink SQL 基础

阶段三(4-8周):进阶实战
├── Flink SQL 高级(Join/Top-N/CDC/性能调优)
├── 异步 I/O(AsyncDataStream 维表关联)
├── Flink CDC 整库同步
├── CEP 复杂事件处理
├── 性能调优与问题排查(反压/数据倾斜/Checkpoint)
├── Metrics 与 Prometheus + Grafana 监控
└── 实时数仓项目实战

阶段四(持续):生产深耕
├── Flink on K8s(Native/Operator)
├── Paimon 流式数据湖深度集成
├── Watermark Alignment 与高级调优
├── 源码阅读(JobManager/TaskManager/Checkpoint/网络栈)
└── Flink 2.0 新技术(Disaggregated State/Autoscaling)

15.2 推荐资源

书籍:

  • 《Flink 基础教程》(Stream Processing with Apache Flink)
  • 《Flink 内核原理与实现》
  • 《Flink SQL 实战》

官方资源:

  • Flink 官方文档:https://flink.apache.org/docs/
  • Flink CDC:https://nightlies.apache.org/flink/flink-cdc-docs-stable/
  • Paimon:https://paimon.apache.org/
  • Flink 中文社区:https://flink-learning.org.cn/

视频:

  • 尚硅谷 Flink 教程(入门推荐)
  • Apache Flink 官方 YouTube 频道
  • Flink Forward 大会演讲

源码:

  • GitHub:https://github.com/apache/flink
  • 从 StreamExecutionEnvironment.execute() 入口阅读

15.3 项目练习建议

  1. WordCount(入门必做):Socket + 文件 + Kafka
  2. 实时热门商品:Kafka 订单流 → 滑动窗口 TopN → Redis/MySQL
  3. 实时风控告警:CEP 检测可疑交易模式
  4. 实时数仓搭建:MySQL CDC → Kafka → Flink SQL → Doris
  5. 整库同步到湖仓:Flink CDC → Paimon → Trino/Presto 查询

15.4 面试高频考点

主题高频问题
架构JobManager/TaskManager 作用?Slot 与并行度关系?Flink on K8s 有哪些部署模式?
时间Watermark 原理?乱序数据如何处理?Watermark Alignment 解决什么问题?
窗口窗口触发时机?Allowed Lateness 机制?
状态Keyed State 与 Operator State 区别?State Backend 选择?
内存模型TaskManager 内存分哪几个区域?Managed Memory 做什么用?容器 OOM Killed 怎么排查?
容错Checkpoint 流程?Barrier 对齐 vs 非对齐?2PC 如何保证端到端 Exactly-Once?
反压Credit-based 反压机制原理?如何通过 Web UI/Metrics 定位反压?
异步 I/OAsync I/O 原理?orderedWait 和 unorderedWait 区别?生产注意什么?
SQL动态表概念?各种 Join 区别?Top-N 实现?数据倾斜三板斧(MiniBatch/Local-Global/Split Distinct)?
SQL 调优MiniBatch 原理?Local-Global 两阶段聚合?COUNT(DISTINCT) 怎么优化?维表 Join 如何加速?
CDC全量+增量如何无缝切换?无锁原理?Flink CDC 3.x Pipeline 有什么优势?
Connector各 Sink 的 Exactly-Once 支持情况?Kafka 事务如何实现?Doris/StarRocks 如何保证 EOS?
数据湖Paimon vs Iceberg vs Hudi 对比?Paimon Changelog Producer 有哪些模式?
监控Flink Metrics 有哪些类型?关键监控指标有哪些?如何集成 Prometheus + Grafana?
调优数据倾斜怎么处理?Checkpoint 超时怎么办?Managed Memory 如何配置?
对比Flink vs Spark Streaming?Lambda vs Kappa?Flink on K8s vs YARN?

总结

Flink 作为当前最主流的流处理引擎,其知识体系涵盖API 使用、核心原理、生态集成、生产调优四大维度。学习 Flink 的关键路径是:

  1. 先跑起来:搭建环境,写通第一个 DataStream 程序
  2. 理解核心:时间、窗口、状态、Checkpoint 四大基石
  3. 掌握 SQL:Flink SQL 是生产效率的关键,也是 2.0 方向
  4. 实战驱动:在真实项目中踩坑、调优、深入
  5. 关注前沿:流批一体、湖仓一体、Disaggregated State、Flink 2.0

最后一句话: 流处理的世界里,Flink 不是工具,而是一种思维方式——用流的视角看待数据,用状态的思维构建应用,用 Checkpoint 的信念保证可靠。祝你在 Flink 的学习之路上越走越远!


本文基于 Apache Flink 1.18/1.19/1.20 版本撰写,部分内容参考 Flink 官方文档。Flink 2.0 特性为社区路线图预览,具体以官方发布为准。


更多推荐