一、Flink 核心特性
- 流批一体:DataStream API 统一处理有界/无界数据
- 事件时间语义:按数据自带时间戳计算,处理乱序数据
- 精确一次(Exactly-Once):Checkpoint + 两阶段提交
- 状态管理:Keyed State / Operator State
二、WordCount(流处理版)
StreamExecutionEnvironment env =
StreamExecutionEnvironment.getExecutionEnvironment();
env.socketTextStream("localhost", 9999)
.flatMap(new FlatMapFunction<String, Tuple2<String, Integer>>() {
@Override
public void flatMap(String value, Collector<Tuple2<String, Integer>> out) {
for (String word : value.split("\\s")) {
out.collect(new Tuple2<>(word, 1));
}
}
})
.keyBy(t -> t.f0)
.sum(1)
.print();
env.execute("Streaming WordCount");
三、事件时间与水位线
WatermarkStrategy<Event> strategy = WatermarkStrategy
.<Event>forBoundedOutOfOrderness(Duration.ofSeconds(5))
.withTimestampAssigner((event, ts) -> event.getTimestamp());
env.addSource(kafkaSource)
.assignTimestampsAndWatermarks(strategy)
.keyBy(Event::getUserId)
.window(TumblingEventTimeWindows.of(Time.minutes(1)))
.aggregate(new CountAgg())
.print();
水位线含义:Watermark(t) 表示不会再有时间戳 ≤ t 的数据到来,用于触发窗口计算。
四、状态管理
public class CountWithState extends RichFlatMapFunction<Event, Tuple2<String, Long>> {
private transient ValueState<Long> countState;
@Override
public void open(Configuration conf) {
ValueStateDescriptor<Long> desc =
new ValueStateDescriptor<>("count", Long.class);
countState = getRuntimeContext().getState(desc);
}
@Override
public void flatMap(Event e, Collector<Tuple2<String, Long>> out) throws Exception {
Long current = countState.value();
if (current == null) current = 0L;
current++;
countState.update(current);
out.collect(new Tuple2<>(e.getUserId(), current));
}
}
五、Checkpoint 配置
env.enableCheckpointing(60000); // 每 60s
env.getCheckpointConfig().setCheckpointingMode(CheckpointingMode.EXACTLY_ONCE);
env.getCheckpointConfig().setMinPauseBetweenCheckpoints(30000);
env.getCheckpointConfig().setCheckpointTimeout(600000);
env.setStateBackend(new RocksDBStateBackend("hdfs://namenode/flink/checkpoints"));
六、常见状态后端
| 后端 | 存储位置 | 适用场景 | | --- | --- | --- | | MemoryStateBackend | JVM 堆 | 本地测试 | | FsStateBackend | 堆 + 远端文件 | 状态较小 | | RocksDBStateBackend | 本地磁盘 | 大状态生产环境 |七、调优要点
- 并行度:等于 Kafka 分区数,避免数据倾斜
- 反压处理:定位到具体算子,常见是外部 IO 慢
- 状态 TTL:给状态设置过期时间,避免无限增长
- 对象复用:
flatMap中复用对象,减少 GC
评论(0)
还没有评论,来说两句吧