【Flink 】从入门到生产实战:一篇文章吃透流处理核心
前言:把自己几年前整理的 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 汇报心跳和状态

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 条触发一次,生产少用 |

窗口 API 用法:
stream
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(5)))
.aggregate(new MyAggregateFunction()) // 增量聚合,性能好
// 或者 .process(new MyWindowFunction()) // 全量聚合,灵活但耗内存
踩坑:
- EventTime 窗口必须设置 Watermark,否则窗口永远不触发。
timeWindow()1.18 已废弃,改用window(TumblingEventTimeWindows.of(...))。- 优先用
AggregateFunction(增量计算,内存友好),WindowFunction是全量加载到内存再算,数据量大容易 OOM。
4.2 EventTime 与 Watermark
Flink 最精妙的设计之一。核心问题:网络延迟导致数据乱序到达,怎么保证算得对?
三种时间语义:
| 时间类型 | 定义 | 适用场景 |
|---|---|---|
| EventTime | 数据真实产生的时间(用户点击的时间戳) | 生产首选,能处理乱序 |
| IngestionTime | 数据进入 Flink Source 的时间 | 极少用 |
| ProcessingTime | 算子执行的系统时间 | 对延迟不敏感的场景 |
Watermark 原理:
Watermark 是一个时间戳标记,告诉系统"所有 EventTime 小于 Watermark 的数据应该都到了"。
Watermark = 当前最大 EventTime - 允许的最大延迟
当 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 恢复。

生产配置(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 调优经验:
- 间隔不要太短(至少 30s),太频繁影响吞吐;也不要太长(超过 5min),恢复时回放太多。
- 反压时 Checkpoint 容易超时(Barrier 对齐要等数据消化完),1.11+ 可开非对齐 Checkpoint:
enableUnalignedCheckpoints(),但会增加状态大小。- RocksDB 增量 Checkpoint 1.18 默认开启。
- 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。
解决反压:
- 增加并行度或资源
- 优化算子逻辑,减少每条数据处理时间
- 调整缓冲大小:
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,性能差):
- public 类
- 有无参构造函数
- 字段 public 或者有 getter/setter
- 字段类型也是 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");
默认所有算子都在
defaultsharing 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 业务需求
- 实时统计每分钟订单金额、订单数量
- 按商品类目分组统计实时销售额 Top10
- 实时检测异常订单(金额超过阈值报警)
- 订单数据来自 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 保证不丢数据。
觉得有用点个赞!
更多推荐




所有评论(0)