Flink 实时计算入门:从 WordCount 到状态管理

一、Flink 核心特性

  • 流批一体:DataStream API 统一处理有界/无界数据
  • 事件时间语义:按数据自带时间戳计算,处理乱序数据
  • 精确一次(Exactly-Once):Checkpoint + 两阶段提交
  • 状态管理:Keyed State / Operator State

二、WordCount(流处理版)

Java
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");

三、事件时间与水位线

Java
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 的数据到来,用于触发窗口计算。

四、状态管理

Java
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 配置

Java
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 | 本地磁盘 | 大状态生产环境 |

七、调优要点

  1. 并行度:等于 Kafka 分区数,避免数据倾斜
  2. 反压处理:定位到具体算子,常见是外部 IO 慢
  3. 状态 TTL:给状态设置过期时间,避免无限增长
  4. 对象复用:flatMap 中复用对象,减少 GC
打赏作者 已有 0 人打赏,共 ¥0.00
我的打赏
大数据老刘
大数据老刘
Lv5 原创 1 粉丝 614

Hadoop / Spark / Flink 离线实时一把梭

  • 1文章
  • 3.3万总阅读
  • 1832获赞
  • 614粉丝

评论(0)

💬

还没有评论,来说两句吧