前言:把自己几年前整理的 Flink 笔记翻出来看了一遍,发现好多东西已经过时了——DataSet API 早就被废弃了、旧的 Source/Sink 接口换成了新的流批一体架构、Watermark 的写法也变了。趁着这次整理,以目前生产环境主流的 Flink 1.18/1.19 为准重新梳理了一遍,去掉过时的东西,补充新特性,顺便把这些年踩过的坑也写进去。不管你是刚入坑的小白,还是想系统回顾一下的老司机,这篇文章应该都能帮到你。


一、Flink 是什么?为什么选它?

Flink 就是一个专门做实时流计算的引擎。和 Spark Streaming 那种"攒一小批再处理"的微批模式不同,Flink 是数据来一条处理一条,延迟能压到毫秒级。

它几个突出的优势:

  • 真正的流处理:逐条处理,不是微批
  • exactly-once 精准一次:配合 Checkpoint 机制,数据既不丢也不重,金融、电商订单场景这是刚需
  • 流批一体:同一套 DataStream API 既能处理实时流也能处理批数据(DataSet API 1.12+ 已废弃)
  • 状态管理:自带 State 状态后端,窗口聚合、连接操作时不用自己操心数据存哪
  • 时间语义丰富:EventTime + Watermark 这套机制是解决乱序数据的核心手段

如果你的业务是"数据源源不断涌进来,我要实时算出结果"——实时监控、实时推荐、实时风控这些场景,Flink 基本是首选。


二、Flink 执行架构与提交流程

2.1 三大核心角色

Client(提交端)
  └── 把代码编译成 JobGraph,提交给 JobManager

JobManager(作业管理器,JM)
  └── 大脑,负责任务调度、协调 Checkpoint、故障恢复
  └── 1.17+ 支持多 JM Standby 高可用(基于 ZooKeeper 或 K8s)

TaskManager(任务管理器,TM)
  └── 干活的,里面有很多 Slot(插槽),每个 Slot 跑一个任务线程
  └── Slot 数量一般配成和 CPU 核心数一致

2.2 提交流程

1. Client 编译代码 → StreamGraph → JobGraph
2. Client 把 JobGraph 提交给 JobManager
3. JobManager 把 JobGraph 转成 ExecutionGraph,调度到各个 TM 的 Slot 上执行
4. TM 定期向 JM 汇报心跳和状态

Flink执行架构与提交流程

2.3 部署模式怎么选

部署模式 适用场景 说明
Standalone 测试/小集群 独立集群,自己管资源,生产不推荐
YARN 传统大数据集群 Session 模式(共享JM)、Per-Job 模式(独立资源)、Application 模式(推荐)
Kubernetes 云原生/生产推荐 Flink 1.18+ 对 K8s 支持很好,配合 Helm Chart 一键部署

生产踩坑:TM 的 Slot 数量默认等于 CPU 核心数,并行度超过总 Slot 数作业会卡住起不来!slot 总数 = TM 个数 × 每个 TM 的 slot 数,提交前先用 flink info your.jar 检查一下。


三、DataStream API:流处理的核心武器

3.1 执行环境(1.18+ 写法)

老写法(已废弃,别用了)

ExecutionEnvironment env = ExecutionEnvironment.getExecutionEnvironment();  // DataSet API,废弃了

新写法(统一走 DataStream)

// 流模式
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
env.setRuntimeMode(RuntimeExecutionMode.STREAMING);

// 批模式(处理有界数据,比如历史数据回补)
env.setRuntimeMode(RuntimeExecutionMode.BATCH);

// 自动模式(根据数据源自动判断,有界走批优化,无界走流)
env.setRuntimeMode(RuntimeExecutionMode.AUTOMATIC);

3.2 Source:数据来源

旧接口(addSource,废弃)

env.addSource(new FlinkKafkaConsumer<>("topic", new SimpleStringSchema(), props));

新接口(fromSource,1.12+ 推荐)

KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers("kafka:9092")
    .setTopics("input-topic")
    .setGroupId("flink-consumer")
    .setStartingOffsets(OffsetsInitializer.earliest())
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .build();

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

常用 Source:

Source 用途
fromElements() / fromCollection() 本地测试
fromSource(FileSource) 读文件,支持监控目录新文件
fromSource(KafkaSource) 生产最常用
socketTextStream() 快速测试,nc -lk 9999 配合

3.3 Transformation:数据转换

// 1. map:一对一转换
DataStream<String> names = stream.map(Event::getName);

// 2. flatMap:一对多展开
DataStream<SubEvent> subEvents = stream.flatMap(new FlatMapFunction<Event, SubEvent>() {
    public void flatMap(Event value, Collector<SubEvent> out) {
        for (SubEvent sub : value.getSubEvents()) {
            out.collect(sub);
        }
    }
});

// 3. filter:过滤
DataStream<Event> filtered = stream.filter(e -> e.getAmount() > 100);

// 4. keyBy:按 key 分组(分组后才能聚合!)
KeyedStream<Event, String> keyed = stream.keyBy(Event::getUserId);

// 5. 聚合(必须在 keyBy 之后)
keyed.sum("amount");
keyed.min("amount");
keyed.maxBy("timestamp");
keyed.reduce((e1, e2) -> { e1.setAmount(e1.getAmount() + e2.getAmount()); return e1; });

// 6. process:最灵活的函数,可以访问定时器和状态
DataStream<Result> result = keyed.process(new MyProcessFunction());

踩坑keyBy 后同一 key 的数据一定在同一个子任务上。如果某个 key 流量占了 80%,就会造成数据倾斜——某个子任务忙死,其他闲死。解法:加随机前缀打散 key,或者用自定义分区器。

3.4 Sink:数据输出

旧接口(废弃)

stream.addSink(new FlinkKafkaProducer<>("topic", new SimpleStringSchema(), props));

新接口(sinkTo,1.12+)

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

stream.sinkTo(sink);

Kafka exactly-once 要注意:transactional.id.prefix 全局唯一,Kafka 的 transaction.max.timeout.ms 要大于 Flink Checkpoint 间隔,否则事务超时。

3.5 物理分区

stream.shuffle();          // 随机分发
stream.rebalance();        // 轮询分发,解决数据倾斜
stream.rescale();          // 本地轮询,比rebalance高效
stream.broadcast();        // 广播到所有子任务(小表join用)
stream.global();           // 全部发第一个子任务(慎用,并行度变1)
stream.partitionCustom(...); // 自定义分区

3.6 Side Output:侧输出流

一条数据可能需要走不同分支处理,比如正常数据走主流程,异常数据走旁路:

OutputTag<String> errorTag = new OutputTag<String>("error"){};

SingleOutputStreamOperator<String> mainStream = stream
    .process(new ProcessFunction<String, String>() {
        @Override
        public void processElement(String value, Context ctx, Collector<String> out) {
            if (value.contains("ERROR")) {
                ctx.output(errorTag, value);  // 异常数据发到侧输出流
            } else {
                out.collect(value);           // 正常数据继续主流程
            }
        }
    });

// 获取侧输出流
DataStream<String> errorStream = mainStream.getSideOutput(errorTag);
errorStream.addSink(new ErrorSink());

3.7 BroadcastState:广播状态(小表 Join 必备)

实时计算里经常遇到这种情况:主流是海量日志,需要关联一张配置表(商品类目、用户标签),配置表不大但会更新。

BroadcastStream<Config> broadcastConfigStream = configStream
    .broadcast(new MapStateDescriptor<>("config", Types.STRING, Types.STRING));

mainStream
    .connect(broadcastConfigStream)
    .process(new BroadcastProcessFunction<Event, Config, Result>() {
        @Override
        public void processElement(Event event, ReadOnlyContext ctx, Collector<Result> out) {
            String category = ctx.getBroadcastState(configDesc).get(event.getItemId());
            out.collect(new Result(event, category));
        }
        
        @Override
        public void processBroadcastElement(Config config, Context ctx, Collector<Result> out) {
            ctx.getBroadcastState(configDesc).put(config.getKey(), config.getValue());
        }
    });

注意:processElement 里只能广播状态,processBroadcastElement 里才能读写。这是 Flink 的设计,保证所有子任务广播状态一致。


四、Flink 四大基石以及进阶扩展

4.1~4.4这四个是 Flink 四大基石部分,面试重点,生产也天天用。

4.1 Window 窗口计算

流数据无限长,做聚合(统计 UV、销售额)需要切分成有限块。

窗口类型 代码 特点
滚动窗口 (Tumbling) .window(TumblingEventTimeWindows.of(Time.seconds(5))) 窗口不重叠,像切香肠
滑动窗口 (Sliding) .window(SlidingEventTimeWindows.of(Time.seconds(10), Time.seconds(5))) 窗口重叠,适合移动平均
会话窗口 (Session) .window(EventTimeSessionWindows.withGap(Time.seconds(30))) 按活动间隔切分,适合用户行为分析
计数窗口 (Count) .countWindow(100) 每 N 条触发一次,生产少用

Flink Window窗口类型全解析

窗口 API 用法:

stream
    .keyBy(Event::getUserId)
    .window(TumblingEventTimeWindows.of(Time.minutes(5)))
    .aggregate(new MyAggregateFunction())    // 增量聚合,性能好
    // 或者 .process(new MyWindowFunction()) // 全量聚合,灵活但耗内存

踩坑

  1. EventTime 窗口必须设置 Watermark,否则窗口永远不触发。
  2. timeWindow() 1.18 已废弃,改用 window(TumblingEventTimeWindows.of(...))
  3. 优先用 AggregateFunction(增量计算,内存友好),WindowFunction 是全量加载到内存再算,数据量大容易 OOM。

4.2 EventTime 与 Watermark

Flink 最精妙的设计之一。核心问题:网络延迟导致数据乱序到达,怎么保证算得对?

三种时间语义

时间类型 定义 适用场景
EventTime 数据真实产生的时间(用户点击的时间戳) 生产首选,能处理乱序
IngestionTime 数据进入 Flink Source 的时间 极少用
ProcessingTime 算子执行的系统时间 对延迟不敏感的场景

Watermark 原理

Watermark 是一个时间戳标记,告诉系统"所有 EventTime 小于 Watermark 的数据应该都到了"。

Watermark = 当前最大 EventTime - 允许的最大延迟

当 Watermark >= 窗口结束时间,窗口才触发计算。

Flink Watermark水位线机制详解

生产代码:

WatermarkStrategy<Event> strategy = WatermarkStrategy
    .<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
    .withTimestampAssigner((event, timestamp) -> event.getTimestamp())
    .withIdleness(Duration.ofMinutes(1));  // 多并行度必加!

dataStream
    .assignTimestampsAndWatermarks(strategy)
    .keyBy(Event::getUserId)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .aggregate(...);

特别注意:并行度 > 1 时,全局 Watermark 取所有并行 Source 的最小值!某个 Source 分区一直没数据,整个作业的 Watermark 就不推进,窗口永远不算!withIdleness() 多并行度场景必加

迟到数据处理

stream
    .assignTimestampsAndWatermarks(strategy)
    .keyBy(Event::getUserId)
    .window(TumblingEventTimeWindows.of(Time.minutes(1)))
    .allowedLateness(Time.seconds(30))    // 允许额外迟到30秒
    .sideOutputLateData(lateOutputTag)    // 超时的扔到侧输出流
    .aggregate(...);

// 侧输出流处理迟到数据
DataStream<Event> lateData = result.getSideOutput(lateOutputTag);
lateData.addSink(new LateDataSink());

4.3 Checkpoint:分布式快照与容错

Checkpoint 是 Flink 可靠性的核心。原理是在数据流中插入 Barrier(栅栏),所有算子对齐 Barrier 后把状态快照存到持久化存储。作业挂了就从最近的 Checkpoint 恢复。

Flink Checkpoint分布式快照机制

生产配置(1.18+)

env.enableCheckpointing(60000, CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(600000);
env.getCheckpointConfig().setTolerableCheckpointFailureNumber(5);
env.getCheckpointConfig().setMaxConcurrentCheckpoints(1);
env.getCheckpointConfig().setExternalizedCheckpointCleanup(
    ExternalizedCheckpointCleanup.RETAIN_ON_CANCELLATION);

// 大状态用 RocksDB
EmbeddedRocksDBStateBackend rocksDbBackend = new EmbeddedRocksDBStateBackend(true);
env.setStateBackend(rocksDbBackend);
env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints");

三种 StateBackend 对比

Backend 存储位置 适用场景
MemoryStateBackend TM 内存 本地测试,快但不安全
HashMapStateBackend TM 内存 + 异步快照 小状态、追求低延迟
EmbeddedRocksDBStateBackend 本地 RocksDB + 增量快照 大状态生产首选,支持 TB 级

Checkpoint 调优经验

  1. 间隔不要太短(至少 30s),太频繁影响吞吐;也不要太长(超过 5min),恢复时回放太多。
  2. 反压时 Checkpoint 容易超时(Barrier 对齐要等数据消化完),1.11+ 可开非对齐 Checkpoint:enableUnalignedCheckpoints(),但会增加状态大小。
  3. RocksDB 增量 Checkpoint 1.18 默认开启。
  4. Checkpoint 一直失败,先排查状态大小和反压。

Savepoint vs Checkpoint

特性 Checkpoint Savepoint
触发 自动周期性 手动
用途 故障恢复 升级、迁移
清理 自动 手动
兼容性 同版本 跨小版本兼容
flink savepoint <jobId> hdfs:///flink/savepoints
flink run -s hdfs:///flink/savepoints/savepoint-xxxxx -c com.example.MyJob my-job.jar

4.4 State 状态管理

Flink 是有状态计算框架。State 就是算子运行过程中的中间结果。

两种状态类型

类型 说明 场景
KeyedState 和 key 绑定,每个 key 独立 keyBy 后的聚合、窗口
OperatorState 和算子实例绑定,不分 key Kafka Consumer offset

KeyedState 四种数据结构

ValueState<T>        // 单值,比如累计金额
ListState<T>         // 列表,比如最近10个IP
MapState<K, V>       // Map,比如商品库存
ReducingState<T>     // 自动归约,增量聚合

代码示例(检测连续3次登录失败报警):

public class LoginFailDetect extends KeyedProcessFunction<String, LoginEvent, Alert> {
    private ValueState<Integer> failCountState;

    @Override
    public void open(Configuration parameters) {
        failCountState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("failCount", Types.INT));
    }

    @Override
    public void processElement(LoginEvent event, Context ctx, Collector<Alert> out) {
        if (event.getResult().equals("FAIL")) {
            int count = failCountState.value() == null ? 0 : failCountState.value();
            failCountState.update(count + 1);
            if (count + 1 >= 3) {
                out.collect(new Alert(event.getUserId(), "连续3次登录失败"));
                failCountState.clear();
            }
        } else {
            failCountState.clear();
        }
    }
}

4.5 ProcessFunction:最灵活的算子

ProcessFunction 可以访问状态定时器,是 Flink 最底层的 API,能实现任何复杂逻辑。

public class CountWithTimeout extends KeyedProcessFunction<String, Event, Result> {
    private ValueState<CountState> state;

    @Override
    public void open(Configuration parameters) {
        state = getRuntimeContext().getState(
            new ValueStateDescriptor<>("myState", CountState.class));
    }

    @Override
    public void processElement(Event event, Context ctx, Collector<Result> out) {
        CountState current = state.value();
        if (current == null) {
            current = new CountState();
            current.key = event.getUserId();
            ctx.timerService().registerEventTimeTimer(event.getTimestamp() + 60000);
        }
        current.count++;
        state.update(current);
    }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<Result> out) {
        CountState result = state.value();
        out.collect(new Result(result.key, result.count));
        state.clear();
    }
}

onTimer 是 ProcessFunction 的灵魂。注册定时器后,当 Watermark 推进到设定时间,onTimer 就会被调用。可以用来做超时检测延迟触发定期清理状态等。


4.6 Async I/O:异步访问外部系统

流计算经常要查外部数据库(Redis、MySQL、HBase),同步查询会阻塞数据流,严重降低吞吐。Async I/O 让你可以并发发请求不阻塞。

AsyncDataStream.unorderedWait(
    stream,
    new AsyncFunction<Event, EnrichedEvent>() {
        @Override
        public void asyncInvoke(Event event, ResultFuture<EnrichedEvent> resultFuture) {
            redisClient.getAsync(event.getUserId())
                .thenAccept(profile -> {
                    resultFuture.complete(
                        Collections.singletonList(new EnrichedEvent(event, profile))
                    );
                });
        }
    },
    1000, TimeUnit.MILLISECONDS, 100
);

unorderedWait 乱序输出吞吐最高,orderedWait 按输入顺序输出吞吐低。业务不要求顺序就用 unorderedWait


4.7 双流 Join

实时计算经常需要把两条流关联起来。

Window Join:两条流数据落在同一个窗口内才关联

stream1.join(stream2)
    .where(Event1::getKey)
    .equalTo(Event2::getKey)
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .apply(new JoinFunction<Event1, Event2, Result>() {
        public Result join(Event1 first, Event2 second) {
            return new Result(first, second);
        }
    });

Interval Join:更灵活,stream1 的元素和 stream2 在特定时间范围内的元素关联

stream1
    .keyBy(Event1::getKey)
    .intervalJoin(stream2.keyBy(Event2::getKey))
    .between(Time.seconds(-5), Time.seconds(5))
    .process(new ProcessJoinFunction<Event1, Event2, Result>() {
        public void processElement(Event1 left, Event2 right, Context ctx, Collector<Result> out) {
            out.collect(new Result(left, right));
        }
    });

Interval Join 比 Window Join 更省资源,匹配到就输出,不用等窗口结束。

CoGroup

stream1.coGroup(stream2)
    .where(Event1::getKey)
    .equalTo(Event2::getKey)
    .window(TumblingEventTimeWindows.of(Time.seconds(10)))
    .apply(new CoGroupFunction<Event1, Event2, Result>() {
        public void coGroup(Iterable<Event1> first, Iterable<Event2> second, Collector<Result> out) {
            for (Event1 e1 : first) {
                for (Event2 e2 : second) {
                    out.collect(new Result(e1, e2));
                }
            }
        }
    });

CoGroup 比 Join 更底层,能处理单边数据为空的情况。

Temporal Table Join(时态表 Join):流表关联维表的历史版本,适合 SCD 场景。

SELECT o.order_id, o.amount, c.currency_rate
FROM orders AS o
LEFT JOIN currency_rates FOR SYSTEM_TIME AS OF o.proc_time AS c
ON o.currency = c.currency;

4.8 State TTL:状态自动过期

状态一直增长会导致内存爆炸。State TTL 让状态自动清理。

StateTtlConfig ttlConfig = StateTtlConfig
    .newBuilder(Time.hours(24))
    .setUpdateType(StateTtlConfig.UpdateType.OnCreateAndWrite)
    .setStateVisibility(StateTtlConfig.StateVisibility.NeverReturnExpired)
    .cleanupIncrementally(10, true)
    .build();

ValueStateDescriptor<MyState> descriptor = new ValueStateDescriptor<>("myState", MyState.class);
descriptor.enableTimeToLive(ttlConfig);
ValueState<MyState> state = getRuntimeContext().getState(descriptor);

RocksDB 后端下,TTL 清理是懒清理——不会立即删,读取时过滤。配合 cleanupIncrementally 可以批量清理过期数据。


4.9 反压(Backpressure)

反压是流处理系统的自我保护。下游处理不过来时,反压信号往上传播,上游减速。

原理:每个算子有输入缓冲,缓冲满了就停止读取上游。1.18+ 支持基于信用的流量控制(Credit-based),比老的反压机制更高效。

发现反压:Flink Web UI 看 Backpressure 标签,显示 OK / LOW / HIGH。或者看指标 task.BackPressuredTimeMsPerSecond

解决反压

  1. 增加并行度或资源
  2. 优化算子逻辑,减少每条数据处理时间
  3. 调整缓冲大小:taskmanager.memory.network.min/max

4.10 TypeInformation 与序列化

Flink 自己管理类型和序列化,这是高性能的关键。

// 显式声明类型
dataStream.returns(TypeInformation.of(new TypeHint<Event>() {}));

// 常用类型
Types.STRING / Types.INT / Types.LONG
Types.TUPLE(Types.STRING, Types.INT)
TypeInformation.of(Event.class)  // POJO

POJO 类要满足的条件(否则退化成 Kryo,性能差):

  1. public 类
  2. 有无参构造函数
  3. 字段 public 或者有 getter/setter
  4. 字段类型也是 Flink 支持的

性能:POJO > Tuple > Kryo。尽量用 POJO 或 Tuple,大数据量下差距很明显。


4.11 Operator Chaining 与 Slot Sharing Group

Flink 默认会把多个算子 chain 在一起执行(在同一个线程里),减少序列化和网络传输开销。但有时候不想 chain,比如某个算子特别耗资源:

stream.map(...).disableChaining();  // 禁止和这个算子 chain
stream.startNewChain();             // 从这里开始一个新 chain

Slot Sharing Group 控制哪些算子可以共享同一个 Slot:

stream.map(...).slotSharingGroup("group1");
stream.filter(...).slotSharingGroup("group2");

默认所有算子都在 default sharing group。合理划分 group 可以避免资源争抢,比如 IO 密集型和 CPU 密集型算子分开。


4.12 CEP:复杂事件处理

CEP 用来从事件流中识别复杂模式。比如检测"用户5分钟内连续登录失败3次然后成功1次"。

Pattern<LoginEvent, ?> pattern = Pattern.<LoginEvent>begin("fail")
    .where(evt -> evt.getResult().equals("FAIL"))
    .times(3)
    .within(Time.minutes(5))
    .next("success")
    .where(evt -> evt.getResult().equals("SUCCESS"));

PatternStream<LoginEvent> patternStream = CEP.pattern(
    loginStream.keyBy(LoginEvent::getUserId), pattern);

patternStream.process(new PatternProcessFunction<LoginEvent, Alert>() {
    @Override
    public void processMatch(Map<String, List<LoginEvent>> match, Context ctx, Collector<Alert> out) {
        LoginEvent success = match.get("success").get(0);
        out.collect(new Alert(success.getUserId(), "连续失败3次后登录成功,疑似撞库"));
    }
});

CEP 模式语法:begin(开始)→ next(严格紧邻)/ followedBy(非严格)/ times(重复次数)/ within(时间窗口)。适合风控、安全检测、用户行为分析。


4.13 内存模型

Flink 1.10+ 引入了新的内存模型,TM 内存被划分为几个区域:

内存区域 配置 用途
JVM Heap taskmanager.memory.framework.heap.size Flink 框架本身、用户代码中的对象
Managed Memory taskmanager.memory.managed.size RocksDB 状态后端、排序、哈希表
Network Memory taskmanager.memory.network.* 算子间数据传输的缓冲
JVM Overhead taskmanager.memory.jvm-overhead.* 其他 JVM 开销(JNI、NIO 等)

生产环境 TM 内存建议 4GB 起步,大状态作业配到 16GB 以上。RocksDB 会用 Managed Memory 做块缓存,大状态场景 Managed Memory 要留够(建议总内存的 30%~40%,具体配置多大也请参考公司集群实际大小)。


五、Table API & SQL

5.1 创建表环境

StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
StreamTableEnvironment tableEnv = StreamTableEnvironment.create(env);

// 纯 SQL 方式
TableEnvironment tableEnv = TableEnvironment.create(EnvironmentSettings.inStreamingMode());

5.2 DDL 建表 + SQL 查询

tableEnv.executeSql("""
    CREATE TABLE user_behavior (
        user_id STRING,
        item_id STRING,
        behavior STRING,
        ts TIMESTAMP(3),
        WATERMARK FOR ts AS ts - INTERVAL '5' SECOND
    ) WITH (
        'connector' = 'kafka',
        'topic' = 'user_behavior',
        'properties.bootstrap.servers' = 'kafka:9092',
        'properties.group.id' = 'flink-sql',
        'format' = 'json',
        'scan.startup.mode' = 'latest-offset'
    )
""");

// Window TVF 语法(1.13+ 推荐)
tableEnv.executeSql("""
    SELECT window_start, window_end, behavior, COUNT(*) as cnt
    FROM TABLE(
        TUMBLE(TABLE user_behavior, DESCRIPTOR(ts), INTERVAL '10' MINUTES)
    )
    GROUP BY window_start, window_end, behavior
""").print();

5.3 Lookup Join(维表关联)

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

5.4 DataStream 和 Table 互转

// DataStream → Table
Table table = tableEnv.fromDataStream(
    dataStream,
    Schema.newBuilder()
        .columnByExpression("rowtime", "CAST(ts AS TIMESTAMP_LTZ(3))")
        .watermark("rowtime", "rowtime - INTERVAL '5' SECOND")
        .build()
);

// Table → DataStream
tableEnv.toDataStream(table).print();

六、Kafka 集成

6.1 Kafka Source

KafkaSource<String> source = KafkaSource.<String>builder()
    .setBootstrapServers("kafka:9092")
    .setTopics("input-topic")
    .setGroupId("flink-consumer")
    .setStartingOffsets(OffsetsInitializer.committedOffsets(OffsetResetStrategy.EARLIEST))
    .setValueOnlyDeserializer(new SimpleStringSchema())
    .setProperty("isolation.level", "read_committed")
    .build();

env.fromSource(source, WatermarkStrategy.forBoundedOutOfOrderness(...), "Kafka Source")
   .setParallelism(4);

6.2 Kafka Sink(exactly-once)

KafkaSink<String> sink = KafkaSink.<String>builder()
    .setBootstrapServers("kafka:9092")
    .setRecordSerializer(...)
    .setDeliveryGuarantee(DeliveryGuarantee.EXACTLY_ONCE)
    .setTransactionalIdPrefix("flink-job-" + jobName)
    .build();

四要素:Sink 开 EXACTLY_ONCE + transactional.id.prefix、Checkpoint 开 EXACTLY_ONCE、Kafka transaction.max.timeout.ms >= 15min、Source isolation.level=read_committed

6.3 Flink CDC

Flink CDC 就是数据库的"实时监控摄像头" —— 它能实时捕捉数据库里的增删改操作,并把这些变更数据实时同步到下游(数据仓库、消息队列、数据湖等),全程不需要你写复杂的轮询代码。

MySqlSource<String> cdcSource = MySqlSource.<String>builder()
    .hostname("mysql-host")
    .port(3306)
    .databaseList("mydb")
    .tableList("mydb.users,mydb.orders")
    .username("flink")
    .password("xxxxx")
    .deserializer(new JsonDebeziumDeserializationSchema())
    .build();

env.fromSource(cdcSource, WatermarkStrategy.noWatermarks(), "MySQL CDC")
    .addSink(...);

标准链路:Flink CDC → Kafka → Flink 计算 → Doris/ClickHouse


七、项目实战:实时订单统计系统

7.1 业务需求

  1. 实时统计每分钟订单金额、订单数量
  2. 按商品类目分组统计实时销售额 Top10
  3. 实时检测异常订单(金额超过阈值报警)
  4. 订单数据来自 MySQL,通过 CDC 实时同步

7.2 技术架构

MySQL(订单表 orders)
    ↓
Flink CDC Connector(无锁读取,实时捕获增删改)
    ↓
Kafka(orders-topic,数据缓冲和解耦)
    ↓
Flink 计算作业
    ├── 分支1:TumblingWindow(1min) → 每分钟订单金额/数量 → Kafka → 实时大屏
    ├── 分支2:KeyedProcessFunction + ValueState → 按类目累计 Top10 → Redis
    └── 分支3:ProcessFunction + 侧输出流 → 异常订单(金额>10000)→ 钉钉报警

7.3 核心代码

订单数据类

public class Order {
    private Long orderId;
    private Long userId;
    private String category;
    private Double amount;
    private Long eventTime;

    public Order() {}
    // getter/setter 省略...
}

主作业

public class OrderStatsJob {
    public static void main(String[] args) throws Exception {
        StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
        env.setRuntimeMode(RuntimeExecutionMode.STREAMING);
        env.setParallelism(4);

        // Checkpoint
        env.enableCheckpointing(60000, CheckpointingMode.EXACTLY_ONCE);
        env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
        env.getCheckpointConfig().setCheckpointStorage("hdfs:///flink/checkpoints/order-stats");
        EmbeddedRocksDBStateBackend rocksDbBackend = new EmbeddedRocksDBStateBackend(true);
        env.setStateBackend(rocksDbBackend);

        // Kafka Source
        KafkaSource<String> source = KafkaSource.<String>builder()
            .setBootstrapServers("kafka:9092")
            .setTopics("orders-topic")
            .setGroupId("order-stats")
            .setStartingOffsets(OffsetsInitializer.latest())
            .setValueOnlyDeserializer(new SimpleStringSchema())
            .build();

        DataStream<Order> orderStream = env
            .fromSource(source, WatermarkStrategy.noWatermarks(), "Orders")
            .map(json -> parseOrder(json))
            .assignTimestampsAndWatermarks(
                WatermarkStrategy.<Order>forBoundedOutOfOrderness(Duration.ofSeconds(3))
                    .withTimestampAssigner((order, ts) -> order.getEventTime())
                    .withIdleness(Duration.ofMinutes(1))
            );

        // 异常订单侧输出
        OutputTag<Order> abnormalTag = new OutputTag<Order>("abnormal"){};

        // 分支1:每分钟统计
        SingleOutputStreamOperator<Order> mainStream = orderStream
            .process(new AbnormalDetectFunction(abnormalTag, 10000.0));

        mainStream
            .keyBy(Order::getCategory)
            .window(TumblingEventTimeWindows.of(Time.minutes(1)))
            .aggregate(new OrderStatsAggregate())
            .addSink(new KafkaSink<>());

        // 分支2:类目 Top10
        orderStream
            .keyBy(Order::getCategory)
            .process(new CategoryTopNFunction(10));

        // 分支3:异常订单报警
        mainStream.getSideOutput(abnormalTag)
            .addSink(new DingTalkAlertSink());

        env.execute("Order Stats Job");
    }
}

异常检测

public class AbnormalDetectFunction extends ProcessFunction<Order, Order> {
    private final OutputTag<Order> abnormalTag;
    private final double threshold;

    @Override
    public void processElement(Order order, Context ctx, Collector<Order> out) {
        if (order.getAmount() > threshold) {
            ctx.output(abnormalTag, order);
        } else {
            out.collect(order);
        }
    }
}

类目 Top10

public class CategoryTopNFunction extends KeyedProcessFunction<String, Order, TopNResult> {
    private final int topN;
    private ValueState<Double> sumState;

    @Override
    public void open(Configuration parameters) {
        sumState = getRuntimeContext().getState(
            new ValueStateDescriptor<>("sum", Types.DOUBLE));
    }

    @Override
    public void processElement(Order order, Context ctx, Collector<TopNResult> out) {
        Double current = sumState.value();
        sumState.update(current == null ? order.getAmount() : current + order.getAmount());
        ctx.timerService().registerEventTimeTimer(ctx.timestamp() + 10000);
    }

    @Override
    public void onTimer(long timestamp, OnTimerContext ctx, Collector<TopNResult> out) {
        out.collect(new TopNResult(ctx.getCurrentKey(), sumState.value()));
    }
}

7.4 部署配置

jobmanager.memory.process.size: 2048m
taskmanager.memory.process.size: 8192m
taskmanager.memory.managed.fraction: 0.4
taskmanager.numberOfTaskSlots: 4
state.backend: rocksdb
state.checkpoints.dir: hdfs:///flink/checkpoints
execution.checkpointing.interval: 60s
execution.checkpointing.mode: EXACTLY_ONCE

八、生产环境踩坑汇总

8.1 常见问题

现象 原因 解法
窗口不触发 Watermark 没推进 加 withIdleness
Checkpoint 超时 反压/状态太大 非对齐CK/增量CK/调大超时
数据倾斜 某个 key 流量太大 随机前缀打散/rebalance
OOM 状态太大 RocksDB/缩小窗口
Kafka 重复消费 事务超时 调大 timeout
作业重启从头消费 Checkpoint 丢失 开外部化CK保留

8.2 性能优化

  • Source 并行度 = Kafka 分区数
  • 聚合用 AggregateFunction 代替 WindowFunction
  • RocksDB + 增量 Checkpoint
  • Checkpoint 间隔 30s~5min
  • Watermark 延迟 0~10s
  • env.getConfig().enableObjectReuse()
  • 用 POJO/Tuple 类型,避免 Kryo
  • 大状态作业定期 Savepoint

8.3 Metrics 监控

接入 Prometheus + Grafana:

metrics.reporter.prom.class: org.apache.flink.metrics.prometheus.PrometheusReporter
metrics.reporter.prom.port: 9249

重点指标:

  • task.BackPressuredTimeMsPerSecond — 反压
  • checkpoint.duration — Checkpoint 耗时
  • records.inPerSecond / records.outPerSecond — 吞吐

九、总结

数据从哪来(Kafka/MySQL CDC)
    ↓
Watermark(EventTime + 延迟容忍)
    ↓
怎么算(keyBy/window/state/process/AsyncIO/CEP)
    ↓
存到哪(Kafka/Doris/Redis)
    ↓
Checkpoint 保证 exactly-once

主线:数据来源 → 处理计算 → 结果输出,Checkpoint 保证不丢数据。


觉得有用点个赞!

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐